From 089491b6a68894c63c09fefd97587489cf9d2f31 Mon Sep 17 00:00:00 2001 From: Sagar Utekar Date: Mon, 6 Jul 2026 22:22:51 +0530 Subject: [PATCH 1/8] test(e2e): add shutdown-on-sigterm failover test Cover the graceful-termination handover deferred from #268: delete a shard primary with the default grace period and assert a replica is promoted before the grace period ends, the shard keeps serving writes, the replaced pod rejoins as a replica, and no keys are lost. Roles are read live from INFO replication rather than ValkeyNode status, which can report stale roles right after cluster formation (see #261). VALKEYCLI_AUTH is unset when running valkey-cli inside the server container, since valkey-cli would otherwise auto-send AUTH as the default user and fail. Closes #270 Signed-off-by: Sagar Utekar --- test/e2e/failover_sigterm_test.go | 203 ++++++++++++++++++++++++++++++ 1 file changed, 203 insertions(+) create mode 100644 test/e2e/failover_sigterm_test.go diff --git a/test/e2e/failover_sigterm_test.go b/test/e2e/failover_sigterm_test.go new file mode 100644 index 00000000..f628e012 --- /dev/null +++ b/test/e2e/failover_sigterm_test.go @@ -0,0 +1,203 @@ +//go:build e2e + +/* +Copyright 2025 Valkey Contributors. + +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 ( + "fmt" + "os" + "os/exec" + "path/filepath" + "strings" + "time" + + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + + "valkey.io/valkey-operator/test/utils" +) + +// keyCount is the number of keys written across the keyspace before the +// failover, and verified afterwards to prove no data was lost. +const keyCount = 50 + +var _ = Describe("Shutdown-on-SIGTERM failover", Label("failover"), func() { + const valkeyName = "failover-sigterm" + + AfterEach(func() { + specReport := CurrentSpecReport() + if specReport.Failed() { + utils.CollectDebugInfo(namespace) + } + + By("cleaning up test resources") + cmd := exec.Command("kubectl", "delete", "valkeycluster", valkeyName, "--ignore-not-found=true") + _, _ = utils.Run(cmd) + }) + + It("promotes a replica before a gracefully terminated primary exits", func() { + By("creating a ValkeyCluster with one replica per shard") + valkeyYaml := fmt.Sprintf(` +apiVersion: valkey.io/v1alpha1 +kind: ValkeyCluster +metadata: + name: %s +spec: + shards: 3 + replicas: 1 +`, valkeyName) + + manifestFile := filepath.Join(os.TempDir(), fmt.Sprintf("%s.yaml", valkeyName)) + err := os.WriteFile(manifestFile, []byte(valkeyYaml), 0644) + Expect(err).NotTo(HaveOccurred(), "Failed to write manifest file") + defer func() { + Expect(os.Remove(manifestFile)).To(Succeed()) + }() + + cmd := exec.Command("kubectl", "create", "-f", manifestFile) + _, err = utils.Run(cmd) + Expect(err).NotTo(HaveOccurred(), "Failed to create ValkeyCluster CR") + + By("waiting for the cluster to become Ready") + Eventually(func(g Gomega) { + cmd := exec.Command("kubectl", "get", "valkeycluster", valkeyName, + "-o", "jsonpath={.status.conditions[?(@.type==\"Ready\")].status}") + output, err := utils.Run(cmd) + g.Expect(err).NotTo(HaveOccurred(), "Failed to get ValkeyCluster Ready condition") + g.Expect(output).To(Equal("True"), "ValkeyCluster is not Ready") + }).WithTimeout(4 * time.Minute).Should(Succeed()) + + By("identifying the primary and replica of shard 0") + var primaryPod, replicaPod string + Eventually(func(g Gomega) { + primaryPod, replicaPod = getShardRoles(g, valkeyName, 0) + g.Expect(primaryPod).NotTo(BeEmpty(), "shard 0 has no primary") + g.Expect(replicaPod).NotTo(BeEmpty(), "shard 0 has no in-sync replica") + }).WithTimeout(3 * time.Minute).Should(Succeed()) + + By("recording the primary pod UID so its replacement can be detected") + cmd = exec.Command("kubectl", "get", "pod", primaryPod, "-o", "jsonpath={.metadata.uid}") + oldPrimaryUID, err := utils.Run(cmd) + Expect(err).NotTo(HaveOccurred(), "Failed to get primary pod UID") + Expect(oldPrimaryUID).NotTo(BeEmpty()) + + By("writing keys across the keyspace") + script := fmt.Sprintf( + "ok=0; for i in $(seq 1 %d); do "+ + "r=$(valkey-cli -c set e2e:failover:$i v$i); "+ + "[ \"$r\" = \"OK\" ] && ok=$((ok+1)); "+ + "done; echo written=$ok", keyCount) + output, err := execValkeyPodShell(primaryPod, script) + Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to write keys: %s", output)) + Expect(output).To(ContainSubstring(fmt.Sprintf("written=%d", keyCount)), + fmt.Sprintf("Not all keys were written: %s", output)) + + By("gracefully terminating the primary pod (SIGTERM with default grace period)") + cmd = exec.Command("kubectl", "delete", "pod", primaryPod, "--wait=false") + _, err = utils.Run(cmd) + Expect(err).NotTo(HaveOccurred(), "Failed to delete primary pod") + + // The handover must complete inside terminationGracePeriodSeconds + // (default 30s), so the replica has to report role:master within that + // window; otherwise SIGKILL would have cut the handoff short. + By("asserting the replica is promoted within the grace period") + Eventually(func(g Gomega) { + output, err := execValkeyPodShell(replicaPod, "valkey-cli INFO replication") + g.Expect(err).NotTo(HaveOccurred(), "Failed to get replication info from replica") + g.Expect(output).To(ContainSubstring("role:master"), + "replica was not promoted to primary") + }).WithTimeout(30 * time.Second).WithPolling(time.Second).Should(Succeed()) + + By("asserting the shard keeps serving writes through the disruption") + Eventually(func(g Gomega) { + output, err := execValkeyPodShell(replicaPod, "valkey-cli -c set e2e:failover:post v-post") + g.Expect(err).NotTo(HaveOccurred(), "Failed to write through promoted primary") + g.Expect(output).To(ContainSubstring("OK")) + }).Should(Succeed()) + + By("waiting for the replaced pod to come back and rejoin as replica") + Eventually(func(g Gomega) { + cmd := exec.Command("kubectl", "get", "pod", primaryPod, "-o", "jsonpath={.metadata.uid}") + uid, err := utils.Run(cmd) + g.Expect(err).NotTo(HaveOccurred(), "replacement pod does not exist yet") + g.Expect(uid).NotTo(Equal(oldPrimaryUID), "old pod is still terminating") + + output, err := execValkeyPodShell(primaryPod, "valkey-cli INFO replication") + g.Expect(err).NotTo(HaveOccurred(), "Failed to get replication info from replaced pod") + g.Expect(output).To(ContainSubstring("role:slave"), + "replaced pod did not rejoin the shard as a replica") + }).WithTimeout(4 * time.Minute).Should(Succeed()) + + By("asserting the cluster reports a healthy state") + Eventually(func(g Gomega) { + output, err := execValkeyPodShell(replicaPod, "valkey-cli CLUSTER INFO") + g.Expect(err).NotTo(HaveOccurred(), "Failed to get cluster info") + g.Expect(output).To(ContainSubstring("cluster_state:ok")) + }).Should(Succeed()) + + By("asserting no keys were lost across the handoff") + script = fmt.Sprintf( + "ok=0; for i in $(seq 1 %d); do "+ + "v=$(valkey-cli -c get e2e:failover:$i); "+ + "[ \"$v\" = \"v$i\" ] && ok=$((ok+1)); "+ + "done; echo readable=$ok", keyCount) + Eventually(func(g Gomega) { + output, err := execValkeyPodShell(replicaPod, script) + g.Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to read keys back: %s", output)) + g.Expect(output).To(ContainSubstring(fmt.Sprintf("readable=%d", keyCount)), + fmt.Sprintf("Some keys were lost across the failover: %s", output)) + }).Should(Succeed()) + }) +}) + +// execValkeyPodShell runs a shell script inside the pod's server container +// with VALKEYCLI_AUTH unset. The operator injects that variable for the probe +// scripts, and valkey-cli would otherwise automatically send AUTH as the +// default user — which fails and pollutes the command output. Commands run as +// the default (nopass) user, like the rest of the e2e suite. +func execValkeyPodShell(pod string, script string) (string, error) { + cmd := exec.Command("kubectl", "exec", pod, "-c", "server", "--", + "sh", "-c", "unset VALKEYCLI_AUTH; "+script) + return utils.Run(cmd) +} + +// getShardRoles returns the pod names of the primary and replica of the given +// shard. Roles are read live from INFO replication on each pod rather than +// from ValkeyNode status, which can lag behind role changes (see #261). The +// replica is only accepted once its replication link is up, so the failover +// is not attempted against a still-syncing replica. +func getShardRoles(g Gomega, clusterName string, shardIndex int) (primaryPod, replicaPod string) { + cmd := exec.Command("kubectl", "get", "pods", + "-l", fmt.Sprintf("valkey.io/cluster=%s,valkey.io/shard-index=%d", clusterName, shardIndex), + "-o", "go-template={{ range .items }}{{ .metadata.name }}{{ \"\\n\" }}{{ end }}") + output, err := utils.Run(cmd) + g.Expect(err).NotTo(HaveOccurred(), "Failed to list shard pods") + + for _, pod := range utils.GetNonEmptyLines(output) { + pod = strings.TrimSpace(pod) + info, err := execValkeyPodShell(pod, "valkey-cli INFO replication") + g.Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to get replication info from %s", pod)) + switch { + case strings.Contains(info, "role:master"): + primaryPod = pod + case strings.Contains(info, "role:slave") && strings.Contains(info, "master_link_status:up"): + replicaPod = pod + } + } + return primaryPod, replicaPod +} From dd5215dcce20ade06d9f182b3f0a473bba96e62c Mon Sep 17 00:00:00 2001 From: Sagar Utekar Date: Wed, 8 Jul 2026 19:12:47 +0530 Subject: [PATCH 2/8] test(e2e): fold StatefulSet-deletion recovery into failover suite Merge the primary-workload-deletion recovery scenario suggested in review into the failover suite as a second spec: delete the shard primary's StatefulSet, assert the replica is promoted, the operator recreates the StatefulSet, the recreated pod rejoins as a replica, and the ValkeyCluster returns to Ready with no keys lost. Cluster creation, key seeding, and health/data assertions are shared between both specs. Signed-off-by: Sagar Utekar --- test/e2e/failover_sigterm_test.go | 194 +++++++++++++++++++++--------- 1 file changed, 137 insertions(+), 57 deletions(-) diff --git a/test/e2e/failover_sigterm_test.go b/test/e2e/failover_sigterm_test.go index f628e012..4c4bed6e 100644 --- a/test/e2e/failover_sigterm_test.go +++ b/test/e2e/failover_sigterm_test.go @@ -37,7 +37,7 @@ import ( const keyCount = 50 var _ = Describe("Shutdown-on-SIGTERM failover", Label("failover"), func() { - const valkeyName = "failover-sigterm" + var valkeyName string AfterEach(func() { specReport := CurrentSpecReport() @@ -51,36 +51,8 @@ var _ = Describe("Shutdown-on-SIGTERM failover", Label("failover"), func() { }) It("promotes a replica before a gracefully terminated primary exits", func() { - By("creating a ValkeyCluster with one replica per shard") - valkeyYaml := fmt.Sprintf(` -apiVersion: valkey.io/v1alpha1 -kind: ValkeyCluster -metadata: - name: %s -spec: - shards: 3 - replicas: 1 -`, valkeyName) - - manifestFile := filepath.Join(os.TempDir(), fmt.Sprintf("%s.yaml", valkeyName)) - err := os.WriteFile(manifestFile, []byte(valkeyYaml), 0644) - Expect(err).NotTo(HaveOccurred(), "Failed to write manifest file") - defer func() { - Expect(os.Remove(manifestFile)).To(Succeed()) - }() - - cmd := exec.Command("kubectl", "create", "-f", manifestFile) - _, err = utils.Run(cmd) - Expect(err).NotTo(HaveOccurred(), "Failed to create ValkeyCluster CR") - - By("waiting for the cluster to become Ready") - Eventually(func(g Gomega) { - cmd := exec.Command("kubectl", "get", "valkeycluster", valkeyName, - "-o", "jsonpath={.status.conditions[?(@.type==\"Ready\")].status}") - output, err := utils.Run(cmd) - g.Expect(err).NotTo(HaveOccurred(), "Failed to get ValkeyCluster Ready condition") - g.Expect(output).To(Equal("True"), "ValkeyCluster is not Ready") - }).WithTimeout(4 * time.Minute).Should(Succeed()) + valkeyName = "failover-sigterm" + createFailoverCluster(valkeyName) By("identifying the primary and replica of shard 0") var primaryPod, replicaPod string @@ -91,21 +63,13 @@ spec: }).WithTimeout(3 * time.Minute).Should(Succeed()) By("recording the primary pod UID so its replacement can be detected") - cmd = exec.Command("kubectl", "get", "pod", primaryPod, "-o", "jsonpath={.metadata.uid}") + cmd := exec.Command("kubectl", "get", "pod", primaryPod, "-o", "jsonpath={.metadata.uid}") oldPrimaryUID, err := utils.Run(cmd) Expect(err).NotTo(HaveOccurred(), "Failed to get primary pod UID") Expect(oldPrimaryUID).NotTo(BeEmpty()) By("writing keys across the keyspace") - script := fmt.Sprintf( - "ok=0; for i in $(seq 1 %d); do "+ - "r=$(valkey-cli -c set e2e:failover:$i v$i); "+ - "[ \"$r\" = \"OK\" ] && ok=$((ok+1)); "+ - "done; echo written=$ok", keyCount) - output, err := execValkeyPodShell(primaryPod, script) - Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to write keys: %s", output)) - Expect(output).To(ContainSubstring(fmt.Sprintf("written=%d", keyCount)), - fmt.Sprintf("Not all keys were written: %s", output)) + writeTestKeys(primaryPod) By("gracefully terminating the primary pod (SIGTERM with default grace period)") cmd = exec.Command("kubectl", "delete", "pod", primaryPod, "--wait=false") @@ -143,28 +107,144 @@ spec: "replaced pod did not rejoin the shard as a replica") }).WithTimeout(4 * time.Minute).Should(Succeed()) - By("asserting the cluster reports a healthy state") + verifyClusterHealthyAndKeysIntact(valkeyName, replicaPod) + }) + + It("recovers when the primary's StatefulSet is deleted", func() { + valkeyName = "failover-sts-delete" + createFailoverCluster(valkeyName) + + By("identifying the primary and replica of shard 0") + var primaryPod, replicaPod string Eventually(func(g Gomega) { - output, err := execValkeyPodShell(replicaPod, "valkey-cli CLUSTER INFO") - g.Expect(err).NotTo(HaveOccurred(), "Failed to get cluster info") - g.Expect(output).To(ContainSubstring("cluster_state:ok")) - }).Should(Succeed()) + primaryPod, replicaPod = getShardRoles(g, valkeyName, 0) + g.Expect(primaryPod).NotTo(BeEmpty(), "shard 0 has no primary") + g.Expect(replicaPod).NotTo(BeEmpty(), "shard 0 has no in-sync replica") + }).WithTimeout(3 * time.Minute).Should(Succeed()) - By("asserting no keys were lost across the handoff") - script = fmt.Sprintf( - "ok=0; for i in $(seq 1 %d); do "+ - "v=$(valkey-cli -c get e2e:failover:$i); "+ - "[ \"$v\" = \"v$i\" ] && ok=$((ok+1)); "+ - "done; echo readable=$ok", keyCount) + By("writing keys across the keyspace") + writeTestKeys(primaryPod) + + By("deleting the primary's StatefulSet to deschedule the whole workload") + // Pods are named -0. + primarySts := strings.TrimSuffix(primaryPod, "-0") + cmd := exec.Command("kubectl", "delete", "statefulset", primarySts, "--wait=false") + _, err := utils.Run(cmd) + Expect(err).NotTo(HaveOccurred(), "Failed to delete primary StatefulSet") + + By("asserting the replica is promoted") Eventually(func(g Gomega) { - output, err := execValkeyPodShell(replicaPod, script) - g.Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to read keys back: %s", output)) - g.Expect(output).To(ContainSubstring(fmt.Sprintf("readable=%d", keyCount)), - fmt.Sprintf("Some keys were lost across the failover: %s", output)) - }).Should(Succeed()) + output, err := execValkeyPodShell(replicaPod, "valkey-cli INFO replication") + g.Expect(err).NotTo(HaveOccurred(), "Failed to get replication info from replica") + g.Expect(output).To(ContainSubstring("role:master"), + "replica was not promoted to primary") + }).WithTimeout(time.Minute).WithPolling(time.Second).Should(Succeed()) + + By("waiting for the operator to recreate the StatefulSet and the pod to rejoin as replica") + Eventually(func(g Gomega) { + cmd := exec.Command("kubectl", "get", "statefulset", primarySts, "-o", "jsonpath={.status.readyReplicas}") + ready, err := utils.Run(cmd) + g.Expect(err).NotTo(HaveOccurred(), "recreated StatefulSet does not exist yet") + g.Expect(ready).To(Equal("1"), "recreated StatefulSet has no ready replicas") + + output, err := execValkeyPodShell(primaryPod, "valkey-cli INFO replication") + g.Expect(err).NotTo(HaveOccurred(), "Failed to get replication info from recreated pod") + g.Expect(output).To(ContainSubstring("role:slave"), + "recreated pod did not rejoin the shard as a replica") + }).WithTimeout(4 * time.Minute).Should(Succeed()) + + verifyClusterHealthyAndKeysIntact(valkeyName, replicaPod) }) }) +// createFailoverCluster creates a ValkeyCluster with one replica per shard and +// waits for it to become Ready. +func createFailoverCluster(valkeyName string) { + GinkgoHelper() + + By("creating a ValkeyCluster with one replica per shard") + valkeyYaml := fmt.Sprintf(` +apiVersion: valkey.io/v1alpha1 +kind: ValkeyCluster +metadata: + name: %s +spec: + shards: 3 + replicas: 1 +`, valkeyName) + + manifestFile := filepath.Join(os.TempDir(), fmt.Sprintf("%s.yaml", valkeyName)) + err := os.WriteFile(manifestFile, []byte(valkeyYaml), 0644) + Expect(err).NotTo(HaveOccurred(), "Failed to write manifest file") + defer func() { + Expect(os.Remove(manifestFile)).To(Succeed()) + }() + + cmd := exec.Command("kubectl", "create", "-f", manifestFile) + _, err = utils.Run(cmd) + Expect(err).NotTo(HaveOccurred(), "Failed to create ValkeyCluster CR") + + By("waiting for the cluster to become Ready") + Eventually(func(g Gomega) { + cmd := exec.Command("kubectl", "get", "valkeycluster", valkeyName, + "-o", "jsonpath={.status.conditions[?(@.type==\"Ready\")].status}") + output, err := utils.Run(cmd) + g.Expect(err).NotTo(HaveOccurred(), "Failed to get ValkeyCluster Ready condition") + g.Expect(output).To(Equal("True"), "ValkeyCluster is not Ready") + }).WithTimeout(4 * time.Minute).Should(Succeed()) +} + +// writeTestKeys writes keyCount keys across the keyspace via the given pod. +func writeTestKeys(pod string) { + GinkgoHelper() + + script := fmt.Sprintf( + "ok=0; for i in $(seq 1 %d); do "+ + "r=$(valkey-cli -c set e2e:failover:$i v$i); "+ + "[ \"$r\" = \"OK\" ] && ok=$((ok+1)); "+ + "done; echo written=$ok", keyCount) + output, err := execValkeyPodShell(pod, script) + Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to write keys: %s", output)) + Expect(output).To(ContainSubstring(fmt.Sprintf("written=%d", keyCount)), + fmt.Sprintf("Not all keys were written: %s", output)) +} + +// verifyClusterHealthyAndKeysIntact asserts the ValkeyCluster returns to +// Ready, the cluster reports a healthy state, and every test key survived the +// disruption. +func verifyClusterHealthyAndKeysIntact(valkeyName, pod string) { + GinkgoHelper() + + By("asserting the ValkeyCluster returns to Ready") + Eventually(func(g Gomega) { + cmd := exec.Command("kubectl", "get", "valkeycluster", valkeyName, + "-o", "jsonpath={.status.conditions[?(@.type==\"Ready\")].status}") + output, err := utils.Run(cmd) + g.Expect(err).NotTo(HaveOccurred(), "Failed to get ValkeyCluster Ready condition") + g.Expect(output).To(Equal("True"), "ValkeyCluster did not return to Ready") + }).Should(Succeed()) + + By("asserting the cluster reports a healthy state") + Eventually(func(g Gomega) { + output, err := execValkeyPodShell(pod, "valkey-cli CLUSTER INFO") + g.Expect(err).NotTo(HaveOccurred(), "Failed to get cluster info") + g.Expect(output).To(ContainSubstring("cluster_state:ok")) + }).Should(Succeed()) + + By("asserting no keys were lost across the disruption") + script := fmt.Sprintf( + "ok=0; for i in $(seq 1 %d); do "+ + "v=$(valkey-cli -c get e2e:failover:$i); "+ + "[ \"$v\" = \"v$i\" ] && ok=$((ok+1)); "+ + "done; echo readable=$ok", keyCount) + Eventually(func(g Gomega) { + output, err := execValkeyPodShell(pod, script) + g.Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to read keys back: %s", output)) + g.Expect(output).To(ContainSubstring(fmt.Sprintf("readable=%d", keyCount)), + fmt.Sprintf("Some keys were lost across the disruption: %s", output)) + }).Should(Succeed()) +} + // execValkeyPodShell runs a shell script inside the pod's server container // with VALKEYCLI_AUTH unset. The operator injects that variable for the probe // scripts, and valkey-cli would otherwise automatically send AUTH as the From 4dcc85614c5bb346100fe89c2e14efd627b4aaa4 Mon Sep 17 00:00:00 2001 From: Sagar Utekar Date: Wed, 8 Jul 2026 19:43:56 +0530 Subject: [PATCH 3/8] test(e2e): authenticate valkey-cli as the default user Give the failover test clusters a password-protected default user via a passwordSecret, following the pattern from #292, and set VALKEYCLI_AUTH to that password for every valkey-cli invocation instead of unsetting it. The tests no longer rely on the default user being passwordless. Signed-off-by: Sagar Utekar --- test/e2e/failover_sigterm_test.go | 41 ++++++++++++++++++++++++------- 1 file changed, 32 insertions(+), 9 deletions(-) diff --git a/test/e2e/failover_sigterm_test.go b/test/e2e/failover_sigterm_test.go index 4c4bed6e..ad27641e 100644 --- a/test/e2e/failover_sigterm_test.go +++ b/test/e2e/failover_sigterm_test.go @@ -19,6 +19,7 @@ limitations under the License. package e2e import ( + "encoding/base64" "fmt" "os" "os/exec" @@ -36,6 +37,11 @@ import ( // failover, and verified afterwards to prove no data was lost. const keyCount = 50 +// failoverDefaultPassword is the password configured for the default user of +// the failover test clusters (via a passwordSecret, following the pattern +// introduced in #292), so valkey-cli commands run authenticated. +const failoverDefaultPassword = "e2eFailoverPassw0rd" + var _ = Describe("Shutdown-on-SIGTERM failover", Label("failover"), func() { var valkeyName string @@ -48,6 +54,8 @@ var _ = Describe("Shutdown-on-SIGTERM failover", Label("failover"), func() { By("cleaning up test resources") cmd := exec.Command("kubectl", "delete", "valkeycluster", valkeyName, "--ignore-not-found=true") _, _ = utils.Run(cmd) + cmd = exec.Command("kubectl", "delete", "secret", valkeyName+"-users", "--ignore-not-found=true") + _, _ = utils.Run(cmd) }) It("promotes a replica before a gracefully terminated primary exits", func() { @@ -157,21 +165,35 @@ var _ = Describe("Shutdown-on-SIGTERM failover", Label("failover"), func() { }) }) -// createFailoverCluster creates a ValkeyCluster with one replica per shard and -// waits for it to become Ready. +// createFailoverCluster creates a ValkeyCluster with one replica per shard +// and a password-protected default user, and waits for it to become Ready. func createFailoverCluster(valkeyName string) { GinkgoHelper() By("creating a ValkeyCluster with one replica per shard") valkeyYaml := fmt.Sprintf(` +apiVersion: v1 +kind: Secret +metadata: + name: %[1]s-users +data: + defaultpw: %[2]s +--- apiVersion: valkey.io/v1alpha1 kind: ValkeyCluster metadata: - name: %s + name: %[1]s spec: shards: 3 replicas: 1 -`, valkeyName) + users: + - name: default + enabled: true + permissions: "+@all ~* &*" + passwordSecret: + name: %[1]s-users + keys: [defaultpw] +`, valkeyName, base64.StdEncoding.EncodeToString([]byte(failoverDefaultPassword))) manifestFile := filepath.Join(os.TempDir(), fmt.Sprintf("%s.yaml", valkeyName)) err := os.WriteFile(manifestFile, []byte(valkeyYaml), 0644) @@ -246,13 +268,14 @@ func verifyClusterHealthyAndKeysIntact(valkeyName, pod string) { } // execValkeyPodShell runs a shell script inside the pod's server container -// with VALKEYCLI_AUTH unset. The operator injects that variable for the probe -// scripts, and valkey-cli would otherwise automatically send AUTH as the -// default user — which fails and pollutes the command output. Commands run as -// the default (nopass) user, like the rest of the e2e suite. +// with VALKEYCLI_AUTH set to the default user's password, so every valkey-cli +// invocation in the script runs authenticated. The operator injects +// VALKEYCLI_AUTH with the _operator user's password for the probe scripts; +// it must be overridden here because valkey-cli auto-sends it as the default +// user's AUTH credential. func execValkeyPodShell(pod string, script string) (string, error) { cmd := exec.Command("kubectl", "exec", pod, "-c", "server", "--", - "sh", "-c", "unset VALKEYCLI_AUTH; "+script) + "sh", "-c", fmt.Sprintf("export VALKEYCLI_AUTH=%q; ", failoverDefaultPassword)+script) return utils.Run(cmd) } From 8271054db9d27106f49cfc039ef5898ac964bd9b Mon Sep 17 00:00:00 2001 From: Sagar Utekar Date: Wed, 8 Jul 2026 20:17:51 +0530 Subject: [PATCH 4/8] test(e2e): drop spec duplicating existing primary-recovery coverage The StatefulSet-deletion recovery scenario is already covered in main by 'should detect and recover when a primary deployment is deleted' in valkeycluster_test.go. Keep this suite focused on what that test does not exercise: the graceful shutdown-on-sigterm handoff of a terminating primary pod. Signed-off-by: Sagar Utekar --- test/e2e/failover_sigterm_test.go | 45 ------------------------------- 1 file changed, 45 deletions(-) diff --git a/test/e2e/failover_sigterm_test.go b/test/e2e/failover_sigterm_test.go index ad27641e..2752f600 100644 --- a/test/e2e/failover_sigterm_test.go +++ b/test/e2e/failover_sigterm_test.go @@ -118,51 +118,6 @@ var _ = Describe("Shutdown-on-SIGTERM failover", Label("failover"), func() { verifyClusterHealthyAndKeysIntact(valkeyName, replicaPod) }) - It("recovers when the primary's StatefulSet is deleted", func() { - valkeyName = "failover-sts-delete" - createFailoverCluster(valkeyName) - - By("identifying the primary and replica of shard 0") - var primaryPod, replicaPod string - Eventually(func(g Gomega) { - primaryPod, replicaPod = getShardRoles(g, valkeyName, 0) - g.Expect(primaryPod).NotTo(BeEmpty(), "shard 0 has no primary") - g.Expect(replicaPod).NotTo(BeEmpty(), "shard 0 has no in-sync replica") - }).WithTimeout(3 * time.Minute).Should(Succeed()) - - By("writing keys across the keyspace") - writeTestKeys(primaryPod) - - By("deleting the primary's StatefulSet to deschedule the whole workload") - // Pods are named -0. - primarySts := strings.TrimSuffix(primaryPod, "-0") - cmd := exec.Command("kubectl", "delete", "statefulset", primarySts, "--wait=false") - _, err := utils.Run(cmd) - Expect(err).NotTo(HaveOccurred(), "Failed to delete primary StatefulSet") - - By("asserting the replica is promoted") - Eventually(func(g Gomega) { - output, err := execValkeyPodShell(replicaPod, "valkey-cli INFO replication") - g.Expect(err).NotTo(HaveOccurred(), "Failed to get replication info from replica") - g.Expect(output).To(ContainSubstring("role:master"), - "replica was not promoted to primary") - }).WithTimeout(time.Minute).WithPolling(time.Second).Should(Succeed()) - - By("waiting for the operator to recreate the StatefulSet and the pod to rejoin as replica") - Eventually(func(g Gomega) { - cmd := exec.Command("kubectl", "get", "statefulset", primarySts, "-o", "jsonpath={.status.readyReplicas}") - ready, err := utils.Run(cmd) - g.Expect(err).NotTo(HaveOccurred(), "recreated StatefulSet does not exist yet") - g.Expect(ready).To(Equal("1"), "recreated StatefulSet has no ready replicas") - - output, err := execValkeyPodShell(primaryPod, "valkey-cli INFO replication") - g.Expect(err).NotTo(HaveOccurred(), "Failed to get replication info from recreated pod") - g.Expect(output).To(ContainSubstring("role:slave"), - "recreated pod did not rejoin the shard as a replica") - }).WithTimeout(4 * time.Minute).Should(Succeed()) - - verifyClusterHealthyAndKeysIntact(valkeyName, replicaPod) - }) }) // createFailoverCluster creates a ValkeyCluster with one replica per shard From e9bb6a700f3f6cc400136c204ffc85b5be2f6cc6 Mon Sep 17 00:00:00 2001 From: Sagar Utekar Date: Thu, 9 Jul 2026 08:06:35 +0530 Subject: [PATCH 5/8] test(e2e): write continuously through the handoff, verify no acked loss Strengthen the sigterm failover spec per review: a continuous writer runs through the termination recording per-attempt acks, and the test asserts every acknowledged write is readable afterwards and that the longest gap between acknowledged writes stays bounded. The failover path (coordinated handoff vs failure detection) is detected from the promoted replica's log and reported. The handoff path is reported rather than hard-asserted: valkey 9.0's clusterAutoFailoverOnShutdown requires exact ack-offset equality when selecting a replica and intermittently falls back to failure-detection promotion (~1 in 3 under write load in local testing), so a hard assertion would flake until that promotion is deterministic. Observed on Kind with the handoff engaged: 7490 acknowledged writes through the disruption, longest writer gap 0.05s, zero lost. Signed-off-by: Sagar Utekar --- test/e2e/failover_sigterm_test.go | 150 ++++++++++++++++++++++++++++-- 1 file changed, 144 insertions(+), 6 deletions(-) diff --git a/test/e2e/failover_sigterm_test.go b/test/e2e/failover_sigterm_test.go index 2752f600..7ed34986 100644 --- a/test/e2e/failover_sigterm_test.go +++ b/test/e2e/failover_sigterm_test.go @@ -79,6 +79,9 @@ var _ = Describe("Shutdown-on-SIGTERM failover", Label("failover"), func() { By("writing keys across the keyspace") writeTestKeys(primaryPod) + By("starting a continuous writer to run through the disruption") + writer := startContinuousWriter(replicaPod) + By("gracefully terminating the primary pod (SIGTERM with default grace period)") cmd = exec.Command("kubectl", "delete", "pod", primaryPod, "--wait=false") _, err = utils.Run(cmd) @@ -95,12 +98,40 @@ var _ = Describe("Shutdown-on-SIGTERM failover", Label("failover"), func() { "replica was not promoted to primary") }).WithTimeout(30 * time.Second).WithPolling(time.Second).Should(Succeed()) - By("asserting the shard keeps serving writes through the disruption") - Eventually(func(g Gomega) { - output, err := execValkeyPodShell(replicaPod, "valkey-cli -c set e2e:failover:post v-post") - g.Expect(err).NotTo(HaveOccurred(), "Failed to write through promoted primary") - g.Expect(output).To(ContainSubstring("OK")) - }).Should(Succeed()) + // The graceful handoff is a coordinated failover driven by the + // terminating primary (CLUSTER FAILOVER FORCE REPLICAID), which logs + // "Forced failover primary request accepted" on the promoted replica. + // The crash path goes through FAIL detection and a rank-based election + // instead. valkey 9.0's replica selection on shutdown is best-effort + // (it requires exact ack-offset equality at the shutdown instant and + // intermittently falls back to the crash path, ~1 in 3 under write + // load in local testing), so the handoff path is reported rather than + // hard-asserted until the promotion is made deterministic. + By("reporting which failover path drove the promotion") + cmd = exec.Command("kubectl", "logs", replicaPod, "-c", "server") + logs, err := utils.Run(cmd) + Expect(err).NotTo(HaveOccurred(), "Failed to get promoted replica logs") + if strings.Contains(strings.ToLower(logs), "forced failover primary request accepted") { + _, _ = fmt.Fprintf(GinkgoWriter, "promotion path: coordinated shutdown handoff (forced failover)\n") + } else { + _, _ = fmt.Fprintf(GinkgoWriter, + "promotion path: failure detection (shutdown handoff did not engage; see 'Unable to find a replica' on the primary)\n") + } + + By("asserting every write acknowledged during the disruption is readable") + acked, maxGap := writer.stop() + Expect(acked).NotTo(BeEmpty(), "continuous writer recorded no acknowledged writes") + _, _ = fmt.Fprintf(GinkgoWriter, + "continuous writer: %d acknowledged writes, longest gap between acknowledged writes: %.2fs\n", + len(acked), maxGap) + verifyAcknowledgedWrites(replicaPod, acked) + + // The orderly handoff keeps a writer available throughout the + // disruption; a crash-style failover leaves the shard writer-less for + // failure detection plus an election. The bound is deliberately + // generous to absorb CI noise. + Expect(maxGap).To(BeNumerically("<", 10.0), + fmt.Sprintf("shard had no writer for %.2fs during the graceful handoff", maxGap)) By("waiting for the replaced pod to come back and rejoin as replica") Eventually(func(g Gomega) { @@ -222,6 +253,113 @@ func verifyClusterHealthyAndKeysIntact(valkeyName, pod string) { }).Should(Succeed()) } +// continuousWriter tracks a background write loop running inside a pod. +type continuousWriter struct { + cmd *exec.Cmd + output *strings.Builder +} + +// startContinuousWriter starts a background loop inside the pod that writes +// uniquely-numbered keys through the cluster for writerDurationSeconds, +// recording a timestamped ack/fail line per attempt. It is started against +// the shard's replica so cluster-mode redirects follow the primary across the +// handoff. +const writerDurationSeconds = 20 + +func startContinuousWriter(pod string) *continuousWriter { + GinkgoHelper() + + // POSIX-sh only: the image's /bin/sh is dash, which has no $SECONDS. + script := fmt.Sprintf( + "end=$(($(date +%%s)+%d)); i=0; "+ + "while [ \"$(date +%%s)\" -lt \"$end\" ]; do "+ + "r=$(valkey-cli -c set e2e:cw:$i v$i 2>/dev/null | tail -n 1); "+ + "if [ \"$r\" = \"OK\" ]; then echo \"ack $i $(date +%%s.%%N)\"; "+ + "else echo \"fail $i $(date +%%s.%%N)\"; fi; "+ + "i=$((i+1)); "+ + "done", writerDurationSeconds) + cmd := exec.Command("kubectl", "exec", pod, "-c", "server", "--", + "sh", "-c", fmt.Sprintf("export VALKEYCLI_AUTH=%q; ", failoverDefaultPassword)+script) + + w := &continuousWriter{cmd: cmd, output: &strings.Builder{}} + cmd.Stdout = w.output + cmd.Stderr = w.output + Expect(cmd.Start()).To(Succeed(), "Failed to start continuous writer") + return w +} + +// stop waits for the writer loop to finish and returns the map of +// acknowledged key indices to their expected values, plus the longest gap in +// seconds between consecutive acknowledged writes (the shard's effective +// write-unavailability window). +func (w *continuousWriter) stop() (acked map[string]string, maxGap float64) { + GinkgoHelper() + + Expect(w.cmd.Wait()).To(Succeed(), "continuous writer failed: %s", w.output.String()) + + acked = map[string]string{} + lastAckTime := -1.0 + for _, line := range utils.GetNonEmptyLines(w.output.String()) { + fields := strings.Fields(line) + if len(fields) != 3 { + continue + } + status, idx := fields[0], fields[1] + var ts float64 + if _, err := fmt.Sscanf(fields[2], "%f", &ts); err != nil { + continue + } + if status != "ack" { + continue + } + if lastAckTime >= 0 && ts-lastAckTime > maxGap { + maxGap = ts - lastAckTime + } + lastAckTime = ts + acked[idx] = "v" + idx + } + return acked, maxGap +} + +// verifyAcknowledgedWrites asserts that every key the continuous writer got +// an OK for is readable with the expected value — i.e. no acknowledged write +// was dropped during the handoff. +func verifyAcknowledgedWrites(pod string, acked map[string]string) { + GinkgoHelper() + + maxIdx := 0 + for idx := range acked { + var i int + _, err := fmt.Sscanf(idx, "%d", &i) + Expect(err).NotTo(HaveOccurred()) + if i > maxIdx { + maxIdx = i + } + } + + script := fmt.Sprintf( + "for i in $(seq 0 %d); do v=$(valkey-cli -c get e2e:cw:$i 2>/dev/null | tail -n 1); echo \"$i $v\"; done", maxIdx) + output, err := execValkeyPodShell(pod, script) + Expect(err).NotTo(HaveOccurred(), "Failed to read back acknowledged writes") + + readable := map[string]string{} + for _, line := range utils.GetNonEmptyLines(output) { + fields := strings.Fields(line) + if len(fields) == 2 { + readable[fields[0]] = fields[1] + } + } + + var lost []string + for idx, want := range acked { + if readable[idx] != want { + lost = append(lost, idx) + } + } + Expect(lost).To(BeEmpty(), + fmt.Sprintf("%d acknowledged write(s) were lost across the handoff: %v", len(lost), lost)) +} + // execValkeyPodShell runs a shell script inside the pod's server container // with VALKEYCLI_AUTH set to the default user's password, so every valkey-cli // invocation in the script runs authenticated. The operator injects From 1bda48e573b88da5b9efa5fbf4329d1c091731d0 Mon Sep 17 00:00:00 2001 From: Sagar Utekar Date: Thu, 9 Jul 2026 09:10:25 +0530 Subject: [PATCH 6/8] test(e2e): fold sigterm handoff coverage into existing failover test Per review, drop the separate failover spec and extend the existing 'should detect and recover when a primary deployment is deleted' test instead, so failover coverage lives in one place. The existing test gains, around its StatefulSet-deletion disruption: - a password-protected default user (passwordSecret, as in #292) so valkey-cli commands run authenticated - 50 keys seeded before the disruption and verified after recovery - a continuous writer through the disruption with per-attempt acks and a 2s connection timeout (a stale MOVED redirect to the terminated primary's IP otherwise hangs a connect for the ~130s TCP SYN timeout) - promotion of the shard's replica asserted within the 30s termination grace period - verification that every write acknowledged during the disruption is readable afterwards - detection of whether the shutdown-on-sigterm handoff engaged (from the promoted replica's log), with the writer-gap bound asserted on the handoff path; the path itself is reported rather than hard-asserted because valkey 9.0's replica selection on shutdown requires exact ack-offset equality and intermittently falls back to failure detection Full e2e suite on Kind: 41 of 41 specs passed. On the handoff path the writer recorded 7904 acknowledged writes with a 0.23s longest gap and zero lost. Signed-off-by: Sagar Utekar --- test/e2e/failover_sigterm_test.go | 399 ------------------------------ test/e2e/valkeycluster_test.go | 297 +++++++++++++++++++++- 2 files changed, 292 insertions(+), 404 deletions(-) delete mode 100644 test/e2e/failover_sigterm_test.go diff --git a/test/e2e/failover_sigterm_test.go b/test/e2e/failover_sigterm_test.go deleted file mode 100644 index 7ed34986..00000000 --- a/test/e2e/failover_sigterm_test.go +++ /dev/null @@ -1,399 +0,0 @@ -//go:build e2e - -/* -Copyright 2025 Valkey Contributors. - -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 ( - "encoding/base64" - "fmt" - "os" - "os/exec" - "path/filepath" - "strings" - "time" - - . "github.com/onsi/ginkgo/v2" - . "github.com/onsi/gomega" - - "valkey.io/valkey-operator/test/utils" -) - -// keyCount is the number of keys written across the keyspace before the -// failover, and verified afterwards to prove no data was lost. -const keyCount = 50 - -// failoverDefaultPassword is the password configured for the default user of -// the failover test clusters (via a passwordSecret, following the pattern -// introduced in #292), so valkey-cli commands run authenticated. -const failoverDefaultPassword = "e2eFailoverPassw0rd" - -var _ = Describe("Shutdown-on-SIGTERM failover", Label("failover"), func() { - var valkeyName string - - AfterEach(func() { - specReport := CurrentSpecReport() - if specReport.Failed() { - utils.CollectDebugInfo(namespace) - } - - By("cleaning up test resources") - cmd := exec.Command("kubectl", "delete", "valkeycluster", valkeyName, "--ignore-not-found=true") - _, _ = utils.Run(cmd) - cmd = exec.Command("kubectl", "delete", "secret", valkeyName+"-users", "--ignore-not-found=true") - _, _ = utils.Run(cmd) - }) - - It("promotes a replica before a gracefully terminated primary exits", func() { - valkeyName = "failover-sigterm" - createFailoverCluster(valkeyName) - - By("identifying the primary and replica of shard 0") - var primaryPod, replicaPod string - Eventually(func(g Gomega) { - primaryPod, replicaPod = getShardRoles(g, valkeyName, 0) - g.Expect(primaryPod).NotTo(BeEmpty(), "shard 0 has no primary") - g.Expect(replicaPod).NotTo(BeEmpty(), "shard 0 has no in-sync replica") - }).WithTimeout(3 * time.Minute).Should(Succeed()) - - By("recording the primary pod UID so its replacement can be detected") - cmd := exec.Command("kubectl", "get", "pod", primaryPod, "-o", "jsonpath={.metadata.uid}") - oldPrimaryUID, err := utils.Run(cmd) - Expect(err).NotTo(HaveOccurred(), "Failed to get primary pod UID") - Expect(oldPrimaryUID).NotTo(BeEmpty()) - - By("writing keys across the keyspace") - writeTestKeys(primaryPod) - - By("starting a continuous writer to run through the disruption") - writer := startContinuousWriter(replicaPod) - - By("gracefully terminating the primary pod (SIGTERM with default grace period)") - cmd = exec.Command("kubectl", "delete", "pod", primaryPod, "--wait=false") - _, err = utils.Run(cmd) - Expect(err).NotTo(HaveOccurred(), "Failed to delete primary pod") - - // The handover must complete inside terminationGracePeriodSeconds - // (default 30s), so the replica has to report role:master within that - // window; otherwise SIGKILL would have cut the handoff short. - By("asserting the replica is promoted within the grace period") - Eventually(func(g Gomega) { - output, err := execValkeyPodShell(replicaPod, "valkey-cli INFO replication") - g.Expect(err).NotTo(HaveOccurred(), "Failed to get replication info from replica") - g.Expect(output).To(ContainSubstring("role:master"), - "replica was not promoted to primary") - }).WithTimeout(30 * time.Second).WithPolling(time.Second).Should(Succeed()) - - // The graceful handoff is a coordinated failover driven by the - // terminating primary (CLUSTER FAILOVER FORCE REPLICAID), which logs - // "Forced failover primary request accepted" on the promoted replica. - // The crash path goes through FAIL detection and a rank-based election - // instead. valkey 9.0's replica selection on shutdown is best-effort - // (it requires exact ack-offset equality at the shutdown instant and - // intermittently falls back to the crash path, ~1 in 3 under write - // load in local testing), so the handoff path is reported rather than - // hard-asserted until the promotion is made deterministic. - By("reporting which failover path drove the promotion") - cmd = exec.Command("kubectl", "logs", replicaPod, "-c", "server") - logs, err := utils.Run(cmd) - Expect(err).NotTo(HaveOccurred(), "Failed to get promoted replica logs") - if strings.Contains(strings.ToLower(logs), "forced failover primary request accepted") { - _, _ = fmt.Fprintf(GinkgoWriter, "promotion path: coordinated shutdown handoff (forced failover)\n") - } else { - _, _ = fmt.Fprintf(GinkgoWriter, - "promotion path: failure detection (shutdown handoff did not engage; see 'Unable to find a replica' on the primary)\n") - } - - By("asserting every write acknowledged during the disruption is readable") - acked, maxGap := writer.stop() - Expect(acked).NotTo(BeEmpty(), "continuous writer recorded no acknowledged writes") - _, _ = fmt.Fprintf(GinkgoWriter, - "continuous writer: %d acknowledged writes, longest gap between acknowledged writes: %.2fs\n", - len(acked), maxGap) - verifyAcknowledgedWrites(replicaPod, acked) - - // The orderly handoff keeps a writer available throughout the - // disruption; a crash-style failover leaves the shard writer-less for - // failure detection plus an election. The bound is deliberately - // generous to absorb CI noise. - Expect(maxGap).To(BeNumerically("<", 10.0), - fmt.Sprintf("shard had no writer for %.2fs during the graceful handoff", maxGap)) - - By("waiting for the replaced pod to come back and rejoin as replica") - Eventually(func(g Gomega) { - cmd := exec.Command("kubectl", "get", "pod", primaryPod, "-o", "jsonpath={.metadata.uid}") - uid, err := utils.Run(cmd) - g.Expect(err).NotTo(HaveOccurred(), "replacement pod does not exist yet") - g.Expect(uid).NotTo(Equal(oldPrimaryUID), "old pod is still terminating") - - output, err := execValkeyPodShell(primaryPod, "valkey-cli INFO replication") - g.Expect(err).NotTo(HaveOccurred(), "Failed to get replication info from replaced pod") - g.Expect(output).To(ContainSubstring("role:slave"), - "replaced pod did not rejoin the shard as a replica") - }).WithTimeout(4 * time.Minute).Should(Succeed()) - - verifyClusterHealthyAndKeysIntact(valkeyName, replicaPod) - }) - -}) - -// createFailoverCluster creates a ValkeyCluster with one replica per shard -// and a password-protected default user, and waits for it to become Ready. -func createFailoverCluster(valkeyName string) { - GinkgoHelper() - - By("creating a ValkeyCluster with one replica per shard") - valkeyYaml := fmt.Sprintf(` -apiVersion: v1 -kind: Secret -metadata: - name: %[1]s-users -data: - defaultpw: %[2]s ---- -apiVersion: valkey.io/v1alpha1 -kind: ValkeyCluster -metadata: - name: %[1]s -spec: - shards: 3 - replicas: 1 - users: - - name: default - enabled: true - permissions: "+@all ~* &*" - passwordSecret: - name: %[1]s-users - keys: [defaultpw] -`, valkeyName, base64.StdEncoding.EncodeToString([]byte(failoverDefaultPassword))) - - manifestFile := filepath.Join(os.TempDir(), fmt.Sprintf("%s.yaml", valkeyName)) - err := os.WriteFile(manifestFile, []byte(valkeyYaml), 0644) - Expect(err).NotTo(HaveOccurred(), "Failed to write manifest file") - defer func() { - Expect(os.Remove(manifestFile)).To(Succeed()) - }() - - cmd := exec.Command("kubectl", "create", "-f", manifestFile) - _, err = utils.Run(cmd) - Expect(err).NotTo(HaveOccurred(), "Failed to create ValkeyCluster CR") - - By("waiting for the cluster to become Ready") - Eventually(func(g Gomega) { - cmd := exec.Command("kubectl", "get", "valkeycluster", valkeyName, - "-o", "jsonpath={.status.conditions[?(@.type==\"Ready\")].status}") - output, err := utils.Run(cmd) - g.Expect(err).NotTo(HaveOccurred(), "Failed to get ValkeyCluster Ready condition") - g.Expect(output).To(Equal("True"), "ValkeyCluster is not Ready") - }).WithTimeout(4 * time.Minute).Should(Succeed()) -} - -// writeTestKeys writes keyCount keys across the keyspace via the given pod. -func writeTestKeys(pod string) { - GinkgoHelper() - - script := fmt.Sprintf( - "ok=0; for i in $(seq 1 %d); do "+ - "r=$(valkey-cli -c set e2e:failover:$i v$i); "+ - "[ \"$r\" = \"OK\" ] && ok=$((ok+1)); "+ - "done; echo written=$ok", keyCount) - output, err := execValkeyPodShell(pod, script) - Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to write keys: %s", output)) - Expect(output).To(ContainSubstring(fmt.Sprintf("written=%d", keyCount)), - fmt.Sprintf("Not all keys were written: %s", output)) -} - -// verifyClusterHealthyAndKeysIntact asserts the ValkeyCluster returns to -// Ready, the cluster reports a healthy state, and every test key survived the -// disruption. -func verifyClusterHealthyAndKeysIntact(valkeyName, pod string) { - GinkgoHelper() - - By("asserting the ValkeyCluster returns to Ready") - Eventually(func(g Gomega) { - cmd := exec.Command("kubectl", "get", "valkeycluster", valkeyName, - "-o", "jsonpath={.status.conditions[?(@.type==\"Ready\")].status}") - output, err := utils.Run(cmd) - g.Expect(err).NotTo(HaveOccurred(), "Failed to get ValkeyCluster Ready condition") - g.Expect(output).To(Equal("True"), "ValkeyCluster did not return to Ready") - }).Should(Succeed()) - - By("asserting the cluster reports a healthy state") - Eventually(func(g Gomega) { - output, err := execValkeyPodShell(pod, "valkey-cli CLUSTER INFO") - g.Expect(err).NotTo(HaveOccurred(), "Failed to get cluster info") - g.Expect(output).To(ContainSubstring("cluster_state:ok")) - }).Should(Succeed()) - - By("asserting no keys were lost across the disruption") - script := fmt.Sprintf( - "ok=0; for i in $(seq 1 %d); do "+ - "v=$(valkey-cli -c get e2e:failover:$i); "+ - "[ \"$v\" = \"v$i\" ] && ok=$((ok+1)); "+ - "done; echo readable=$ok", keyCount) - Eventually(func(g Gomega) { - output, err := execValkeyPodShell(pod, script) - g.Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to read keys back: %s", output)) - g.Expect(output).To(ContainSubstring(fmt.Sprintf("readable=%d", keyCount)), - fmt.Sprintf("Some keys were lost across the disruption: %s", output)) - }).Should(Succeed()) -} - -// continuousWriter tracks a background write loop running inside a pod. -type continuousWriter struct { - cmd *exec.Cmd - output *strings.Builder -} - -// startContinuousWriter starts a background loop inside the pod that writes -// uniquely-numbered keys through the cluster for writerDurationSeconds, -// recording a timestamped ack/fail line per attempt. It is started against -// the shard's replica so cluster-mode redirects follow the primary across the -// handoff. -const writerDurationSeconds = 20 - -func startContinuousWriter(pod string) *continuousWriter { - GinkgoHelper() - - // POSIX-sh only: the image's /bin/sh is dash, which has no $SECONDS. - script := fmt.Sprintf( - "end=$(($(date +%%s)+%d)); i=0; "+ - "while [ \"$(date +%%s)\" -lt \"$end\" ]; do "+ - "r=$(valkey-cli -c set e2e:cw:$i v$i 2>/dev/null | tail -n 1); "+ - "if [ \"$r\" = \"OK\" ]; then echo \"ack $i $(date +%%s.%%N)\"; "+ - "else echo \"fail $i $(date +%%s.%%N)\"; fi; "+ - "i=$((i+1)); "+ - "done", writerDurationSeconds) - cmd := exec.Command("kubectl", "exec", pod, "-c", "server", "--", - "sh", "-c", fmt.Sprintf("export VALKEYCLI_AUTH=%q; ", failoverDefaultPassword)+script) - - w := &continuousWriter{cmd: cmd, output: &strings.Builder{}} - cmd.Stdout = w.output - cmd.Stderr = w.output - Expect(cmd.Start()).To(Succeed(), "Failed to start continuous writer") - return w -} - -// stop waits for the writer loop to finish and returns the map of -// acknowledged key indices to their expected values, plus the longest gap in -// seconds between consecutive acknowledged writes (the shard's effective -// write-unavailability window). -func (w *continuousWriter) stop() (acked map[string]string, maxGap float64) { - GinkgoHelper() - - Expect(w.cmd.Wait()).To(Succeed(), "continuous writer failed: %s", w.output.String()) - - acked = map[string]string{} - lastAckTime := -1.0 - for _, line := range utils.GetNonEmptyLines(w.output.String()) { - fields := strings.Fields(line) - if len(fields) != 3 { - continue - } - status, idx := fields[0], fields[1] - var ts float64 - if _, err := fmt.Sscanf(fields[2], "%f", &ts); err != nil { - continue - } - if status != "ack" { - continue - } - if lastAckTime >= 0 && ts-lastAckTime > maxGap { - maxGap = ts - lastAckTime - } - lastAckTime = ts - acked[idx] = "v" + idx - } - return acked, maxGap -} - -// verifyAcknowledgedWrites asserts that every key the continuous writer got -// an OK for is readable with the expected value — i.e. no acknowledged write -// was dropped during the handoff. -func verifyAcknowledgedWrites(pod string, acked map[string]string) { - GinkgoHelper() - - maxIdx := 0 - for idx := range acked { - var i int - _, err := fmt.Sscanf(idx, "%d", &i) - Expect(err).NotTo(HaveOccurred()) - if i > maxIdx { - maxIdx = i - } - } - - script := fmt.Sprintf( - "for i in $(seq 0 %d); do v=$(valkey-cli -c get e2e:cw:$i 2>/dev/null | tail -n 1); echo \"$i $v\"; done", maxIdx) - output, err := execValkeyPodShell(pod, script) - Expect(err).NotTo(HaveOccurred(), "Failed to read back acknowledged writes") - - readable := map[string]string{} - for _, line := range utils.GetNonEmptyLines(output) { - fields := strings.Fields(line) - if len(fields) == 2 { - readable[fields[0]] = fields[1] - } - } - - var lost []string - for idx, want := range acked { - if readable[idx] != want { - lost = append(lost, idx) - } - } - Expect(lost).To(BeEmpty(), - fmt.Sprintf("%d acknowledged write(s) were lost across the handoff: %v", len(lost), lost)) -} - -// execValkeyPodShell runs a shell script inside the pod's server container -// with VALKEYCLI_AUTH set to the default user's password, so every valkey-cli -// invocation in the script runs authenticated. The operator injects -// VALKEYCLI_AUTH with the _operator user's password for the probe scripts; -// it must be overridden here because valkey-cli auto-sends it as the default -// user's AUTH credential. -func execValkeyPodShell(pod string, script string) (string, error) { - cmd := exec.Command("kubectl", "exec", pod, "-c", "server", "--", - "sh", "-c", fmt.Sprintf("export VALKEYCLI_AUTH=%q; ", failoverDefaultPassword)+script) - return utils.Run(cmd) -} - -// getShardRoles returns the pod names of the primary and replica of the given -// shard. Roles are read live from INFO replication on each pod rather than -// from ValkeyNode status, which can lag behind role changes (see #261). The -// replica is only accepted once its replication link is up, so the failover -// is not attempted against a still-syncing replica. -func getShardRoles(g Gomega, clusterName string, shardIndex int) (primaryPod, replicaPod string) { - cmd := exec.Command("kubectl", "get", "pods", - "-l", fmt.Sprintf("valkey.io/cluster=%s,valkey.io/shard-index=%d", clusterName, shardIndex), - "-o", "go-template={{ range .items }}{{ .metadata.name }}{{ \"\\n\" }}{{ end }}") - output, err := utils.Run(cmd) - g.Expect(err).NotTo(HaveOccurred(), "Failed to list shard pods") - - for _, pod := range utils.GetNonEmptyLines(output) { - pod = strings.TrimSpace(pod) - info, err := execValkeyPodShell(pod, "valkey-cli INFO replication") - g.Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to get replication info from %s", pod)) - switch { - case strings.Contains(info, "role:master"): - primaryPod = pod - case strings.Contains(info, "role:slave") && strings.Contains(info, "master_link_status:up"): - replicaPod = pod - } - } - return primaryPod, replicaPod -} diff --git a/test/e2e/valkeycluster_test.go b/test/e2e/valkeycluster_test.go index 2e11891f..4dfd144a 100644 --- a/test/e2e/valkeycluster_test.go +++ b/test/e2e/valkeycluster_test.go @@ -907,26 +907,48 @@ spec: // fix (shardExistsInTopology + findShardPrimary) now handles this: when // Valkey promotes the replica, the replacement node-index=0 pod joins // as a replica of the promoted primary instead of trying to claim slots. + // + // It also covers the shutdown-on-sigterm handoff (#268/#270): the + // StatefulSet deletion terminates the primary pod gracefully, and the + // test verifies the replica is promoted within the termination grace + // period, that writes acknowledged during the disruption survive it, + // and that keys written before the disruption remain readable. It("should detect and recover when a primary deployment is deleted", func() { - By("creating a ValkeyCluster") - failoverClusterManifest := `apiVersion: valkey.io/v1alpha1 + By("creating a ValkeyCluster with a password-protected default user") + failoverClusterName := "valkeycluster-failover-test" + failoverClusterManifest := fmt.Sprintf(`apiVersion: v1 +kind: Secret +metadata: + name: %[1]s-users +data: + defaultpw: %[2]s +--- +apiVersion: valkey.io/v1alpha1 kind: ValkeyCluster metadata: - name: valkeycluster-failover-test + name: %[1]s spec: shards: 3 replicas: 1 -` + users: + - name: default + enabled: true + permissions: "+@all ~* &*" + passwordSecret: + name: %[1]s-users + keys: [defaultpw] +`, failoverClusterName, base64.StdEncoding.EncodeToString([]byte(failoverDefaultPassword))) manifestFile := filepath.Join(os.TempDir(), "valkeycluster-failover.yaml") err := os.WriteFile(manifestFile, []byte(failoverClusterManifest), 0644) Expect(err).NotTo(HaveOccurred(), "Failed to write manifest file") defer os.Remove(manifestFile) - failoverClusterName := "valkeycluster-failover-test" defer func() { cmd := exec.Command("kubectl", "delete", "valkeycluster", failoverClusterName, "--ignore-not-found=true", "--wait=false") _, _ = utils.Run(cmd) + cmd = exec.Command("kubectl", "delete", "secret", failoverClusterName+"-users", "--ignore-not-found=true") + _, _ = utils.Run(cmd) }() By("applying the CR") @@ -958,11 +980,69 @@ spec: } Eventually(getPrimaryStatefulset).Should(Succeed()) + By("identifying the shard's primary and replica pods") + cmd = exec.Command("kubectl", "get", "statefulset", primaryStatefulset, + "-o", "jsonpath={.metadata.labels.valkey\\.io/shard-index}") + shardIndex, err := utils.Run(cmd) + Expect(err).NotTo(HaveOccurred(), "Failed to get shard index of the primary statefulset") + var primaryPod, replicaPod string + Eventually(func(g Gomega) { + primaryPod, replicaPod = getShardRoles(g, failoverClusterName, shardIndex) + g.Expect(primaryPod).NotTo(BeEmpty(), "shard has no primary") + g.Expect(replicaPod).NotTo(BeEmpty(), "shard has no in-sync replica") + }).WithTimeout(3 * time.Minute).Should(Succeed()) + + By("writing keys across the keyspace") + writeTestKeys(primaryPod) + + By("starting a continuous writer to run through the disruption") + writer := startContinuousWriter(replicaPod) + By(fmt.Sprintf("deleting primary statefulset %s to trigger Valkey failover", primaryStatefulset)) cmd = exec.Command("kubectl", "delete", "statefulset", primaryStatefulset, "--wait=false") _, err = utils.Run(cmd) Expect(err).NotTo(HaveOccurred(), "Failed to delete primary statefulset") + // With shutdown-on-sigterm failover the handover must complete + // inside terminationGracePeriodSeconds (default 30s), so the + // replica has to report role:master within that window. + By("asserting the replica is promoted within the termination grace period") + Eventually(func(g Gomega) { + output, err := execValkeyPodShell(replicaPod, "valkey-cli INFO replication") + g.Expect(err).NotTo(HaveOccurred(), "Failed to get replication info from replica") + g.Expect(output).To(ContainSubstring("role:master"), + "replica was not promoted to primary") + }).WithTimeout(30 * time.Second).WithPolling(time.Second).Should(Succeed()) + + // The graceful handoff is a coordinated failover driven by the + // terminating primary (CLUSTER FAILOVER FORCE REPLICAID), which + // logs "Forced failover primary request accepted" on the promoted + // replica; the crash path goes through FAIL detection instead. + // valkey 9.0's replica selection on shutdown is best-effort (it + // requires exact ack-offset equality at the shutdown instant and + // intermittently falls back to the crash path, ~1 in 3 under + // write load in local testing), so the path is reported rather + // than hard-asserted until the promotion is deterministic. + By("detecting which failover path drove the promotion") + cmd = exec.Command("kubectl", "logs", replicaPod, "-c", "server") + logs, err := utils.Run(cmd) + Expect(err).NotTo(HaveOccurred(), "Failed to get promoted replica logs") + handoffEngaged := strings.Contains(strings.ToLower(logs), "forced failover primary request accepted") + + By("asserting every write acknowledged during the disruption is readable") + acked, maxGap := writer.stop() + Expect(acked).NotTo(BeEmpty(), "continuous writer recorded no acknowledged writes") + _, _ = fmt.Fprintf(GinkgoWriter, + "handoff engaged=%t, %d acknowledged writes, longest writer gap: %.2fs\n", + handoffEngaged, len(acked), maxGap) + verifyAcknowledgedWrites(replicaPod, acked) + if handoffEngaged { + // The orderly handoff keeps a writer available throughout + // the disruption; the bound is generous to absorb CI noise. + Expect(maxGap).To(BeNumerically("<", 10.0), + fmt.Sprintf("shard had no writer for %.2fs during the graceful handoff", maxGap)) + } + By("waiting for the operator to recreate the deployment and the cluster to recover") verifyClusterRecovery := func(g Gomega) { cmd := exec.Command("kubectl", "get", "valkeynodes", @@ -1000,6 +1080,9 @@ spec: } } Eventually(verifyClusterRecovery).Should(Succeed()) + + By("asserting the keys written before the disruption are still readable") + verifySeededKeys(replicaPod) }) }) @@ -1508,3 +1591,207 @@ spec: }) }) }) + +// --------------------------------------------------------------------------- +// Failover test helpers (shutdown-on-sigterm handoff instrumentation, #270). +// --------------------------------------------------------------------------- + +// failoverDefaultPassword is the password configured for the default user of +// the failover test cluster (via a passwordSecret, following the pattern +// introduced in #292), so valkey-cli commands run authenticated. +const failoverDefaultPassword = "e2eFailoverPassw0rd" + +// failoverKeyCount is the number of keys written across the keyspace before +// the disruption, and verified afterwards to prove no data was lost. +const failoverKeyCount = 50 + +// writerDurationSeconds is how long the continuous writer samples the shard. +const writerDurationSeconds = 20 + +// execValkeyPodShell runs a shell script inside the pod's server container +// with VALKEYCLI_AUTH set to the default user's password, so every valkey-cli +// invocation in the script runs authenticated. The operator injects +// VALKEYCLI_AUTH with the _operator user's password for the probe scripts; +// it must be overridden here because valkey-cli auto-sends it as the default +// user's AUTH credential. +func execValkeyPodShell(pod string, script string) (string, error) { + cmd := exec.Command("kubectl", "exec", pod, "-c", "server", "--", + "sh", "-c", fmt.Sprintf("export VALKEYCLI_AUTH=%q; ", failoverDefaultPassword)+script) + return utils.Run(cmd) +} + +// getShardRoles returns the pod names of the primary and replica of the given +// shard. Roles are read live from INFO replication on each pod rather than +// from ValkeyNode status, which can lag behind role changes (see #261). The +// replica is only accepted once its replication link is up, so the failover +// is not attempted against a still-syncing replica. +func getShardRoles(g Gomega, clusterName string, shardIndex string) (primaryPod, replicaPod string) { + cmd := exec.Command("kubectl", "get", "pods", + "-l", fmt.Sprintf("valkey.io/cluster=%s,valkey.io/shard-index=%s", clusterName, shardIndex), + "-o", "go-template={{ range .items }}{{ .metadata.name }}{{ \"\\n\" }}{{ end }}") + output, err := utils.Run(cmd) + g.Expect(err).NotTo(HaveOccurred(), "Failed to list shard pods") + + for _, pod := range utils.GetNonEmptyLines(output) { + pod = strings.TrimSpace(pod) + info, err := execValkeyPodShell(pod, "valkey-cli INFO replication") + g.Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to get replication info from %s", pod)) + switch { + case strings.Contains(info, "role:master"): + primaryPod = pod + case strings.Contains(info, "role:slave") && strings.Contains(info, "master_link_status:up"): + replicaPod = pod + } + } + return primaryPod, replicaPod +} + +// writeTestKeys writes failoverKeyCount keys across the keyspace via the +// given pod. +func writeTestKeys(pod string) { + GinkgoHelper() + + script := fmt.Sprintf( + "ok=0; for i in $(seq 1 %d); do "+ + "r=$(valkey-cli -t 2 -c set e2e:failover:$i v$i 2>/dev/null | tail -n 1); "+ + "[ \"$r\" = \"OK\" ] && ok=$((ok+1)); "+ + "done; echo written=$ok", failoverKeyCount) + output, err := execValkeyPodShell(pod, script) + Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to write keys: %s", output)) + Expect(output).To(ContainSubstring(fmt.Sprintf("written=%d", failoverKeyCount)), + fmt.Sprintf("Not all keys were written: %s", output)) +} + +// verifySeededKeys asserts every key written by writeTestKeys is still +// readable with the expected value. +func verifySeededKeys(pod string) { + GinkgoHelper() + + script := fmt.Sprintf( + "ok=0; for i in $(seq 1 %d); do "+ + "v=$(valkey-cli -t 2 -c get e2e:failover:$i 2>/dev/null | tail -n 1); "+ + "[ \"$v\" = \"v$i\" ] && ok=$((ok+1)); "+ + "done; echo readable=$ok", failoverKeyCount) + Eventually(func(g Gomega) { + output, err := execValkeyPodShell(pod, script) + g.Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to read keys back: %s", output)) + g.Expect(output).To(ContainSubstring(fmt.Sprintf("readable=%d", failoverKeyCount)), + fmt.Sprintf("Some keys were lost across the disruption: %s", output)) + }).Should(Succeed()) +} + +// continuousWriter tracks a background write loop running inside a pod. +type continuousWriter struct { + cmd *exec.Cmd + output *strings.Builder +} + +// startContinuousWriter starts a background loop inside the pod that writes +// uniquely-numbered keys through the cluster for writerDurationSeconds, +// recording a timestamped ack/fail line per attempt. It is started against +// the shard's replica so cluster-mode redirects follow the primary across the +// handoff. +func startContinuousWriter(pod string) *continuousWriter { + GinkgoHelper() + + // POSIX-sh only: the image's /bin/sh is dash, which has no $SECONDS. + // The final "end" sentinel line closes the observation window so a + // write outage lasting until the end of the loop still counts as a gap. + // The -t 2 connection timeout keeps the writer sampling when a stale + // MOVED redirect points at the terminated primary's unroutable IP; + // without it a single connect can hang for the ~130s TCP SYN timeout. + script := fmt.Sprintf( + "end=$(($(date +%%s)+%d)); i=0; "+ + "while [ \"$(date +%%s)\" -lt \"$end\" ]; do "+ + "r=$(valkey-cli -t 2 -c set e2e:cw:$i v$i 2>/dev/null | tail -n 1); "+ + "if [ \"$r\" = \"OK\" ]; then echo \"ack $i $(date +%%s.%%N)\"; "+ + "else echo \"fail $i $(date +%%s.%%N)\"; fi; "+ + "i=$((i+1)); "+ + "done; echo \"end - $(date +%%s.%%N)\"", writerDurationSeconds) + cmd := exec.Command("kubectl", "exec", pod, "-c", "server", "--", + "sh", "-c", fmt.Sprintf("export VALKEYCLI_AUTH=%q; ", failoverDefaultPassword)+script) + + w := &continuousWriter{cmd: cmd, output: &strings.Builder{}} + cmd.Stdout = w.output + cmd.Stderr = w.output + Expect(cmd.Start()).To(Succeed(), "Failed to start continuous writer") + return w +} + +// stop waits for the writer loop to finish and returns the map of +// acknowledged key indices to their expected values, plus the longest gap in +// seconds between consecutive acknowledged writes (the shard's effective +// write-unavailability window). +func (w *continuousWriter) stop() (acked map[string]string, maxGap float64) { + GinkgoHelper() + + Expect(w.cmd.Wait()).To(Succeed(), "continuous writer failed: %s", w.output.String()) + + acked = map[string]string{} + lastAckTime := -1.0 + for _, line := range utils.GetNonEmptyLines(w.output.String()) { + fields := strings.Fields(line) + if len(fields) != 3 { + continue + } + status, idx := fields[0], fields[1] + var ts float64 + if _, err := fmt.Sscanf(fields[2], "%f", &ts); err != nil { + continue + } + // The "end" sentinel closes the window: a write outage running + // through the end of the loop counts as a gap instead of being + // silently dropped. + if status != "ack" && status != "end" { + continue + } + if lastAckTime >= 0 && ts-lastAckTime > maxGap { + maxGap = ts - lastAckTime + } + if status == "end" { + break + } + lastAckTime = ts + acked[idx] = "v" + idx + } + return acked, maxGap +} + +// verifyAcknowledgedWrites asserts that every key the continuous writer got +// an OK for is readable with the expected value — i.e. no acknowledged write +// was dropped during the handoff. +func verifyAcknowledgedWrites(pod string, acked map[string]string) { + GinkgoHelper() + + maxIdx := 0 + for idx := range acked { + var i int + _, err := fmt.Sscanf(idx, "%d", &i) + Expect(err).NotTo(HaveOccurred()) + if i > maxIdx { + maxIdx = i + } + } + + script := fmt.Sprintf( + "for i in $(seq 0 %d); do v=$(valkey-cli -t 2 -c get e2e:cw:$i 2>/dev/null | tail -n 1); echo \"$i $v\"; done", maxIdx) + output, err := execValkeyPodShell(pod, script) + Expect(err).NotTo(HaveOccurred(), "Failed to read back acknowledged writes") + + readable := map[string]string{} + for _, line := range utils.GetNonEmptyLines(output) { + fields := strings.Fields(line) + if len(fields) == 2 { + readable[fields[0]] = fields[1] + } + } + + var lost []string + for idx, want := range acked { + if readable[idx] != want { + lost = append(lost, idx) + } + } + Expect(lost).To(BeEmpty(), + fmt.Sprintf("%d acknowledged write(s) were lost across the handoff: %v", len(lost), lost)) +} From 8484fe2a0fb05e24b531bf024a3736380b36aaf4 Mon Sep 17 00:00:00 2001 From: Sagar Utekar Date: Mon, 13 Jul 2026 12:06:44 +0530 Subject: [PATCH 7/8] test(e2e): batch key read-back through a single valkey-cli instance Per review, stop spawning one valkey-cli process per key when verifying data after the disruption. Both read-back paths now pipe their GETs into one valkey-cli over stdin, with an ECHO KEY: marker before each GET so the raw-mode output is correlated per key without relying on line ordering. Missing keys print an empty line and simply leave no entry. The continuous writer intentionally keeps one process per attempt: each attempt samples fresh-connection availability (what a refilling client pool experiences during the disruption) and records a per-attempt timestamp between commands, which a single long-lived instance would not measure. Verified on Kind: spec passes with the handoff engaged, 7597 acknowledged writes, all read back through the batched path. Signed-off-by: Sagar Utekar --- test/e2e/valkeycluster_test.go | 58 ++++++++++++++++++++++------------ 1 file changed, 38 insertions(+), 20 deletions(-) diff --git a/test/e2e/valkeycluster_test.go b/test/e2e/valkeycluster_test.go index 4dfd144a..58a87669 100644 --- a/test/e2e/valkeycluster_test.go +++ b/test/e2e/valkeycluster_test.go @@ -1662,21 +1662,49 @@ func writeTestKeys(pod string) { fmt.Sprintf("Not all keys were written: %s", output)) } +// readKeysBatch reads the given key prefix for indices [from, to] through a +// single valkey-cli instance: the GETs are piped over stdin with an +// "echo KEY:" marker before each one, so one process serves the whole +// batch and the output can be correlated per key regardless of extra lines. +func readKeysBatch(pod, keyPrefix string, from, to int) (map[string]string, error) { + script := fmt.Sprintf( + "for i in $(seq %d %d); do echo \"echo KEY:$i\"; echo \"get %s$i\"; done | valkey-cli -t 2 -c 2>/dev/null", + from, to, keyPrefix) + output, err := execValkeyPodShell(pod, script) + if err != nil { + return nil, fmt.Errorf("reading keys back: %w (output: %s)", err, output) + } + + values := map[string]string{} + current := "" + for _, line := range strings.Split(output, "\n") { + line = strings.TrimSpace(line) + if after, ok := strings.CutPrefix(line, "KEY:"); ok { + current = after + continue + } + // The last non-empty line before the next marker is the value; + // a missing key prints an empty line and leaves no entry. + if current != "" && line != "" { + values[current] = line + } + } + return values, nil +} + // verifySeededKeys asserts every key written by writeTestKeys is still // readable with the expected value. func verifySeededKeys(pod string) { GinkgoHelper() - script := fmt.Sprintf( - "ok=0; for i in $(seq 1 %d); do "+ - "v=$(valkey-cli -t 2 -c get e2e:failover:$i 2>/dev/null | tail -n 1); "+ - "[ \"$v\" = \"v$i\" ] && ok=$((ok+1)); "+ - "done; echo readable=$ok", failoverKeyCount) Eventually(func(g Gomega) { - output, err := execValkeyPodShell(pod, script) - g.Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to read keys back: %s", output)) - g.Expect(output).To(ContainSubstring(fmt.Sprintf("readable=%d", failoverKeyCount)), - fmt.Sprintf("Some keys were lost across the disruption: %s", output)) + values, err := readKeysBatch(pod, "e2e:failover:", 1, failoverKeyCount) + g.Expect(err).NotTo(HaveOccurred()) + for i := 1; i <= failoverKeyCount; i++ { + idx := fmt.Sprintf("%d", i) + g.Expect(values[idx]).To(Equal("v"+idx), + fmt.Sprintf("seeded key e2e:failover:%s was lost across the disruption", idx)) + } }).Should(Succeed()) } @@ -1773,19 +1801,9 @@ func verifyAcknowledgedWrites(pod string, acked map[string]string) { } } - script := fmt.Sprintf( - "for i in $(seq 0 %d); do v=$(valkey-cli -t 2 -c get e2e:cw:$i 2>/dev/null | tail -n 1); echo \"$i $v\"; done", maxIdx) - output, err := execValkeyPodShell(pod, script) + readable, err := readKeysBatch(pod, "e2e:cw:", 0, maxIdx) Expect(err).NotTo(HaveOccurred(), "Failed to read back acknowledged writes") - readable := map[string]string{} - for _, line := range utils.GetNonEmptyLines(output) { - fields := strings.Fields(line) - if len(fields) == 2 { - readable[fields[0]] = fields[1] - } - } - var lost []string for idx, want := range acked { if readable[idx] != want { From c6c2582c50bd1f87e5bd7e1669b20677852d1cb1 Mon Sep 17 00:00:00 2001 From: Sagar Utekar Date: Mon, 13 Jul 2026 12:13:01 +0530 Subject: [PATCH 8/8] test(e2e): batch the seeded writes through a single valkey-cli too Apply the same single-instance pattern to writeTestKeys: all seed SETs are piped through one valkey-cli, counting the OK responses. Only the continuous writer keeps one process per attempt, deliberately, to sample fresh-connection availability during the disruption. Signed-off-by: Sagar Utekar --- test/e2e/valkeycluster_test.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/test/e2e/valkeycluster_test.go b/test/e2e/valkeycluster_test.go index 58a87669..ac81c9ae 100644 --- a/test/e2e/valkeycluster_test.go +++ b/test/e2e/valkeycluster_test.go @@ -1651,11 +1651,11 @@ func getShardRoles(g Gomega, clusterName string, shardIndex string) (primaryPod, func writeTestKeys(pod string) { GinkgoHelper() + // All SETs are piped through a single valkey-cli instance; each + // successful SET prints exactly "OK" on its own line in raw mode. script := fmt.Sprintf( - "ok=0; for i in $(seq 1 %d); do "+ - "r=$(valkey-cli -t 2 -c set e2e:failover:$i v$i 2>/dev/null | tail -n 1); "+ - "[ \"$r\" = \"OK\" ] && ok=$((ok+1)); "+ - "done; echo written=$ok", failoverKeyCount) + "ok=$(for i in $(seq 1 %d); do echo \"set e2e:failover:$i v$i\"; done "+ + "| valkey-cli -t 2 -c 2>/dev/null | grep -c '^OK$'); echo written=$ok", failoverKeyCount) output, err := execValkeyPodShell(pod, script) Expect(err).NotTo(HaveOccurred(), fmt.Sprintf("Failed to write keys: %s", output)) Expect(output).To(ContainSubstring(fmt.Sprintf("written=%d", failoverKeyCount)),