Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
315 changes: 310 additions & 5 deletions test/e2e/valkeycluster_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1092,26 +1092,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)
Comment on lines 1132 to +1136

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick win

Delete the Secret before the create step, not only after the test.

The cleanup deletes the Secret at line 1135. The pre-create cleanup at line 1140 deletes only the ValkeyCluster. If a previous run aborted before the deferred cleanup ran, the Secret survives. kubectl create -f then fails with AlreadyExists and the whole spec fails on rerun. Add the Secret to the pre-create cleanup, or use kubectl apply -f.

🛠️ Proposed fix
 			By("applying the CR")
 			cmd := exec.Command("kubectl", "delete", "valkeycluster", failoverClusterName, "--ignore-not-found=true")
 			_, _ = utils.Run(cmd)
+			cmd = exec.Command("kubectl", "delete", "secret", failoverClusterName+"-users", "--ignore-not-found=true")
+			_, _ = utils.Run(cmd)
 			cmd = exec.Command("kubectl", "create", "-f", manifestFile)
🧰 Tools
🪛 ast-grep (0.45.0)

[error] 1134-1134: An argument passed to exec.Command/exec.CommandContext is built by concatenating a string literal with dynamic input. If that input is attacker-controlled (and especially when the command is a shell such as sh -c/bash -c), this enables OS command injection. Pass untrusted data as separate, fixed arguments instead of interpolating it into a command string, avoid invoking a shell, and validate/escape the input where a shell is unavoidable.
Context: exec.Command("kubectl", "delete", "secret", failoverClusterName+"-users", "--ignore-not-found=true")
Note: [CWE-78] Improper Neutralization of Special Elements used in an OS Command ('OS Command Injection').

(command-injection-exec-concat-arg-go)

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@test/e2e/valkeycluster_test.go` around lines 1132 - 1136, Update the
pre-create cleanup for the failover cluster to delete both the ValkeyCluster and
its associated users Secret before running kubectl create. Reuse the existing
failoverClusterName and ignore-not-found behavior, while preserving the deferred
cleanup after the test.

}()

By("applying the CR")
Expand Down Expand Up @@ -1143,11 +1165,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())
Comment on lines +1169 to +1178

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick win

Assert shardIndex is non-empty before you use it as a label selector.

If the valkey.io/shard-index label is absent, utils.Run returns an empty string with no error. getShardRoles then builds the selector valkey.io/shard-index=, which matches no pods, and the Eventually block fails after 3 minutes with "shard has no primary". That message hides the real cause. Add an explicit check.

🛠️ Proposed fix
 			shardIndex, err := utils.Run(cmd)
 			Expect(err).NotTo(HaveOccurred(), "Failed to get shard index of the primary statefulset")
+			shardIndex = strings.TrimSpace(shardIndex)
+			Expect(shardIndex).NotTo(BeEmpty(),
+				"statefulset %s has no valkey.io/shard-index label", primaryStatefulset)
 			var primaryPod, replicaPod string
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
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())
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")
shardIndex = strings.TrimSpace(shardIndex)
Expect(shardIndex).NotTo(BeEmpty(),
"statefulset %s has no valkey.io/shard-index label")
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())
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@test/e2e/valkeycluster_test.go` around lines 1169 - 1178, Validate that
shardIndex is non-empty immediately after utils.Run succeeds and before invoking
getShardRoles in the failover flow. Add an explicit assertion with a clear
message identifying the missing valkey.io/shard-index label, while preserving
the existing Eventually checks for primary and replica pods.


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",
Expand Down Expand Up @@ -1185,6 +1265,9 @@ spec:
}
}
Eventually(verifyClusterRecovery).Should(Succeed())

By("asserting the keys written before the disruption are still readable")
verifySeededKeys(replicaPod)
})
})

Expand Down Expand Up @@ -1968,3 +2051,225 @@ 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()

// 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=$(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)),
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:<index>" 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()

Eventually(func(g Gomega) {
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())
}

