From 91c5e4f842e9055817d7a0a32726c499a3d95b2b Mon Sep 17 00:00:00 2001 From: Thomas Jungblut Date: Thu, 11 May 2023 15:45:26 +0200 Subject: [PATCH 1/3] remove legacy recovery tests --- test/extended/dr/backup_restore.go | 308 ------------- test/extended/dr/common.go | 138 +----- test/extended/dr/machine_recover.go | 426 ------------------ test/extended/dr/quorum_restore.go | 287 ------------ .../generated/zz_generated.annotations.go | 12 - 5 files changed, 19 insertions(+), 1152 deletions(-) delete mode 100644 test/extended/dr/backup_restore.go delete mode 100644 test/extended/dr/machine_recover.go delete mode 100644 test/extended/dr/quorum_restore.go diff --git a/test/extended/dr/backup_restore.go b/test/extended/dr/backup_restore.go deleted file mode 100644 index c9dd4893819b..000000000000 --- a/test/extended/dr/backup_restore.go +++ /dev/null @@ -1,308 +0,0 @@ -package dr - -import ( - "context" - "fmt" - "io/ioutil" - "os" - "strings" - "time" - - g "github.com/onsi/ginkgo/v2" - o "github.com/onsi/gomega" - - corev1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/labels" - "k8s.io/apimachinery/pkg/util/wait" - "k8s.io/client-go/kubernetes" - corev1client "k8s.io/client-go/kubernetes/typed/core/v1" - "k8s.io/kubernetes/test/e2e/framework" - - exutil "github.com/openshift/origin/test/extended/util" - "github.com/openshift/origin/test/extended/util/disruption" -) - -// Disabled / Obsolete: rewritten in recovery.go due to issues with ssh access -var _ = g.Describe("[sig-etcd][Feature:DisasterRecovery][Disruptive][Disabled:Broken]", func() { - defer g.GinkgoRecover() - - f := framework.NewDefaultFramework("backup-restore") - f.SkipNamespaceCreation = true - - oc := exutil.NewCLIWithoutNamespace("backup-restore") - - // Validate the documented backup and restore procedure as closely as possible: - // - // backup: https://docs.openshift.com/container-platform/4.6/backup_and_restore/backing-up-etcd.html - // restore: https://docs.openshift.com/container-platform/4.6/backup_and_restore/disaster_recovery/scenario-2-restoring-cluster-state.html - // - // Comments like 'Backup 2' and 'Restore '1a' indicate where a test step - // corresponds to a step in the documentation. - // - // Backing up and recovering on the same node is tested by quorum_restore.go. - g.It("[Feature:EtcdRecovery] Cluster should recover from a backup taken on one node and recovered on another [apigroup:operator.openshift.io]", func() { - masters := masterNodes(oc) - // Need one node to backup from and another to restore to - o.Expect(len(masters)).To(o.BeNumerically(">=", 2)) - - // Pick one node to backup on - backupNode := masters[0] - framework.Logf("Selecting node %q as the backup host", backupNode.Name) - - // Recovery 1 - // Pick a different node to recover on - recoveryNode := masters[1] - framework.Logf("Selecting node %q as the recovery host", recoveryNode.Name) - - // Recovery 2 - g.By("Verifying that all masters are reachable via ssh") - for _, master := range masters { - checkSSH(master) - } - - disruptionFunc := func() { - // Backup 4 - // - // The backup has to be taken after the upgrade tests have done - // their pre-disruption setup to ensure that the api state that is - // restored includes those changes. - g.By(fmt.Sprintf("Running the backup script on node %q", backupNode.Name)) - sudoExecOnNodeOrFail(backupNode, "rm -rf /home/core/backup && /usr/local/bin/cluster-backup.sh /home/core/backup && chown -R core /home/core/backup") - - // Recovery 3 - // Copy the backup data from the backup node to the test host and - // from there to the recovery node. - // - // Another solution could be enabling the recovery node to connect - // directly to the backup node. It seemed simpler to use the test host - // as an intermediary rather than enabling agent forwarding or copying - // the private ssh key to the recovery node. - - g.By("Creating a local temporary directory") - tempDir, err := ioutil.TempDir("", "e2e-backup-restore") - o.Expect(err).NotTo(o.HaveOccurred()) - - // Define the ssh configuration necessary to invoke scp, which does - // not appear to be supported by the golang ssh client. - commonOpts := "-o StrictHostKeyChecking=no -o LogLevel=error -o ServerAliveInterval=30 -o ConnectionAttempts=100 -o ConnectTimeout=30" - authOpt := fmt.Sprintf("-i %s", os.Getenv("KUBE_SSH_KEY_PATH")) - bastionHost := os.Getenv("KUBE_SSH_BASTION") - proxyOpt := "" - if len(bastionHost) > 0 { - framework.Logf("Bastion host %s will be used to proxy scp to cluster nodes", bastionHost) - // The bastion host is expected to be of the form address:port - hostParts := strings.Split(bastionHost, ":") - o.Expect(len(hostParts)).To(o.Equal(2)) - address := hostParts[0] - port := hostParts[1] - // A proxy command is required for a bastion host - proxyOpt = fmt.Sprintf("-o ProxyCommand='ssh -A -W %%h:%%p %s %s -p %s core@%s'", commonOpts, authOpt, port, address) - } - - g.By(fmt.Sprintf("Copying the backup directory from backup node %q to the test host", backupNode.Name)) - backupNodeAddress := addressForNode(backupNode) - o.Expect(backupNodeAddress).NotTo(o.BeEmpty()) - copyFromBackupNodeCmd := fmt.Sprintf(`scp -v %s %s %s -r core@%s:backup %s`, commonOpts, authOpt, proxyOpt, backupNodeAddress, tempDir) - runCommandAndRetry(copyFromBackupNodeCmd) - - g.By(fmt.Sprintf("Cleaning the backup path on recovery node %q", recoveryNode.Name)) - sudoExecOnNodeOrFail(recoveryNode, "rm -rf /home/core/backup") - - g.By(fmt.Sprintf("Copying the backup directory from the test host to recovery node %q", recoveryNode.Name)) - recoveryNodeAddress := addressForNode(recoveryNode) - o.Expect(recoveryNodeAddress).NotTo(o.BeEmpty()) - copyToRecoveryNodeCmd := fmt.Sprintf(`scp %s %s %s -r %s/backup core@%s:`, commonOpts, authOpt, proxyOpt, tempDir, recoveryNodeAddress) - runCommandAndRetry(copyToRecoveryNodeCmd) - - // Stop etcd static pods on non-recovery masters. - for _, master := range masters { - // The restore script will stop static pods on the recovery node - if master.Name == recoveryNode.Name { - continue - } - // Recovery 4b - g.By(fmt.Sprintf("Stopping etcd static pod on node %q", master.Name)) - manifest := "/etc/kubernetes/manifests/etcd-pod.yaml" - // Move only if present to ensure idempotent behavior during debugging. - sudoExecOnNodeOrFail(master, fmt.Sprintf("test -f %s && mv -f %s /tmp || true", manifest, manifest)) - - // Recovery 4c - g.By(fmt.Sprintf("Waiting for etcd to exit on node %q", master.Name)) - // Look for 'etcd ' (with trailing space) to be missing to - // differentiate from pods like etcd-operator. - sudoExecOnNodeOrFail(master, "crictl ps | grep 'etcd ' | wc -l | grep -q 0") - - // Recovery 4f - g.By(fmt.Sprintf("Moving etcd data directory on node %q", master.Name)) - // Move only if present to ensure idempotent behavior during debugging. - sudoExecOnNodeOrFail(master, "test -d /var/lib/etcd && (rm -rf /tmp/etcd && mv /var/lib/etcd/ /tmp) || true") - } - - // Recovery 4d - // Trigger stop of kube-apiserver static pods on non-recovery - // masters, without waiting, to minimize the test time required for - // graceful termination to complete. - for _, master := range masters { - // The restore script will stop static pods on the recovery node - if master.Name == recoveryNode.Name { - continue - } - g.By(fmt.Sprintf("Stopping kube-apiserver static pod on node %q", master.Name)) - manifest := "/etc/kubernetes/manifests/kube-apiserver-pod.yaml" - // Move only if present to ensure idempotent behavior during debugging. - sudoExecOnNodeOrFail(master, fmt.Sprintf("test -f %s && mv -f %s /tmp || true", manifest, manifest)) - } - - // Recovery 4e - // Wait for kube-apiserver pods to exit - for _, master := range masters { - // The restore script will stop static pods on the recovery node - if master.Name == recoveryNode.Name { - continue - } - g.By(fmt.Sprintf("Waiting for kube-apiserver to exit on node %q", master.Name)) - // Look for 'kube-apiserver ' (with trailing space) to be missing - // to differentiate from pods like kube-apiserver-operator. - sudoExecOnNodeOrFail(master, "crictl ps | grep -q 'kube-apiserver ' | wc -l | grep -q 0") - } - - // Recovery 7 - restoreFromBackup(recoveryNode) - - // Recovery 8 - for _, master := range masters { - restartKubelet(master) - } - - // Recovery 9a, 9b - waitForAPIServer(oc.AdminKubeClient(), recoveryNode) - - // Recovery 10,11,12 - forceOperandRedeployment(oc.AdminOperatorClient().OperatorV1()) - - // Recovery 13 - waitForReadyEtcdPods(oc.AdminKubeClient(), len(masters)) - - waitForOperatorsToSettle() - } - - disruption.Run(f, "Backup from one node and recover on another", "restore_different_node", - disruption.TestData{}, - disruptionTests, - disruptionFunc, - ) - }) -}) - -// addressForNode looks for an ssh-accessible ip address for a node in case the -// node name doesn't resolve in the test environment. An empty string will be -// returned if an address could not be determined. -func addressForNode(node *corev1.Node) string { - for _, a := range node.Status.Addresses { - if a.Type == corev1.NodeExternalIP && a.Address != "" { - return a.Address - } - } - // No external IPs were found, let's try to use internal as plan B - for _, a := range node.Status.Addresses { - if a.Type == corev1.NodeInternalIP && a.Address != "" { - return a.Address - } - } - return "" -} - -// What follows are helper functions corresponding to steps in the recovery -// procedure. They are defined in a granular fashion to allow reuse by the -// quorum restore test. The quorum restore test needs to interleave the -// standard commands with commands related to master recreation. - -// Recovery 7 -func restoreFromBackup(node *corev1.Node) { - g.By(fmt.Sprintf("Running restore script on recovery node %q", node.Name)) - sudoExecOnNodeOrFail(node, "/usr/local/bin/cluster-restore.sh /home/core/backup") -} - -// Recovery 8 -func restartKubelet(node *corev1.Node) { - g.By(fmt.Sprintf("Restarting the kubelet service on node %q", node.Name)) - sudoExecOnNodeOrFail(node, "systemctl restart kubelet.service") -} - -// Recovery 9a -func waitForEtcdContainer(node *corev1.Node) { - g.By(fmt.Sprintf("Verifying that the etcd container is running on recovery node %q", node.Name)) - // Look for 'etcd ' (with trailing space) to differentiate from pods - // like etcd-operator. - sudoExecOnNodeOrFail(node, "crictl ps | grep -q 'etcd '") -} - -// Recovery 9b -func waitForEtcdPod(node *corev1.Node) { - // The success of this check also ensures that the kube apiserver on - // the recovery node is accepting connections. - g.By(fmt.Sprintf("Verifying that the etcd pod is running on recovery node %q", node.Name)) - // Look for a single running etcd pod - runningEtcdPodCmd := "oc get pods -n openshift-etcd -l k8s-app=etcd --no-headers=true | grep Running | wc -l | grep -q 1" - // The kubeconfig on the node is only readable by root and usage requires sudo. - nodeKubeConfig := "/etc/kubernetes/static-pod-resources/kube-apiserver-certs/secrets/node-kubeconfigs/localhost.kubeconfig" - sudoExecOnNodeOrFail(node, fmt.Sprintf("KUBECONFIG=%s %s", nodeKubeConfig, runningEtcdPodCmd)) -} - -func waitForAPIServerAvailability(client kubernetes.Interface) { - g.By("Waiting for API server to become available") - err := wait.PollImmediate(10*time.Second, 30*time.Minute, func() (done bool, err error) { - _, err = client.CoreV1().Namespaces().Get(context.Background(), "default", metav1.GetOptions{}) - if err != nil { - framework.Logf("Observed an error waiting for apiserver availability outside the cluster: %v", err) - } - return err == nil, nil - }) - o.Expect(err).NotTo(o.HaveOccurred()) -} - -// waitForAPIServer waits for the etcd container and pod running on the -// recovery node and then waits for the apiserver to be accessible outside -// the cluster. -func waitForAPIServer(client kubernetes.Interface, node *corev1.Node) { - // Recovery 9a - waitForEtcdContainer(node) - - // Recovery 9b - waitForEtcdPod(node) - - // Even with the apiserver available on the recovery node, it may - // take additional time for the api to become available externally - // to the cluster. - waitForAPIServerAvailability(client) -} - -// Recovery 13 -func waitForReadyEtcdPods(client kubernetes.Interface, masterCount int) { - g.By(fmt.Sprintf("Waiting for all %d etcd pods to become ready", masterCount)) - waitForPodsTolerateClientTimeout( - client.CoreV1().Pods("openshift-etcd"), - exutil.ParseLabelsOrDie("k8s-app=etcd"), - exutil.CheckPodIsReady, - masterCount, - 40*time.Minute, - ) -} - -func waitForPodsTolerateClientTimeout(c corev1client.PodInterface, label labels.Selector, predicate func(corev1.Pod) bool, count int, timeout time.Duration) { - err := wait.Poll(60*time.Second, timeout, func() (bool, error) { - p, e := exutil.GetPodNamesByFilter(c, label, predicate) - if e != nil { - framework.Logf("Saw an error waiting for etcd pods to become available: %v", e) - // TODO tolerate transient etcd timeout only and fail other errors - return false, nil - } - if len(p) != count { - framework.Logf("Only %d of %d expected pods are ready", len(p), count) - return false, nil - } - return true, nil - }) - o.Expect(err).NotTo(o.HaveOccurred()) -} diff --git a/test/extended/dr/common.go b/test/extended/dr/common.go index 45e4746a9e59..a89d4543c0e7 100644 --- a/test/extended/dr/common.go +++ b/test/extended/dr/common.go @@ -8,7 +8,6 @@ import ( "crypto/x509" "encoding/pem" "fmt" - "os/exec" "strings" "text/tabwriter" "time" @@ -16,13 +15,13 @@ import ( g "github.com/onsi/ginkgo/v2" o "github.com/onsi/gomega" "github.com/openshift/library-go/test/library" - "github.com/openshift/origin/test/e2e/upgrade" exutil "github.com/openshift/origin/test/extended/util" "github.com/openshift/origin/test/extended/util/image" "github.com/stretchr/objx" xssh "golang.org/x/crypto/ssh" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/util/intstr" "k8s.io/apimachinery/pkg/util/wait" @@ -31,10 +30,10 @@ import ( applymetav1 "k8s.io/client-go/applyconfigurations/meta/v1" "k8s.io/client-go/dynamic" "k8s.io/client-go/kubernetes" + corev1client "k8s.io/client-go/kubernetes/typed/core/v1" "k8s.io/kubernetes/test/e2e/framework" e2e "k8s.io/kubernetes/test/e2e/framework" e2epod "k8s.io/kubernetes/test/e2e/framework/pod" - e2essh "k8s.io/kubernetes/test/e2e/framework/ssh" "k8s.io/kubernetes/test/utils" ) @@ -56,30 +55,6 @@ const ( ` ) -func runCommandAndRetry(command string) string { - const ( - maxRetries = 10 - pause = 10 - ) - var ( - retryCount = 0 - out []byte - err error - ) - e2e.Logf("command '%s'", command) - for retryCount = 0; retryCount <= maxRetries; retryCount++ { - out, err = exec.Command("bash", "-c", command).CombinedOutput() - e2e.Logf("output:\n%s", out) - if err == nil { - break - } - e2e.Logf("%v", err) - time.Sleep(time.Second * pause) - } - o.Expect(retryCount).NotTo(o.Equal(maxRetries + 1)) - return string(out) -} - func masterNodes(oc *exutil.CLI) []*corev1.Node { masterNodes, err := oc.AdminKubeClient().CoreV1().Nodes().List(context.Background(), metav1.ListOptions{ LabelSelector: "node-role.kubernetes.io/master", @@ -93,20 +68,6 @@ func masterNodes(oc *exutil.CLI) []*corev1.Node { return nodes } -func clusterNodes(oc *exutil.CLI) (masters, workers []*corev1.Node) { - nodes, err := oc.AdminKubeClient().CoreV1().Nodes().List(context.Background(), metav1.ListOptions{}) - o.Expect(err).NotTo(o.HaveOccurred()) - for i := range nodes.Items { - node := &nodes.Items[i] - if _, ok := node.Labels["node-role.kubernetes.io/master"]; ok { - masters = append(masters, node) - } else { - workers = append(workers, node) - } - } - return -} - func internalIP(restoreNode *corev1.Node) string { internalIp := restoreNode.Name for _, address := range restoreNode.Status.Addresses { @@ -117,15 +78,6 @@ func internalIP(restoreNode *corev1.Node) string { return internalIp } -func waitForMastersToUpdate(oc *exutil.CLI, mcps dynamic.NamespaceableResourceInterface) { - e2e.Logf("Waiting for MachineConfig master to finish rolling out") - err := wait.Poll(30*time.Second, 30*time.Minute, func() (done bool, err error) { - done, _ = upgrade.IsPoolUpdated(mcps, "master") - return done, nil - }) - o.Expect(err).NotTo(o.HaveOccurred()) -} - func waitForOperatorsToSettle() { g.By("Waiting for operators to settle before performing post-disruption testing") config, err := framework.LoadConfig() @@ -247,75 +199,6 @@ func condition(cv objx.Map, condition string) objx.Map { return objx.Map(nil) } -func nodeConditionStatus(conditions []corev1.NodeCondition, conditionType corev1.NodeConditionType) corev1.ConditionStatus { - for _, condition := range conditions { - if condition.Type == conditionType { - return condition.Status - } - } - return corev1.ConditionUnknown -} - -func countReady(items []corev1.Node) int { - ready := 0 - for _, item := range items { - if nodeConditionStatus(item.Status.Conditions, corev1.NodeReady) == corev1.ConditionTrue { - ready++ - } - } - return ready -} - -func fetchFileContents(node *corev1.Node, path string) string { - e2e.Logf("Fetching %s file contents from %s", path, node.Name) - out := execOnNodeWithOutputOrFail(node, fmt.Sprintf("cat %q", path)) - return out.Stdout -} - -// execOnNodeWithOutputOrFail executes a command via ssh against a -// node in a poll loop to ensure reliable execution in a disrupted -// environment. The calling test will be failed if the command cannot -// be executed successfully before the provided timeout. -func execOnNodeWithOutputOrFail(node *corev1.Node, cmd string) *e2essh.Result { - var out *e2essh.Result - var err error - waitErr := wait.PollImmediate(5*time.Second, defaultSSHTimeout, func() (bool, error) { - out, err = e2essh.IssueSSHCommandWithResult(context.TODO(), cmd, e2e.TestContext.Provider, node) - // IssueSSHCommandWithResult logs output - if err != nil { - e2e.Logf("Failed to exec cmd [%s] on node %s: %v", cmd, node.Name, err) - } - return err == nil, nil - }) - o.Expect(waitErr).NotTo(o.HaveOccurred()) - return out -} - -// execOnNodeOrFail executes a command via ssh against a node in a -// poll loop until success or timeout. The output is ignored. The -// calling test will be failed if the command cannot be executed -// successfully before the timeout. -func execOnNodeOrFail(node *corev1.Node, cmd string) { - _ = execOnNodeWithOutputOrFail(node, cmd) -} - -// sudoExecOnNodeOrFail executes a command under sudo with execOnNodeOrFail. -func sudoExecOnNodeOrFail(node *corev1.Node, cmd string) { - sudoCmd := fmt.Sprintf(`sudo -i /bin/bash -cx "%s"`, cmd) - execOnNodeOrFail(node, sudoCmd) -} - -// checkSSH repeatedly attempts to establish an ssh connection to a -// node and fails the calling test if unable to establish the -// connection before the default timeout. -func checkSSH(node *corev1.Node) { - _ = execOnNodeWithOutputOrFail(node, "true") -} - -func ssh(cmd string, node *corev1.Node) (*e2essh.Result, error) { - return e2essh.IssueSSHCommandWithResult(context.TODO(), cmd, e2e.TestContext.Provider, node) -} - func waitForReadyEtcdStaticPods(client kubernetes.Interface, masterCount int) { g.By("Waiting for all etcd static pods to become ready") waitForPodsTolerateClientTimeout( @@ -327,6 +210,23 @@ func waitForReadyEtcdStaticPods(client kubernetes.Interface, masterCount int) { ) } +func waitForPodsTolerateClientTimeout(c corev1client.PodInterface, label labels.Selector, predicate func(corev1.Pod) bool, count int, timeout time.Duration) { + err := wait.Poll(60*time.Second, timeout, func() (bool, error) { + p, e := exutil.GetPodNamesByFilter(c, label, predicate) + if e != nil { + framework.Logf("Saw an error waiting for etcd pods to become available: %v", e) + // TODO tolerate transient etcd timeout only and fail other errors + return false, nil + } + if len(p) != count { + framework.Logf("Only %d of %d expected pods are ready", len(p), count) + return false, nil + } + return true, nil + }) + o.Expect(err).NotTo(o.HaveOccurred()) +} + // InstallSSHKeyOnControlPlaneNodes will create a new private/public ssh keypair, // create a new secret for both in the openshift-etcd-operator namespace. Then it // will append the public key on the host core user authorized_keys file with a daemon set. diff --git a/test/extended/dr/machine_recover.go b/test/extended/dr/machine_recover.go deleted file mode 100644 index a8f9ea33d82f..000000000000 --- a/test/extended/dr/machine_recover.go +++ /dev/null @@ -1,426 +0,0 @@ -package dr - -import ( - "context" - "fmt" - "io/ioutil" - "math/rand" - "net" - "os" - "path" - "strings" - "time" - - g "github.com/onsi/ginkgo/v2" - o "github.com/onsi/gomega" - "go.etcd.io/etcd/api/v3/etcdserverpb" - "go.etcd.io/etcd/client/pkg/v3/transport" - clientv3 "go.etcd.io/etcd/client/v3" - "google.golang.org/grpc" - - corev1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" - "k8s.io/apimachinery/pkg/runtime/schema" - "k8s.io/apimachinery/pkg/util/sets" - "k8s.io/apimachinery/pkg/util/wait" - "k8s.io/client-go/dynamic" - "k8s.io/client-go/kubernetes" - "k8s.io/client-go/rest" - "k8s.io/kubernetes/test/e2e/framework" - "k8s.io/kubernetes/test/e2e/upgrades" - "k8s.io/kubernetes/test/e2e/upgrades/network" - - exutil "github.com/openshift/origin/test/extended/util" - "github.com/openshift/origin/test/extended/util/disruption" - "github.com/openshift/origin/test/extended/util/disruption/controlplane" - "github.com/openshift/origin/test/extended/util/disruption/frontends" -) - -var _ = g.Describe("[sig-cluster-lifecycle][Feature:DisasterRecovery][Disruptive]", func() { - f := framework.NewDefaultFramework("machine-recovery") - f.SkipNamespaceCreation = true - - oc := exutil.NewCLIWithoutNamespace("machine-recovery") - - g.It("[Feature:NodeRecovery] Cluster should survive master and worker failure and recover with machine health checks [apigroup:machine.openshift.io]", func() { - - framework.Logf("Verify SSH is available before restart") - masters, workers := clusterNodes(oc) - o.Expect(len(masters)).To(o.BeNumerically(">=", 3)) - o.Expect(len(workers)).To(o.BeNumerically(">=", 2)) - - replacedMaster := masters[rand.Intn(len(masters))] - checkSSH(replacedMaster) - - replacedWorker := workers[rand.Intn(len(workers))] - checkSSH(replacedWorker) - - disruption.Run(f, "Machine Shutdown and Restore", "machine_failure", - disruption.TestData{}, - []upgrades.Test{ - &network.ServiceUpgradeTest{}, - controlplane.NewKubeAvailableWithNewConnectionsTest(), - controlplane.NewOpenShiftAvailableNewConnectionsTest(), - controlplane.NewOAuthAvailableNewConnectionsTest(), - frontends.NewOAuthRouteAvailableWithNewConnectionsTest(), - frontends.NewOAuthRouteAvailableWithConnectionReuseTest(), - frontends.NewConsoleRouteAvailableWithNewConnectionsTest(), - frontends.NewConsoleRouteAvailableWithConnectionReuseTest(), - }, - func() { - - config, err := framework.LoadConfig() - o.Expect(err).NotTo(o.HaveOccurred()) - dynamicClient := dynamic.NewForConfigOrDie(config) - ms := dynamicClient.Resource(schema.GroupVersionResource{ - Group: "machine.openshift.io", - Version: "v1beta1", - Resource: "machines", - }).Namespace("openshift-machine-api") - - // framework.Logf("Verify etcd endpoints are healthy") - // certsDir := "/tmp/etcd-client-certs/" - // dumpEtcdCertsOnDisk(oc, certsDir) - // defer os.RemoveAll(certsDir) - // unhealthyEtcds := getUnhealthyEtcds(certsDir, masters) - // o.Expect(len(unhealthyEtcds)).To(o.BeNumerically("==", 0)) - - // createMachineHealthCheckForRole("master") - // defer deleteMachineCheckForRole("master") - createMachineHealthCheckForRole("worker") - defer deleteMachineCheckForRole("worker") - - replacedWorkerMachineName := getMachineNameByNodeName(oc, replacedWorker.Name) - // replacedMasterMachineName := getMachineNameByNodeName(oc, replacedMaster.Name) - // replacedMasterMachine, err := ms.Get(context.Background(), replacedMasterMachineName, metav1.GetOptions{}) - // o.Expect(err).NotTo(o.HaveOccurred()) - - // add controller reference - // replacedMasterMachineCopy := replacedMasterMachine.DeepCopy() - // ownerReferences := getOwnerReferenceForMasterMachine(replacedMasterMachineCopy) - // replacedMasterMachineCopy.SetOwnerReferences(ownerReferences) - // _, err = ms.Update(context.Background(), replacedMasterMachineCopy, metav1.UpdateOptions{}) - // o.Expect(err).NotTo(o.HaveOccurred()) - - targets := []*corev1.Node{ /* replacedMaster , */ replacedWorker} - targetMachineNames := []string{ /* replacedMasterMachineName, */ replacedWorkerMachineName} - - // we use a hard shutdown to simulate a poweroff - for _, target := range targets { - framework.Logf("Forcing shutdown of node %s", target.Name) - _, err = ssh("sudo -i systemctl poweroff --force --force", target) - } - - pollConfig := rest.CopyConfig(config) - pollConfig.Timeout = 5 * time.Second - pollClient, err := kubernetes.NewForConfig(pollConfig) - o.Expect(err).NotTo(o.HaveOccurred()) - - framework.Logf("Wait for nodes to go unready") - err = wait.Poll(30*time.Second, 30*time.Minute, func() (done bool, err error) { - nodes, err := pollClient.CoreV1().Nodes().List(context.Background(), metav1.ListOptions{}) - if err != nil || nodes.Items == nil { - framework.Logf("return false - err %v nodes.Items %v", err, nodes.Items) - return false, nil - } - notReadyNodes := sets.NewString() - for _, node := range nodes.Items { - nodeReady := true - for _, t := range node.Spec.Taints { - if t.Key == "node.kubernetes.io/unreachable" { - nodeReady = false - break - } - } - if !nodeReady { - notReadyNodes.Insert(node.Name) - } - } - if notReadyNodes.Len() != len(targets) { - framework.Logf("Nodes waiting to go unready: %v", notReadyNodes.List()) - return false, nil - } - return true, nil - }) - o.Expect(err).NotTo(o.HaveOccurred()) - - // etcdMemberToRemove := getEtcdMemberToRemove(oc, replacedMaster.Name) - // framework.Logf("remove etcd member with ID %s", etcdMemberToRemove) - // removeMember(oc, etcdMemberToRemove) - - // framework.Logf("Recreating master %s", replacedMasterMachineName+"-duplicate") - // newMaster := replacedMasterMachine.DeepCopy() - // // The providerID is relied upon by the machine controller to determine a machine - // // has been provisioned - // // https://github.com/openshift/cluster-api/blob/c4a461a19efb8a25b58c630bed0829512d244ba7/pkg/controller/machine/controller.go#L306-L308 - // unstructured.SetNestedField(newMaster.Object, "", "spec", "providerID") - // newMaster.SetName(replacedMasterMachineName + "-duplicate") - // newMaster.SetResourceVersion("") - // newMaster.SetSelfLink("") - // newMaster.SetUID("") - // newMaster.SetCreationTimestamp(metav1.NewTime(time.Time{})) - // // retry until the machine gets created - // err = wait.PollImmediate(5*time.Second, 10*time.Minute, func() (bool, error) { - // _, err := ms.Create(context.Background(), newMaster, metav1.CreateOptions{}) - // if errors.IsAlreadyExists(err) { - // framework.Logf("Waiting for old machine object %s to be deleted so we can create a new one", replacedMaster.Name) - // return false, nil - // } - // if err != nil { - // return false, err - // } - // return true, nil - // }) - // o.Expect(err).NotTo(o.HaveOccurred()) - - // framework.Logf("Wait for masters to join as nodes and go ready") - // err = wait.Poll(30*time.Second, 30*time.Minute, func() (done bool, err error) { - // defer func() { - // if r := recover(); r != nil { - // fmt.Println("Recovered from panic", r) - // } - // }() - // nodes, err := oc.AdminKubeClient().CoreV1().Nodes().List(context.Background(), metav1.ListOptions{LabelSelector: "node-role.kubernetes.io/master="}) - // if err != nil { - // return false, err - // } - // ready := countReady(nodes.Items) - // if ready != len(masters) { - // framework.Logf("%d master nodes still unready", len(masters)-ready) - // return false, nil - // } - // return true, nil - // }) - // o.Expect(err).NotTo(o.HaveOccurred()) - - framework.Logf("Wait for worker to join as nodes and go ready") - err = wait.Poll(30*time.Second, 30*time.Minute, func() (done bool, err error) { - defer func() { - if r := recover(); r != nil { - fmt.Println("Recovered from panic", r) - } - }() - nodes, err := oc.AdminKubeClient().CoreV1().Nodes().List(context.Background(), metav1.ListOptions{LabelSelector: "node-role.kubernetes.io/worker="}) - if err != nil { - return false, err - } - ready := countReady(nodes.Items) - if ready != len(workers) { - framework.Logf("%d worker nodes still unready", len(workers)-ready) - return false, nil - } - return true, nil - }) - o.Expect(err).NotTo(o.HaveOccurred()) - - framework.Logf("Wait for old machines to be deleted") - err = wait.Poll(30*time.Second, 30*time.Minute, func() (done bool, err error) { - machines, err := ms.List(context.Background(), metav1.ListOptions{}) - if err != nil || machines.Items == nil { - framework.Logf("return false - err %v nodes.Items %v", err, machines.Items) - return false, nil - } - vanishedMachines := sets.NewString() - for _, machine := range targetMachineNames { - vanishedMachines.Insert(machine) - } - for _, machine := range machines.Items { - vanishedMachines.Delete(machine.GetName()) - } - if vanishedMachines.Len() != len(targetMachineNames) { - framework.Logf("Machines waiting to go be deleted: %v", vanishedMachines.List()) - return false, nil - } - return true, nil - }) - o.Expect(err).NotTo(o.HaveOccurred()) - }) - }, - ) -}) - -func getEtcdMemberToRemove(oc *exutil.CLI, unhealthyNodeName string) string { - var healthyEtcdPod string - nodes, err := oc.AdminKubeClient().CoreV1().Nodes().List(context.Background(), metav1.ListOptions{LabelSelector: "node-role.kubernetes.io/master="}) - o.Expect(err).NotTo(o.HaveOccurred()) - for _, node := range nodes.Items { - nodeReady := true - for _, t := range node.Spec.Taints { - if t.Key == "node.kubernetes.io/unreachable" { - nodeReady = false - break - } - } - if nodeReady { - healthyEtcdPod = "etcd-" + node.Name - break - } - } - o.Expect(err).NotTo(o.HaveOccurred()) - - var memberListOutput string - // give 2 mins for api to be up and retry - err = wait.Poll(2*time.Second, 2*time.Minute, func() (done bool, err error) { - memberListOutput, err = oc.AsAdmin().Run("exec").Args("-n", "openshift-etcd", healthyEtcdPod, "-c", "etcdctl", "--", "etcdctl", "memberListOutput", "list").Output() - if err != nil { - return false, nil - } - return true, nil - }) - for _, memberLine := range strings.Split(memberListOutput, "\n") { - if strings.Contains(memberLine, unhealthyNodeName) { - return strings.Split(memberLine, ", ")[0] - } - } - o.Expect(fmt.Errorf("could not find memberListOutput name %s in memberListOutput output %s", unhealthyNodeName, memberListOutput)).NotTo(o.HaveOccurred()) - return "" -} - -func deleteMachineCheckForRole(role string) { - config, err := framework.LoadConfig() - o.Expect(err).NotTo(o.HaveOccurred()) - dynamicClient := dynamic.NewForConfigOrDie(config) - mhc := dynamicClient.Resource(schema.GroupVersionResource{ - Group: "machine.openshift.io", - Version: "v1beta1", - Resource: "machinehealthchecks", - }).Namespace("openshift-machine-api") - err = mhc.Delete(context.Background(), "e2e-health-check-"+role, metav1.DeleteOptions{}) - o.Expect(err).ToNot(o.HaveOccurred()) -} - -func createMachineHealthCheckForRole(role string) { - config, err := framework.LoadConfig() - o.Expect(err).NotTo(o.HaveOccurred()) - dynamicClient := dynamic.NewForConfigOrDie(config) - mhc := dynamicClient.Resource(schema.GroupVersionResource{ - Group: "machine.openshift.io", - Version: "v1beta1", - Resource: "machinehealthchecks", - }).Namespace("openshift-machine-api") - u := &unstructured.Unstructured{} - u.SetGroupVersionKind(schema.GroupVersionKind{ - Group: "machine.openshift.io", - Version: "v1beta1", - Kind: "MachineHealthCheck", - }) - u.SetName("e2e-health-check-" + role) - u.SetNamespace("openshift-machine-api") - err = unstructured.SetNestedField(u.Object, role, "spec", "selector", "matchLabels", "machine.openshift.io/cluster-api-machine-role") - o.Expect(err).ToNot(o.HaveOccurred()) - err = unstructured.SetNestedField(u.Object, []interface{}{ - map[string]interface{}{ - "type": "Ready", - "timeout": "5m", - "status": "False", - }, - map[string]interface{}{ - "type": "Ready", - "timeout": "5m", - "status": "Unknown", - }, - }, "spec", "unhealthyConditions") - o.Expect(err).ToNot(o.HaveOccurred()) - _, err = mhc.Create(context.Background(), u, metav1.CreateOptions{}) - o.Expect(err).ToNot(o.HaveOccurred()) -} - -func getOwnerReferenceForMasterMachine(obj metav1.Object) []metav1.OwnerReference { - o := metav1.NewControllerRef(obj, schema.GroupVersionKind{ - Group: "machine.openshift.io", - Version: "v1beta1", - Kind: "MachineSet", - }) - return []metav1.OwnerReference{*o} -} - -func getUnhealthyEtcds(certsDir string, masters []*corev1.Node) []*etcdserverpb.Member { - endpoints := getEtcdEndpoints(masters) - etcdClient, err := getEtcdClient(certsDir, endpoints) - o.Expect(err).ToNot(o.HaveOccurred()) - memberListResp, err := etcdClient.MemberList(context.Background()) - o.Expect(err).ToNot(o.HaveOccurred()) - - unhealthEtcds := []*etcdserverpb.Member{} - - for _, m := range memberListResp.Members { - _, err := etcdClient.Status(context.Background(), m.ClientURLs[0]) - if err == nil { - unhealthEtcds = append(unhealthEtcds, m) - } - } - return unhealthEtcds -} - -func getEtcdClient(certsDir string, endpoints []string) (*clientv3.Client, error) { - dialOptions := []grpc.DialOption{ - grpc.WithBlock(), // block until the underlying connection is up - } - - tlsInfo := transport.TLSInfo{ - CertFile: path.Join(certsDir, "tls.crt"), - KeyFile: path.Join(certsDir, "tls.key"), - TrustedCAFile: path.Join(certsDir, "ca-bundle.crt"), - } - tlsConfig, err := tlsInfo.ClientConfig() - - cfg := &clientv3.Config{ - DialOptions: dialOptions, - Endpoints: endpoints, - DialTimeout: 15 * time.Second, - TLS: tlsConfig, - } - - cli, err := clientv3.New(*cfg) - if err != nil { - return nil, err - } - return cli, err -} - -func dumpEtcdCertsOnDisk(oc *exutil.CLI, dir string) { - err := os.MkdirAll(dir, os.ModePerm) - o.Expect(err).ToNot(o.HaveOccurred()) - - etcdCA, err := oc.AdminKubeClient().CoreV1().ConfigMaps("openshift-config").Get(context.Background(), "etcd-ca-bundle", metav1.GetOptions{}) - o.Expect(err).ToNot(o.HaveOccurred()) - caData, ok := etcdCA.Data["ca-bundle.crt"] - if !ok { - o.Expect(fmt.Errorf("etcd CA data missing in configmap openshift-config/etcd-ca-bundle")).ToNot(o.HaveOccurred()) - } - - etcdClientCerts, err := oc.AdminKubeClient().CoreV1().Secrets("openshift-config").Get(context.Background(), "etcd-client", metav1.GetOptions{}) - o.Expect(err).ToNot(o.HaveOccurred()) - clientCert, ok := etcdClientCerts.Data["tls.crt"] - if !ok { - o.Expect(fmt.Errorf("etcd client Certificate data missing in secret openshift-config/etcd-client")).ToNot(o.HaveOccurred()) - } - clientKey, ok := etcdClientCerts.Data["tls.key"] - if !ok { - o.Expect(fmt.Errorf("etcd client Private Key data missing in secret openshift-config/etcd-client")).ToNot(o.HaveOccurred()) - } - - err = ioutil.WriteFile(path.Join(dir, "ca-bundle.crt"), []byte(caData), 0600) - o.Expect(err).NotTo(o.HaveOccurred()) - - err = ioutil.WriteFile(path.Join(dir, "tls.crt"), []byte(clientCert), 0600) - o.Expect(err).NotTo(o.HaveOccurred()) - - err = ioutil.WriteFile(path.Join(dir, "tls.key"), []byte(clientKey), 0600) - o.Expect(err).NotTo(o.HaveOccurred()) -} - -func getEtcdEndpoints(masters []*corev1.Node) []string { - endpoints := []string{} - for _, m := range masters { - for _, addr := range m.Status.Addresses { - if addr.Type == corev1.NodeInternalIP { - endpoints = append(endpoints, "https://"+net.JoinHostPort(addr.Address, "2379")) - break - } - } - } - - return endpoints -} diff --git a/test/extended/dr/quorum_restore.go b/test/extended/dr/quorum_restore.go deleted file mode 100644 index 1f48a63ca28c..000000000000 --- a/test/extended/dr/quorum_restore.go +++ /dev/null @@ -1,287 +0,0 @@ -package dr - -import ( - "context" - "fmt" - "math/rand" - "strings" - "time" - - g "github.com/onsi/ginkgo/v2" - o "github.com/onsi/gomega" - - "k8s.io/apimachinery/pkg/api/errors" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" - "k8s.io/apimachinery/pkg/runtime/schema" - "k8s.io/apimachinery/pkg/util/wait" - "k8s.io/client-go/dynamic" - "k8s.io/client-go/kubernetes" - "k8s.io/client-go/rest" - "k8s.io/kubernetes/test/e2e/framework" - e2eskipper "k8s.io/kubernetes/test/e2e/framework/skipper" - "k8s.io/kubernetes/test/e2e/upgrades" - "k8s.io/kubernetes/test/e2e/upgrades/apps" - "k8s.io/kubernetes/test/e2e/upgrades/network" - "k8s.io/kubernetes/test/e2e/upgrades/node" - - exutil "github.com/openshift/origin/test/extended/util" - "github.com/openshift/origin/test/extended/util/disruption" -) - -const ( - machineAnnotationName = "machine.openshift.io/machine" -) - -var disruptionTests []upgrades.Test = []upgrades.Test{ - &network.ServiceUpgradeTest{}, - &node.SecretUpgradeTest{}, - &apps.ReplicaSetUpgradeTest{}, - &apps.StatefulSetUpgradeTest{}, - &apps.DeploymentUpgradeTest{}, - &apps.DaemonSetUpgradeTest{}, -} - -var _ = g.Describe("[sig-etcd][Feature:DisasterRecovery][Disruptive][Disabled:Broken]", func() { - defer g.GinkgoRecover() - - f := framework.NewDefaultFramework("disaster-recovery") - f.SkipNamespaceCreation = true - - oc := exutil.NewCLIWithoutNamespace("disaster-recovery") - - // Validate backing up and restoring to the same node on a cluster - // that has lost quorum after the backup was taken. - g.It("[Feature:EtcdRecovery] Cluster should restore itself after quorum loss [apigroup:machine.openshift.io][apigroup:operator.openshift.io]", func() { - config, err := framework.LoadConfig() - o.Expect(err).NotTo(o.HaveOccurred()) - dynamicClient := dynamic.NewForConfigOrDie(config) - ms := dynamicClient.Resource(schema.GroupVersionResource{ - Group: "machine.openshift.io", - Version: "v1beta1", - Resource: "machines", - }).Namespace("openshift-machine-api") - mcps := dynamicClient.Resource(schema.GroupVersionResource{ - Group: "machineconfiguration.openshift.io", - Version: "v1", - Resource: "machineconfigpools", - }) - - // test for machines as a proxy for "can we recover a master" - machines, err := dynamicClient.Resource(schema.GroupVersionResource{ - Group: "machine.openshift.io", - Version: "v1beta1", - Resource: "machines", - }).List(context.Background(), metav1.ListOptions{}) - o.Expect(err).NotTo(o.HaveOccurred()) - if len(machines.Items) == 0 { - e2eskipper.Skipf("machine API is not enabled and automatic recovery test is not possible") - } - - disruption.Run(f, "Quorum Loss and Restore", "quorum_restore", - disruption.TestData{}, - disruptionTests, - func() { - - framework.Logf("Verify SSH is available before restart") - masters := masterNodes(oc) - o.Expect(len(masters)).To(o.BeNumerically(">=", 1)) - survivingNode := masters[rand.Intn(len(masters))] - survivingNodeName := survivingNode.Name - checkSSH(survivingNode) - - expectedNumberOfMasters := len(masters) - survivingMachineName := getMachineNameByNodeName(oc, survivingNodeName) - survivingMachine, err := ms.Get(context.Background(), survivingMachineName, metav1.GetOptions{}) - o.Expect(err).NotTo(o.HaveOccurred()) - - // The backup script only supports taking a backup of a cluster - // that still has quorum, so the backup must be performed before - // quorum-destroying machine deletion. - framework.Logf("Perform etcd backup on node %s (machine %s) while quorum still exists", survivingNodeName, survivingMachineName) - execOnNodeOrFail(survivingNode, "sudo -i /bin/bash -cx 'rm -rf /home/core/backup && /usr/local/bin/cluster-backup.sh /home/core/backup'") - - framework.Logf("Destroy %d masters", len(masters)-1) - var masterMachines []string - for _, node := range masters { - masterMachine := getMachineNameByNodeName(oc, node.Name) - masterMachines = append(masterMachines, masterMachine) - - if node.Name == survivingNodeName { - continue - } - - framework.Logf("Destroying %s", masterMachine) - err = ms.Delete(context.Background(), masterMachine, metav1.DeleteOptions{}) - o.Expect(err).NotTo(o.HaveOccurred()) - } - - // All API calls for the remainder of the test should be performed in - // a polling loop to insure against transient failures. API calls can - // only be assumed to succeed without polling against a healthy - // cluster, and only at the successful exit of this function will - // that once again be the case. - - pollConfig := rest.CopyConfig(config) - pollConfig.Timeout = 5 * time.Second - pollClient, err := kubernetes.NewForConfig(pollConfig) - o.Expect(err).NotTo(o.HaveOccurred()) - - if len(masters) != 1 { - framework.Logf("Wait for control plane to become unresponsive (may take several minutes)") - failures := 0 - err = wait.Poll(5*time.Second, 30*time.Minute, func() (done bool, err error) { - _, err = pollClient.CoreV1().Nodes().List(context.Background(), metav1.ListOptions{}) - if err != nil { - framework.Logf("Error seen checking for unresponsive control plane: %v", err) - failures++ - } else { - failures = 0 - } - - // wait to see the control plane go down for good to avoid a transient failure - return failures > 4, nil - }) - } - - // Recovery 7 - restoreFromBackup(survivingNode) - - // Recovery 8 - restartKubelet(survivingNode) - - // Recovery 9a, 9b - waitForAPIServer(oc.AdminKubeClient(), survivingNode) - - // Restoring brings back machines and nodes deleted since the - // backup was taken. Those machines and nodes need to be removed - // before they can be created again. - // - // TODO(marun) Ensure the mechanics of node replacement around - // disaster recovery are documented. - for _, master := range masterMachines { - if master == survivingMachineName { - continue - } - err := wait.PollImmediate(5*time.Second, 5*time.Minute, func() (bool, error) { - framework.Logf("Initiating deletion of machine removed after the backup was taken: %s", master) - err := ms.Delete(context.Background(), master, metav1.DeleteOptions{}) - if err != nil && !errors.IsNotFound(err) { - framework.Logf("Error seen when attempting to remove restored machine %s: %v", master, err) - return false, nil - } - return true, nil - }) - o.Expect(err).NotTo(o.HaveOccurred()) - } - for _, node := range masters { - if node.Name == survivingNodeName { - continue - } - err := wait.PollImmediate(5*time.Second, 5*time.Minute, func() (bool, error) { - framework.Logf("Initiating deletion of node removed after the backup was taken: %s", node.Name) - err := oc.AdminKubeClient().CoreV1().Nodes().Delete(context.Background(), node.Name, metav1.DeleteOptions{}) - if err != nil && !errors.IsNotFound(err) { - framework.Logf("Error seen when attempting to remove restored node %s: %v", node.Name, err) - return false, nil - } - return true, nil - }) - o.Expect(err).NotTo(o.HaveOccurred()) - } - - if expectedNumberOfMasters == 1 { - framework.Logf("Cannot create new masters, you must manually create masters and update their DNS entries according to the docs") - } else { - framework.Logf("Create new masters") - for _, master := range masterMachines { - if master == survivingMachineName { - continue - } - framework.Logf("Creating master %s", master) - newMaster := survivingMachine.DeepCopy() - // The providerID is relied upon by the machine controller to determine a machine - // has been provisioned - // https://github.com/openshift/cluster-api/blob/c4a461a19efb8a25b58c630bed0829512d244ba7/pkg/controller/machine/controller.go#L306-L308 - unstructured.SetNestedField(newMaster.Object, "", "spec", "providerID") - newMaster.SetName(master) - newMaster.SetResourceVersion("") - newMaster.SetSelfLink("") - newMaster.SetUID("") - newMaster.SetCreationTimestamp(metav1.NewTime(time.Time{})) - // retry until the machine gets created - err := wait.PollImmediate(5*time.Second, 10*time.Minute, func() (bool, error) { - _, err := ms.Create(context.Background(), newMaster, metav1.CreateOptions{}) - if errors.IsAlreadyExists(err) { - framework.Logf("Waiting for old machine object %s to be deleted so we can create a new one", master) - return false, nil - } - if err != nil { - framework.Logf("Error seen when re-creating machines: %v", err) - return false, nil - } - return true, nil - }) - o.Expect(err).NotTo(o.HaveOccurred()) - } - - framework.Logf("Waiting for machines to be created") - err = wait.Poll(30*time.Second, 20*time.Minute, func() (done bool, err error) { - mastersList, err := ms.List(context.Background(), metav1.ListOptions{ - LabelSelector: "machine.openshift.io/cluster-api-machine-role=master", - }) - if err != nil { - framework.Logf("Failed to check that machines are created: %v", err) - return false, nil - } - if mastersList.Items == nil { - return false, nil - } - return len(mastersList.Items) == expectedNumberOfMasters, nil - }) - o.Expect(err).NotTo(o.HaveOccurred()) - - framework.Logf("Wait for masters to join as nodes and go ready") - err = wait.Poll(30*time.Second, 50*time.Minute, func() (done bool, err error) { - defer func() { - if r := recover(); r != nil { - fmt.Println("Recovered from panic", r) - } - }() - nodes, err := oc.AdminKubeClient().CoreV1().Nodes().List(context.Background(), metav1.ListOptions{LabelSelector: "node-role.kubernetes.io/master="}) - if err != nil { - // scale up to 2nd etcd will make this error inevitable - framework.Logf("Error seen attempting to list master nodes: %v", err) - return false, nil - } - ready := countReady(nodes.Items) - if ready != expectedNumberOfMasters { - framework.Logf("%d nodes still unready", expectedNumberOfMasters-ready) - return false, nil - } - return true, nil - }) - o.Expect(err).NotTo(o.HaveOccurred()) - } - - // Recovery 10,11,12 - forceOperandRedeployment(oc.AdminOperatorClient().OperatorV1()) - - // Recovery 13 - waitForReadyEtcdPods(oc.AdminKubeClient(), expectedNumberOfMasters) - - waitForMastersToUpdate(oc, mcps) - waitForOperatorsToSettle() - }) - }, - ) -}) - -func getMachineNameByNodeName(oc *exutil.CLI, name string) string { - masterNode, err := oc.AdminKubeClient().CoreV1().Nodes().Get(context.Background(), name, metav1.GetOptions{}) - o.Expect(err).NotTo(o.HaveOccurred()) - - annotations := masterNode.GetAnnotations() - o.Expect(annotations).To(o.HaveKey(machineAnnotationName)) - return strings.Split(annotations[machineAnnotationName], "/")[1] -} diff --git a/test/extended/util/annotate/generated/zz_generated.annotations.go b/test/extended/util/annotate/generated/zz_generated.annotations.go index 5cd11d3e8c6d..fb8ac93cd5cf 100644 --- a/test/extended/util/annotate/generated/zz_generated.annotations.go +++ b/test/extended/util/annotate/generated/zz_generated.annotations.go @@ -701,12 +701,6 @@ var Annotations = map[string]string{ "[sig-arch][Early] Managed cluster should [apigroup:config.openshift.io] start all core operators": " [Skipped:Disconnected] [Suite:openshift/conformance/parallel]", - "[sig-arch][Feature:ClusterUpgrade] Cluster should be upgradeable after finishing upgrade [Late][Suite:upgrade]": "", - - "[sig-arch][Feature:ClusterUpgrade] Cluster should be upgradeable before beginning upgrade [Early][Suite:upgrade]": "", - - "[sig-arch][Feature:ClusterUpgrade] Cluster should remain functional during upgrade [Disruptive]": " [Serial]", - "[sig-arch][Late] clients should not use APIs that are removed in upcoming releases [apigroup:apiserver.openshift.io]": " [Suite:openshift/conformance/parallel]", "[sig-arch][Late] operators should not create watch channels very often [apigroup:apiserver.openshift.io]": " [Suite:openshift/conformance/parallel]", @@ -1715,8 +1709,6 @@ var Annotations = map[string]string{ "[sig-cluster-lifecycle] TestAdminAck should succeed [apigroup:config.openshift.io]": " [Suite:openshift/conformance/parallel]", - "[sig-cluster-lifecycle][Feature:DisasterRecovery][Disruptive] [Feature:NodeRecovery] Cluster should survive master and worker failure and recover with machine health checks [apigroup:machine.openshift.io]": " [Serial]", - "[sig-cluster-lifecycle][Feature:Machines] Managed cluster should have machine resources [apigroup:machine.openshift.io]": " [Suite:openshift/conformance/parallel]", "[sig-cluster-lifecycle][Feature:Machines][Disruptive] Managed cluster should recover from deleted worker machines [apigroup:machine.openshift.io]": " [Serial]", @@ -1873,10 +1865,6 @@ var Annotations = map[string]string{ "[sig-etcd] etcd record the start revision of the etcd-operator [Early]": " [Suite:openshift/conformance/parallel]", - "[sig-etcd][Feature:DisasterRecovery][Disruptive][Disabled:Broken] [Feature:EtcdRecovery] Cluster should recover from a backup taken on one node and recovered on another [apigroup:operator.openshift.io]": " [Serial]", - - "[sig-etcd][Feature:DisasterRecovery][Disruptive][Disabled:Broken] [Feature:EtcdRecovery] Cluster should restore itself after quorum loss [apigroup:machine.openshift.io][apigroup:operator.openshift.io]": " [Serial]", - "[sig-etcd][Feature:DisasterRecovery][Suite:openshift/etcd/recovery][Timeout:2h] [Feature:EtcdRecovery][Disruptive] Recover with snapshot with two unhealthy nodes and lost quorum": " [Serial]", "[sig-etcd][Feature:DisasterRecovery][Suite:openshift/etcd/recovery][Timeout:30m] [Feature:EtcdRecovery][Disruptive] Restore snapshot from node on another single unhealthy node": " [Serial]", From 29715ed82ee67179187f48dec1dd9e75f84b227a Mon Sep 17 00:00:00 2001 From: Thomas Jungblut Date: Mon, 15 May 2023 11:40:15 +0200 Subject: [PATCH 2/3] explicitly add upgrade tests back --- test/extended/include.go | 1 + .../util/annotate/generated/zz_generated.annotations.go | 6 ++++++ 2 files changed, 7 insertions(+) diff --git a/test/extended/include.go b/test/extended/include.go index f29dd3d9160b..e7c5b5c35c77 100644 --- a/test/extended/include.go +++ b/test/extended/include.go @@ -7,6 +7,7 @@ import ( _ "k8s.io/kubernetes/openshift-hack/e2e" + _ "github.com/openshift/origin/test/e2e/upgrade" _ "github.com/openshift/origin/test/extended/adminack" _ "github.com/openshift/origin/test/extended/apiserver" _ "github.com/openshift/origin/test/extended/authentication" diff --git a/test/extended/util/annotate/generated/zz_generated.annotations.go b/test/extended/util/annotate/generated/zz_generated.annotations.go index fb8ac93cd5cf..e0c0d36b23c1 100644 --- a/test/extended/util/annotate/generated/zz_generated.annotations.go +++ b/test/extended/util/annotate/generated/zz_generated.annotations.go @@ -701,6 +701,12 @@ var Annotations = map[string]string{ "[sig-arch][Early] Managed cluster should [apigroup:config.openshift.io] start all core operators": " [Skipped:Disconnected] [Suite:openshift/conformance/parallel]", + "[sig-arch][Feature:ClusterUpgrade] Cluster should be upgradeable after finishing upgrade [Late][Suite:upgrade]": "", + + "[sig-arch][Feature:ClusterUpgrade] Cluster should be upgradeable before beginning upgrade [Early][Suite:upgrade]": "", + + "[sig-arch][Feature:ClusterUpgrade] Cluster should remain functional during upgrade [Disruptive]": " [Serial]", + "[sig-arch][Late] clients should not use APIs that are removed in upcoming releases [apigroup:apiserver.openshift.io]": " [Suite:openshift/conformance/parallel]", "[sig-arch][Late] operators should not create watch channels very often [apigroup:apiserver.openshift.io]": " [Suite:openshift/conformance/parallel]", From 57f70062327daa0301468d6629381a5d24397303 Mon Sep 17 00:00:00 2001 From: Thomas Jungblut Date: Tue, 16 May 2023 08:27:45 +0200 Subject: [PATCH 3/3] update go.mod and vendor --- go.mod | 2 +- .../upgrades/network/kube_proxy_migration.go | 231 ------------------ .../test/e2e/upgrades/network/services.go | 133 ---------- vendor/modules.txt | 1 - 4 files changed, 1 insertion(+), 366 deletions(-) delete mode 100644 vendor/k8s.io/kubernetes/test/e2e/upgrades/network/kube_proxy_migration.go delete mode 100644 vendor/k8s.io/kubernetes/test/e2e/upgrades/network/services.go diff --git a/go.mod b/go.mod index 8575d76adeb6..be0acf3a1c13 100644 --- a/go.mod +++ b/go.mod @@ -38,7 +38,6 @@ require ( github.com/spf13/viper v1.8.1 github.com/stretchr/objx v0.5.0 github.com/stretchr/testify v1.8.1 - go.etcd.io/etcd/api/v3 v3.5.7 go.etcd.io/etcd/client/pkg/v3 v3.5.7 go.etcd.io/etcd/client/v3 v3.5.7 golang.org/x/crypto v0.1.0 @@ -219,6 +218,7 @@ require ( github.com/xiang90/probing v0.0.0-20190116061207-43a291ad63a2 // indirect github.com/xlab/treeprint v1.1.0 // indirect go.etcd.io/bbolt v1.3.6 // indirect + go.etcd.io/etcd/api/v3 v3.5.7 // indirect go.etcd.io/etcd/client/v2 v2.305.7 // indirect go.etcd.io/etcd/pkg/v3 v3.5.7 // indirect go.etcd.io/etcd/raft/v3 v3.5.7 // indirect diff --git a/vendor/k8s.io/kubernetes/test/e2e/upgrades/network/kube_proxy_migration.go b/vendor/k8s.io/kubernetes/test/e2e/upgrades/network/kube_proxy_migration.go deleted file mode 100644 index 443c63be03db..000000000000 --- a/vendor/k8s.io/kubernetes/test/e2e/upgrades/network/kube_proxy_migration.go +++ /dev/null @@ -1,231 +0,0 @@ -/* -Copyright 2017 The Kubernetes 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 network - -import ( - "context" - "fmt" - "time" - - appsv1 "k8s.io/api/apps/v1" - v1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/labels" - "k8s.io/apimachinery/pkg/util/wait" - clientset "k8s.io/client-go/kubernetes" - "k8s.io/kubernetes/test/e2e/framework" - e2edaemonset "k8s.io/kubernetes/test/e2e/framework/daemonset" - e2enode "k8s.io/kubernetes/test/e2e/framework/node" - "k8s.io/kubernetes/test/e2e/upgrades" - - "github.com/onsi/ginkgo/v2" -) - -const ( - defaultTestTimeout = time.Duration(5 * time.Minute) - clusterAddonLabelKey = "k8s-app" - clusterComponentKey = "component" - kubeProxyLabelName = "kube-proxy" -) - -// KubeProxyUpgradeTest tests kube-proxy static pods -> DaemonSet upgrade path. -type KubeProxyUpgradeTest struct { -} - -// Name returns the tracking name of the test. -func (KubeProxyUpgradeTest) Name() string { return "[sig-network] kube-proxy-upgrade" } - -// Setup verifies kube-proxy static pods is running before upgrade. -func (t *KubeProxyUpgradeTest) Setup(ctx context.Context, f *framework.Framework) { - ginkgo.By("Waiting for kube-proxy static pods running and ready") - err := waitForKubeProxyStaticPodsRunning(ctx, f.ClientSet) - framework.ExpectNoError(err) -} - -// Test validates if kube-proxy is migrated from static pods to DaemonSet. -func (t *KubeProxyUpgradeTest) Test(ctx context.Context, f *framework.Framework, done <-chan struct{}, upgrade upgrades.UpgradeType) { - c := f.ClientSet - - // Block until upgrade is done. - ginkgo.By("Waiting for upgrade to finish") - <-done - - ginkgo.By("Waiting for kube-proxy static pods disappear") - err := waitForKubeProxyStaticPodsDisappear(ctx, c) - framework.ExpectNoError(err) - - ginkgo.By("Waiting for kube-proxy DaemonSet running and ready") - err = waitForKubeProxyDaemonSetRunning(ctx, f, c) - framework.ExpectNoError(err) -} - -// Teardown does nothing. -func (t *KubeProxyUpgradeTest) Teardown(ctx context.Context, f *framework.Framework) { -} - -// KubeProxyDowngradeTest tests kube-proxy DaemonSet -> static pods downgrade path. -type KubeProxyDowngradeTest struct { -} - -// Name returns the tracking name of the test. -func (KubeProxyDowngradeTest) Name() string { return "[sig-network] kube-proxy-downgrade" } - -// Setup verifies kube-proxy DaemonSet is running before upgrade. -func (t *KubeProxyDowngradeTest) Setup(ctx context.Context, f *framework.Framework) { - ginkgo.By("Waiting for kube-proxy DaemonSet running and ready") - err := waitForKubeProxyDaemonSetRunning(ctx, f, f.ClientSet) - framework.ExpectNoError(err) -} - -// Test validates if kube-proxy is migrated from DaemonSet to static pods. -func (t *KubeProxyDowngradeTest) Test(ctx context.Context, f *framework.Framework, done <-chan struct{}, upgrade upgrades.UpgradeType) { - c := f.ClientSet - - // Block until upgrade is done. - ginkgo.By("Waiting for upgrade to finish") - <-done - - ginkgo.By("Waiting for kube-proxy DaemonSet disappear") - err := waitForKubeProxyDaemonSetDisappear(ctx, c) - framework.ExpectNoError(err) - - ginkgo.By("Waiting for kube-proxy static pods running and ready") - err = waitForKubeProxyStaticPodsRunning(ctx, c) - framework.ExpectNoError(err) -} - -// Teardown does nothing. -func (t *KubeProxyDowngradeTest) Teardown(ctx context.Context, f *framework.Framework) { -} - -func waitForKubeProxyStaticPodsRunning(ctx context.Context, c clientset.Interface) error { - framework.Logf("Waiting up to %v for kube-proxy static pods running", defaultTestTimeout) - - condition := func() (bool, error) { - pods, err := getKubeProxyStaticPods(ctx, c) - if err != nil { - framework.Logf("Failed to get kube-proxy static pods: %v", err) - return false, nil - } - - nodes, err := e2enode.GetReadySchedulableNodes(ctx, c) - if err != nil { - framework.Logf("Failed to get nodes: %v", err) - return false, nil - } - - numberSchedulableNodes := len(nodes.Items) - numberkubeProxyPods := 0 - for _, pod := range pods.Items { - if pod.Status.Phase == v1.PodRunning { - numberkubeProxyPods = numberkubeProxyPods + 1 - } - } - if numberkubeProxyPods != numberSchedulableNodes { - framework.Logf("Expect %v kube-proxy static pods running, got %v running, %v in total", numberSchedulableNodes, numberkubeProxyPods, len(pods.Items)) - return false, nil - } - return true, nil - } - - if err := wait.PollImmediate(5*time.Second, defaultTestTimeout, condition); err != nil { - return fmt.Errorf("error waiting for kube-proxy static pods running: %w", err) - } - return nil -} - -func waitForKubeProxyStaticPodsDisappear(ctx context.Context, c clientset.Interface) error { - framework.Logf("Waiting up to %v for kube-proxy static pods disappear", defaultTestTimeout) - - condition := func() (bool, error) { - pods, err := getKubeProxyStaticPods(ctx, c) - if err != nil { - framework.Logf("Failed to get kube-proxy static pods: %v", err) - return false, nil - } - - if len(pods.Items) != 0 { - framework.Logf("Expect kube-proxy static pods to disappear, got %v pods", len(pods.Items)) - return false, nil - } - return true, nil - } - - if err := wait.PollImmediate(5*time.Second, defaultTestTimeout, condition); err != nil { - return fmt.Errorf("error waiting for kube-proxy static pods disappear: %w", err) - } - return nil -} - -func waitForKubeProxyDaemonSetRunning(ctx context.Context, f *framework.Framework, c clientset.Interface) error { - framework.Logf("Waiting up to %v for kube-proxy DaemonSet running", defaultTestTimeout) - - condition := func() (bool, error) { - daemonSets, err := getKubeProxyDaemonSet(ctx, c) - if err != nil { - framework.Logf("Failed to get kube-proxy DaemonSet: %v", err) - return false, nil - } - - if len(daemonSets.Items) != 1 { - framework.Logf("Expect only one kube-proxy DaemonSet, got %v", len(daemonSets.Items)) - return false, nil - } - - return e2edaemonset.CheckRunningOnAllNodes(ctx, f, &daemonSets.Items[0]) - } - - if err := wait.PollImmediate(5*time.Second, defaultTestTimeout, condition); err != nil { - return fmt.Errorf("error waiting for kube-proxy DaemonSet running: %w", err) - } - return nil -} - -func waitForKubeProxyDaemonSetDisappear(ctx context.Context, c clientset.Interface) error { - framework.Logf("Waiting up to %v for kube-proxy DaemonSet disappear", defaultTestTimeout) - - condition := func() (bool, error) { - daemonSets, err := getKubeProxyDaemonSet(ctx, c) - if err != nil { - framework.Logf("Failed to get kube-proxy DaemonSet: %v", err) - return false, nil - } - - if len(daemonSets.Items) != 0 { - framework.Logf("Expect kube-proxy DaemonSet to disappear, got %v DaemonSet", len(daemonSets.Items)) - return false, nil - } - return true, nil - } - - if err := wait.PollImmediate(5*time.Second, defaultTestTimeout, condition); err != nil { - return fmt.Errorf("error waiting for kube-proxy DaemonSet disappear: %w", err) - } - return nil -} - -func getKubeProxyStaticPods(ctx context.Context, c clientset.Interface) (*v1.PodList, error) { - label := labels.SelectorFromSet(labels.Set(map[string]string{clusterComponentKey: kubeProxyLabelName})) - listOpts := metav1.ListOptions{LabelSelector: label.String()} - return c.CoreV1().Pods(metav1.NamespaceSystem).List(ctx, listOpts) -} - -func getKubeProxyDaemonSet(ctx context.Context, c clientset.Interface) (*appsv1.DaemonSetList, error) { - label := labels.SelectorFromSet(labels.Set(map[string]string{clusterAddonLabelKey: kubeProxyLabelName})) - listOpts := metav1.ListOptions{LabelSelector: label.String()} - return c.AppsV1().DaemonSets(metav1.NamespaceSystem).List(ctx, listOpts) -} diff --git a/vendor/k8s.io/kubernetes/test/e2e/upgrades/network/services.go b/vendor/k8s.io/kubernetes/test/e2e/upgrades/network/services.go deleted file mode 100644 index 83f6d407d325..000000000000 --- a/vendor/k8s.io/kubernetes/test/e2e/upgrades/network/services.go +++ /dev/null @@ -1,133 +0,0 @@ -/* -Copyright 2017 The Kubernetes 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 network - -import ( - "context" - - v1 "k8s.io/api/core/v1" - "k8s.io/apimachinery/pkg/util/wait" - "k8s.io/kubernetes/test/e2e/framework" - e2eservice "k8s.io/kubernetes/test/e2e/framework/service" - "k8s.io/kubernetes/test/e2e/upgrades" - - "github.com/onsi/ginkgo/v2" -) - -// ServiceUpgradeTest tests that a service is available before and -// after a cluster upgrade. During a master-only upgrade, it will test -// that a service remains available during the upgrade. -type ServiceUpgradeTest struct { - jig *e2eservice.TestJig - tcpService *v1.Service - tcpIngressIP string - svcPort int -} - -// Name returns the tracking name of the test. -func (ServiceUpgradeTest) Name() string { return "service-upgrade" } - -func shouldTestPDBs() bool { return true } - -// Setup creates a service with a load balancer and makes sure it's reachable. -func (t *ServiceUpgradeTest) Setup(ctx context.Context, f *framework.Framework) { - serviceName := "service-test" - jig := e2eservice.NewTestJig(f.ClientSet, f.Namespace.Name, serviceName) - - ns := f.Namespace - cs := f.ClientSet - - ginkgo.By("creating a TCP service " + serviceName + " with type=LoadBalancer in namespace " + ns.Name) - _, err := jig.CreateTCPService(ctx, func(s *v1.Service) { - s.Spec.Type = v1.ServiceTypeLoadBalancer - }) - framework.ExpectNoError(err) - tcpService, err := jig.WaitForLoadBalancer(ctx, e2eservice.GetServiceLoadBalancerCreationTimeout(ctx, cs)) - framework.ExpectNoError(err) - - // Get info to hit it with - tcpIngressIP := e2eservice.GetIngressPoint(&tcpService.Status.LoadBalancer.Ingress[0]) - svcPort := int(tcpService.Spec.Ports[0].Port) - - ginkgo.By("creating pod to be part of service " + serviceName) - rc, err := jig.Run(ctx, jig.AddRCAntiAffinity) - framework.ExpectNoError(err) - - if shouldTestPDBs() { - ginkgo.By("creating a PodDisruptionBudget to cover the ReplicationController") - _, err = jig.CreatePDB(ctx, rc) - framework.ExpectNoError(err) - } - - // Hit it once before considering ourselves ready - ginkgo.By("hitting the pod through the service's LoadBalancer") - timeout := e2eservice.LoadBalancerLagTimeoutDefault - if framework.ProviderIs("aws") { - timeout = e2eservice.LoadBalancerLagTimeoutAWS - } - e2eservice.TestReachableHTTP(ctx, tcpIngressIP, svcPort, timeout) - - t.jig = jig - t.tcpService = tcpService - t.tcpIngressIP = tcpIngressIP - t.svcPort = svcPort -} - -// Test runs a connectivity check to the service. -func (t *ServiceUpgradeTest) Test(ctx context.Context, f *framework.Framework, done <-chan struct{}, upgrade upgrades.UpgradeType) { - switch upgrade { - case upgrades.MasterUpgrade, upgrades.ClusterUpgrade: - t.test(ctx, f, done, true, true) - case upgrades.NodeUpgrade: - // Node upgrades should test during disruption only on GCE/GKE for now. - t.test(ctx, f, done, shouldTestPDBs(), false) - default: - t.test(ctx, f, done, false, false) - } -} - -// Teardown cleans up any remaining resources. -func (t *ServiceUpgradeTest) Teardown(ctx context.Context, f *framework.Framework) { - // rely on the namespace deletion to clean up everything -} - -func (t *ServiceUpgradeTest) test(ctx context.Context, f *framework.Framework, done <-chan struct{}, testDuringDisruption, testFinalizer bool) { - if testDuringDisruption { - // Continuous validation - ginkgo.By("continuously hitting the pod through the service's LoadBalancer") - // TODO (pohly): add context support - wait.Until(func() { - e2eservice.TestReachableHTTP(ctx, t.tcpIngressIP, t.svcPort, e2eservice.LoadBalancerLagTimeoutDefault) - }, framework.Poll, done) - } else { - // Block until upgrade is done - ginkgo.By("waiting for upgrade to finish without checking if service remains up") - <-done - } - - // Hit it once more - ginkgo.By("hitting the pod through the service's LoadBalancer") - e2eservice.TestReachableHTTP(ctx, t.tcpIngressIP, t.svcPort, e2eservice.LoadBalancerLagTimeoutDefault) - if testFinalizer { - defer func() { - ginkgo.By("Check that service can be deleted with finalizer") - e2eservice.WaitForServiceDeletedWithFinalizer(ctx, t.jig.Client, t.tcpService.Namespace, t.tcpService.Name) - }() - ginkgo.By("Check that finalizer is present on loadBalancer type service") - e2eservice.WaitForServiceUpdatedWithFinalizer(ctx, t.jig.Client, t.tcpService.Namespace, t.tcpService.Name, true) - } -} diff --git a/vendor/modules.txt b/vendor/modules.txt index 0d8ae27d77f9..f3de6c350997 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -3311,7 +3311,6 @@ k8s.io/kubernetes/test/e2e/storage/vsphere k8s.io/kubernetes/test/e2e/testing-manifests k8s.io/kubernetes/test/e2e/upgrades k8s.io/kubernetes/test/e2e/upgrades/apps -k8s.io/kubernetes/test/e2e/upgrades/network k8s.io/kubernetes/test/e2e/upgrades/node k8s.io/kubernetes/test/fixtures k8s.io/kubernetes/test/integration