// 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
}
Comment on lines +2219 to +2238

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Initial write outage is excluded from availability accounting

lastAckTime is initialized only after an ack, while preceding fail records are discarded. If the primary is deleted before the writer's first successful command, an outage longer than the ten-second bound can be followed by one successful write and still report an under-threshold maxGap. Initialize the measurement window from the first valid writer timestamp, including failures, so the first acknowledgement is compared with the start of the observed window.

Artifacts

Authored focused reproduction source

  • A runnable Go program feeds identical >10-second initial-failure writer output to current and window-start accounting, showing the comparison scope.

Reference accounting output (before)

  • The executed reference run measures the 11.500-second outage before the first acknowledgement and reports a nonempty acknowledged map, showing the expected accounting.

Current accounting output (after)

  • The executed current-code run returns `acked=map[3:v3]` and `maxGap=0.100s` for the same 11.500-second initial outage, confirming the defect.

Current source lines

  • The captured target source lines show `lastAckTime` is initialized only by acknowledged writes while failures are skipped, explaining the omission.

View artifacts

T-Rex Ran code and verified through T-Rex

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
}
}

readable, err := readKeysBatch(pod, "e2e:cw:", 0, maxIdx)
Expect(err).NotTo(HaveOccurred(), "Failed to read back acknowledged writes")

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))
}
Comment on lines +2251 to +2275

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win

Retry the read-back, as verifySeededKeys does.

verifyAcknowledgedWrites runs at line 1223, immediately after the disruption and before the cluster-recovery wait at line 1267. At that moment the shard can still return transient errors or stale MOVED redirects, and readKeysBatch uses a 2-second connection timeout. A single failed batch read reports every acknowledged key as lost, which makes the spec flaky. verifySeededKeys already guards against this with Eventually. Apply the same pattern here.

🛠️ Proposed fix
 	maxIdx := 0
 	for idx := range acked {
 		var i int
 		_, err := fmt.Sscanf(idx, "%d", &i)
 		Expect(err).NotTo(HaveOccurred())
 		if i > maxIdx {
 			maxIdx = i
 		}
 	}
 
-	readable, err := readKeysBatch(pod, "e2e:cw:", 0, maxIdx)
-	Expect(err).NotTo(HaveOccurred(), "Failed to read back acknowledged writes")
-
-	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))
+	Eventually(func(g Gomega) {
+		readable, err := readKeysBatch(pod, "e2e:cw:", 0, maxIdx)
+		g.Expect(err).NotTo(HaveOccurred(), "Failed to read back acknowledged writes")
+
+		var lost []string
+		for idx, want := range acked {
+			if readable[idx] != want {
+				lost = append(lost, idx)
+			}
+		}
+		g.Expect(lost).To(BeEmpty(),
+			fmt.Sprintf("%d acknowledged write(s) were lost across the handoff: %v", len(lost), lost))
+	}).Should(Succeed())
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
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
}
}
readable, err := readKeysBatch(pod, "e2e:cw:", 0, maxIdx)
Expect(err).NotTo(HaveOccurred(), "Failed to read back acknowledged writes")
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))
}
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
}
}
Eventually(func(g Gomega) {
readable, err := readKeysBatch(pod, "e2e:cw:", 0, maxIdx)
g.Expect(err).NotTo(HaveOccurred(), "Failed to read back acknowledged writes")
var lost []string
for idx, want := range acked {
if readable[idx] != want {
lost = append(lost, idx)
}
}
g.Expect(lost).To(BeEmpty(),
fmt.Sprintf("%d acknowledged write(s) were lost across the handoff: %v", len(lost), lost))
}).Should(Succeed())
}
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@test/e2e/valkeycluster_test.go` around lines 2251 - 2275, Update
verifyAcknowledgedWrites to retry the read-back using the same Eventually
pattern and timing used by verifySeededKeys. Retry readKeysBatch until transient
errors and stale MOVED responses resolve, while preserving the existing
acknowledged-key comparison and failure message after retries are exhausted.