From 0e45c6b24ddace846f3ab199d605f4a7c04df004 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Marsza=C5=82ek?= Date: Mon, 6 Jul 2026 14:10:38 +0200 Subject: [PATCH] test(e2e): fix TOCTOU race in CSI reload waits The CSI e2e tests wait for the SPCPS version change before calling WaitReloaded/WaitEnvVar, but Reloader reacts to that same SPCPS update. When Reloader won the race, WaitReloaded captured the already-reloaded annotation as its baseline and then timed out waiting for a further change (seen in CI: "Init container with CSI volume should reload..."). Add WaitReloadedFrom/WaitEnvVarFrom adapter variants that take a caller-supplied baseline, and have the CSI tests capture that baseline before updating the Vault secret. Negative tests also benefit: an erroneous reload that lands during the CSI sync wait is now detected instead of silently absorbed into the baseline. Co-Authored-By: Claude Fable 5 --- test/e2e/advanced/multi_container_test.go | 17 ++++++-- test/e2e/annotations/auto_reload_test.go | 27 +++++++++--- test/e2e/annotations/exclude_test.go | 18 ++++++-- test/e2e/core/workloads_test.go | 50 ++++++++++++++++------- test/e2e/csi/csi_test.go | 26 +++++++++--- test/e2e/utils/workload_adapter.go | 15 +++++++ test/e2e/utils/workload_argo.go | 12 ++++++ test/e2e/utils/workload_cronjob.go | 11 +++++ test/e2e/utils/workload_daemonset.go | 12 ++++++ test/e2e/utils/workload_deployment.go | 12 ++++++ test/e2e/utils/workload_job.go | 11 +++++ test/e2e/utils/workload_openshift.go | 12 ++++++ test/e2e/utils/workload_statefulset.go | 12 ++++++ 13 files changed, 201 insertions(+), 34 deletions(-) diff --git a/test/e2e/advanced/multi_container_test.go b/test/e2e/advanced/multi_container_test.go index 98ed6391..c977720e 100644 --- a/test/e2e/advanced/multi_container_test.go +++ b/test/e2e/advanced/multi_container_test.go @@ -152,6 +152,11 @@ var _ = Describe("Multi-Container Tests", Serial, func() { initialVersion, err := utils.GetSPCPSVersion(ctx, csiClient, testNamespace, spcpsName) Expect(err).NotTo(HaveOccurred()) + // Capture the reload-annotation baseline before the trigger: Reloader reacts to the + // same SPCPS update the test waits on below, so it may reload before WaitReloaded runs. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, deploymentName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret") err = utils.UpdateVaultSecret(ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{ "api_key": "updated-init-value", @@ -163,8 +168,8 @@ var _ = Describe("Multi-Container Tests", Serial, func() { Expect(err).NotTo(HaveOccurred()) By("Waiting for Deployment to be reloaded") - reloaded, err := adapter.WaitReloaded(ctx, testNamespace, deploymentName, - utils.AnnotationLastReloadedFrom, utils.ReloadTimeout) + reloaded, err := adapter.WaitReloadedFrom(ctx, testNamespace, deploymentName, + utils.AnnotationLastReloadedFrom, priorReload, utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue(), "Deployment with init container using CSI volume should be reloaded") }) @@ -199,6 +204,10 @@ var _ = Describe("Multi-Container Tests", Serial, func() { initialVersion, err := utils.GetSPCPSVersion(ctx, csiClient, testNamespace, spcpsName) Expect(err).NotTo(HaveOccurred()) + // Capture the baseline before the trigger — see the sibling test above. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, deploymentName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret") err = utils.UpdateVaultSecret(ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{ "api_key": "updated-init-auto-value", @@ -210,8 +219,8 @@ var _ = Describe("Multi-Container Tests", Serial, func() { Expect(err).NotTo(HaveOccurred()) By("Waiting for Deployment to be reloaded") - reloaded, err := adapter.WaitReloaded(ctx, testNamespace, deploymentName, - utils.AnnotationLastReloadedFrom, utils.ReloadTimeout) + reloaded, err := adapter.WaitReloadedFrom(ctx, testNamespace, deploymentName, + utils.AnnotationLastReloadedFrom, priorReload, utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue(), "Deployment with init container CSI volume and auto=true should be reloaded") }) diff --git a/test/e2e/annotations/auto_reload_test.go b/test/e2e/annotations/auto_reload_test.go index c407fa39..63d2244a 100644 --- a/test/e2e/annotations/auto_reload_test.go +++ b/test/e2e/annotations/auto_reload_test.go @@ -243,6 +243,11 @@ var _ = Describe("Auto Reload Annotation Tests", func() { Expect(err).NotTo(HaveOccurred()) GinkgoWriter.Printf("Initial SPCPS version: %s\n", initialVersion) + // Capture the reload-annotation baseline before the trigger: Reloader reacts to the + // same SPCPS update the test waits on below, so it may reload before WaitReloaded runs. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, deploymentName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret") err = utils.UpdateVaultSecret(ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{"api_key": "updated-value-v2"}) Expect(err).NotTo(HaveOccurred()) @@ -253,8 +258,8 @@ var _ = Describe("Auto Reload Annotation Tests", func() { GinkgoWriter.Println("CSI driver synced new secret version") By("Waiting for Deployment to be reloaded") - reloaded, err := adapter.WaitReloaded(ctx, testNamespace, deploymentName, - utils.AnnotationLastReloadedFrom, utils.ReloadTimeout) + reloaded, err := adapter.WaitReloadedFrom(ctx, testNamespace, deploymentName, + utils.AnnotationLastReloadedFrom, priorReload, utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue(), "Deployment should have been reloaded for Vault secret change") }) @@ -305,6 +310,11 @@ var _ = Describe("Auto Reload Annotation Tests", func() { initialVersion, err := utils.GetSPCPSVersion(ctx, csiClient, testNamespace, spcpsName) Expect(err).NotTo(HaveOccurred()) + // Capture the baseline before the trigger to avoid racing Reloader's own reaction + // to the SPCPS update below. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, deploymentName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret (should trigger reload)") err = utils.UpdateVaultSecret(ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{"api_key": "updated-value-v2"}) Expect(err).NotTo(HaveOccurred()) @@ -314,8 +324,8 @@ var _ = Describe("Auto Reload Annotation Tests", func() { Expect(err).NotTo(HaveOccurred()) By("Waiting for Deployment to be reloaded for SPC change") - reloaded, err = adapter.WaitReloaded(ctx, testNamespace, deploymentName, - utils.AnnotationLastReloadedFrom, utils.ReloadTimeout) + reloaded, err = adapter.WaitReloadedFrom(ctx, testNamespace, deploymentName, + utils.AnnotationLastReloadedFrom, priorReload, utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue(), "Deployment should have been reloaded for Vault secret change") }) @@ -349,6 +359,11 @@ var _ = Describe("Auto Reload Annotation Tests", func() { initialVersion, err := utils.GetSPCPSVersion(ctx, csiClient, testNamespace, spcpsName) Expect(err).NotTo(HaveOccurred()) + // Capture the baseline before the trigger to avoid racing Reloader's own reaction + // to the SPCPS update below. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, deploymentName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret") err = utils.UpdateVaultSecret(ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{"api_key": "updated-value-v2"}) Expect(err).NotTo(HaveOccurred()) @@ -358,8 +373,8 @@ var _ = Describe("Auto Reload Annotation Tests", func() { Expect(err).NotTo(HaveOccurred()) By("Waiting for Deployment to be reloaded") - reloaded, err := adapter.WaitReloaded(ctx, testNamespace, deploymentName, - utils.AnnotationLastReloadedFrom, utils.ReloadTimeout) + reloaded, err := adapter.WaitReloadedFrom(ctx, testNamespace, deploymentName, + utils.AnnotationLastReloadedFrom, priorReload, utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue(), "Deployment with auto=true should have been reloaded for Vault secret change") }) diff --git a/test/e2e/annotations/exclude_test.go b/test/e2e/annotations/exclude_test.go index 73e0e8f0..37a05abc 100644 --- a/test/e2e/annotations/exclude_test.go +++ b/test/e2e/annotations/exclude_test.go @@ -304,6 +304,11 @@ var _ = Describe("Exclude Annotation Tests", func() { initialVersion, err := utils.GetSPCPSVersion(ctx, csiClient, testNamespace, spcpsName) Expect(err).NotTo(HaveOccurred()) + // Capture the baseline before the trigger so an erroneous reload happening while we + // wait for the CSI sync below is still detected. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, deploymentName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret for excluded SPC") err = utils.UpdateVaultSecret(ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{ "api_key": "updated-excluded-value", @@ -316,8 +321,8 @@ var _ = Describe("Exclude Annotation Tests", func() { By("Verifying Deployment was NOT reloaded (excluded SPC)") time.Sleep(utils.NegativeTestWait) - reloaded, err := adapter.WaitReloaded(ctx, testNamespace, deploymentName, - utils.AnnotationLastReloadedFrom, utils.ShortTimeout) + reloaded, err := adapter.WaitReloadedFrom(ctx, testNamespace, deploymentName, + utils.AnnotationLastReloadedFrom, priorReload, utils.ShortTimeout) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeFalse(), "Deployment should NOT reload when excluded SecretProviderClassPodStatus changes") }) @@ -365,6 +370,11 @@ var _ = Describe("Exclude Annotation Tests", func() { initialVersion, err := utils.GetSPCPSVersion(ctx, csiClient, testNamespace, spcpsName2) Expect(err).NotTo(HaveOccurred()) + // Capture the baseline before the trigger to avoid racing Reloader's own reaction + // to the SPCPS update below. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, deploymentName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret for non-excluded SPC") err = utils.UpdateVaultSecret(ctx, kubeClient, restConfig, vaultSecretPath2, map[string]string{ "api_key": "updated-nonexcluded-value", @@ -376,8 +386,8 @@ var _ = Describe("Exclude Annotation Tests", func() { Expect(err).NotTo(HaveOccurred()) By("Waiting for Deployment to be reloaded") - reloaded, err := adapter.WaitReloaded(ctx, testNamespace, deploymentName, - utils.AnnotationLastReloadedFrom, utils.ReloadTimeout) + reloaded, err := adapter.WaitReloadedFrom(ctx, testNamespace, deploymentName, + utils.AnnotationLastReloadedFrom, priorReload, utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue(), "Deployment should reload when non-excluded SecretProviderClassPodStatus changes") }) diff --git a/test/e2e/core/workloads_test.go b/test/e2e/core/workloads_test.go index 1a7f7b37..da789dff 100644 --- a/test/e2e/core/workloads_test.go +++ b/test/e2e/core/workloads_test.go @@ -178,6 +178,11 @@ var _ = Describe("Workload Reload Tests", Serial, func() { Expect(err).NotTo(HaveOccurred()) GinkgoWriter.Printf("Initial SPCPS version: %s\n", initialVersion) + // Capture the reload-annotation baseline before the trigger: Reloader reacts to the + // same SPCPS update the test waits on below, so it may reload before WaitReloaded runs. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, workloadName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret") err = utils.UpdateVaultSecret(ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{"api_key": "updated-value-v2"}) Expect(err).NotTo(HaveOccurred()) @@ -189,8 +194,8 @@ var _ = Describe("Workload Reload Tests", Serial, func() { GinkgoWriter.Println("CSI driver synced new secret version") By("Waiting for workload to be reloaded") - reloaded, err := adapter.WaitReloaded(ctx, testNamespace, workloadName, utils.AnnotationLastReloadedFrom, - utils.ReloadTimeout) + reloaded, err := adapter.WaitReloadedFrom(ctx, testNamespace, workloadName, utils.AnnotationLastReloadedFrom, + priorReload, utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue(), "%s should have been reloaded when Vault secret changed", workloadType) }, Entry("Deployment", Label("csi"), utils.WorkloadDeployment), @@ -1041,6 +1046,11 @@ var _ = Describe("Workload Reload Tests", Serial, func() { initialVersion, err := utils.GetSPCPSVersion(ctx, csiClient, testNamespace, spcpsName) Expect(err).NotTo(HaveOccurred()) + // Capture the baseline before the trigger to avoid racing Reloader's own + // reaction to the SPCPS update below. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, workloadName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret") err = utils.UpdateVaultSecret(ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{"api_key": "updated-value-v2"}) Expect(err).NotTo(HaveOccurred()) @@ -1051,8 +1061,8 @@ var _ = Describe("Workload Reload Tests", Serial, func() { Expect(err).NotTo(HaveOccurred()) By("Waiting for workload to be reloaded") - reloaded, err := adapter.WaitReloaded(ctx, testNamespace, workloadName, - utils.AnnotationLastReloadedFrom, utils.ReloadTimeout) + reloaded, err := adapter.WaitReloadedFrom(ctx, testNamespace, workloadName, + utils.AnnotationLastReloadedFrom, priorReload, utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue(), "%s should reload with SPC annotation on pod template", workloadType) }, @@ -1108,6 +1118,11 @@ var _ = Describe("Workload Reload Tests", Serial, func() { initialVersion, err := utils.GetSPCPSVersion(ctx, csiClient, testNamespace, spcpsName) Expect(err).NotTo(HaveOccurred()) + // Capture the baseline before the trigger to avoid racing Reloader's own + // reaction to the SPCPS update below. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, workloadName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret") err = utils.UpdateVaultSecret(ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{"api_key": "updated-value-v2"}) Expect(err).NotTo(HaveOccurred()) @@ -1118,8 +1133,8 @@ var _ = Describe("Workload Reload Tests", Serial, func() { Expect(err).NotTo(HaveOccurred()) By("Waiting for workload to be reloaded") - reloaded, err := adapter.WaitReloaded(ctx, testNamespace, workloadName, - utils.AnnotationLastReloadedFrom, utils.ReloadTimeout) + reloaded, err := adapter.WaitReloadedFrom(ctx, testNamespace, workloadName, + utils.AnnotationLastReloadedFrom, priorReload, utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue(), "%s should reload with SPC auto on pod template", workloadType) }, @@ -1411,8 +1426,11 @@ var _ = Describe("Workload Reload Tests", Serial, func() { Expect(err).NotTo(HaveOccurred()) By("Waiting for workload to have STAKATER_ env var") - found, err := adapter.WaitEnvVar(ctx, testNamespace, workloadName, utils.StakaterEnvVarPrefix, - utils.ReloadTimeout) + // Baseline "" (workload is fresh, no STAKATER_ env var yet): Reloader may have + // reloaded already while we waited for the CSI sync above, so the wait must not + // capture its own baseline now. + found, err := adapter.WaitEnvVarFrom(ctx, testNamespace, workloadName, utils.StakaterEnvVarPrefix, + "", utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(found).To(BeTrue(), "%s should have STAKATER_ env var after Vault secret change", workloadType) }, Entry("Deployment", Label("csi"), utils.WorkloadDeployment), @@ -1630,8 +1648,9 @@ var _ = Describe("Workload Reload Tests", Serial, func() { Expect(err).NotTo(HaveOccurred()) By("Waiting for Deployment to have STAKATER_ env var") - found, err := adapter.WaitEnvVar(ctx, testNamespace, workloadName, utils.StakaterEnvVarPrefix, - utils.ReloadTimeout) + // Baseline "" (fresh Deployment): avoids racing Reloader's reaction to the SPCPS update. + found, err := adapter.WaitEnvVarFrom(ctx, testNamespace, workloadName, utils.StakaterEnvVarPrefix, + "", utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(found).To(BeTrue(), "Deployment with SPC auto annotation should have STAKATER_ env var") }) @@ -1691,8 +1710,10 @@ var _ = Describe("Workload Reload Tests", Serial, func() { By("Verifying Deployment does NOT have STAKATER_ env var") time.Sleep(utils.NegativeTestWait) - found, err := adapter.WaitEnvVar(ctx, testNamespace, workloadName, utils.StakaterEnvVarPrefix, - utils.ShortTimeout) + // Baseline "" (fresh Deployment): an erroneous reload that already happened during the + // CSI sync wait above must still be detected as a failure. + found, err := adapter.WaitEnvVarFrom(ctx, testNamespace, workloadName, utils.StakaterEnvVarPrefix, + "", utils.ShortTimeout) Expect(err).NotTo(HaveOccurred()) Expect(found).To(BeFalse(), "Deployment should NOT have STAKATER_ env var for excluded SPCPS change") }) @@ -1747,8 +1768,9 @@ var _ = Describe("Workload Reload Tests", Serial, func() { Expect(err).NotTo(HaveOccurred()) By("Waiting for Deployment to have STAKATER_ env var") - found, err := adapter.WaitEnvVar(ctx, testNamespace, workloadName, - utils.StakaterEnvVarPrefix, utils.ReloadTimeout) + // Baseline "" (fresh Deployment): avoids racing Reloader's reaction to the SPCPS update. + found, err := adapter.WaitEnvVarFrom(ctx, testNamespace, workloadName, + utils.StakaterEnvVarPrefix, "", utils.ReloadTimeout) Expect(err).NotTo(HaveOccurred()) Expect(found).To(BeTrue(), "Deployment with init container CSI should have STAKATER_ env var") }) diff --git a/test/e2e/csi/csi_test.go b/test/e2e/csi/csi_test.go index ef55f2bd..2c37845a 100644 --- a/test/e2e/csi/csi_test.go +++ b/test/e2e/csi/csi_test.go @@ -72,6 +72,11 @@ var _ = Describe("CSI SecretProviderClass Tests", Label("csi"), Serial, func() { Expect(err).NotTo(HaveOccurred()) GinkgoWriter.Printf("Initial SPCPS version: %s\n", initialVersion) + // Capture the reload-annotation baseline before the trigger: Reloader reacts to the + // same SPCPS update the test waits on below, so it may reload before WaitReloaded runs. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, deploymentName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret") err = utils.UpdateVaultSecret( ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{"api_key": "updated-value-v2"}) @@ -83,9 +88,9 @@ var _ = Describe("CSI SecretProviderClass Tests", Label("csi"), Serial, func() { GinkgoWriter.Println("CSI driver synced new secret version") By("Waiting for Deployment to be reloaded by Reloader") - reloaded, err := adapter.WaitReloaded( + reloaded, err := adapter.WaitReloadedFrom( ctx, testNamespace, deploymentName, - utils.AnnotationLastReloadedFrom, utils.ReloadTimeout, + utils.AnnotationLastReloadedFrom, priorReload, utils.ReloadTimeout, ) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue(), "Deployment should have been reloaded after Vault secret change") @@ -124,6 +129,10 @@ var _ = Describe("CSI SecretProviderClass Tests", Label("csi"), Serial, func() { By("First update to Vault secret") initialVersion, _ := utils.GetSPCPSVersion(ctx, csiClient, testNamespace, spcpsName) + // Capture the baseline before the trigger to avoid racing Reloader's own reaction + // to the SPCPS update below. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, deploymentName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) err = utils.UpdateVaultSecret( ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{"password": "pass-v2"}) Expect(err).NotTo(HaveOccurred()) @@ -133,9 +142,9 @@ var _ = Describe("CSI SecretProviderClass Tests", Label("csi"), Serial, func() { Expect(err).NotTo(HaveOccurred()) By("Waiting for first reload") - reloaded, err := adapter.WaitReloaded( + reloaded, err := adapter.WaitReloadedFrom( ctx, testNamespace, deploymentName, - utils.AnnotationLastReloadedFrom, utils.ReloadTimeout, + utils.AnnotationLastReloadedFrom, priorReload, utils.ReloadTimeout, ) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue()) @@ -230,6 +239,11 @@ var _ = Describe("CSI SecretProviderClass Tests", Label("csi"), Serial, func() { By("Getting SPCPS version before Vault update") initialVersion, _ := utils.GetSPCPSVersion(ctx, csiClient, testNamespace, spcpsName) + // Capture the baseline before the trigger to avoid racing Reloader's own reaction + // to the SPCPS update below. + priorReload, err := adapter.GetPodTemplateAnnotation(ctx, testNamespace, deploymentName, utils.AnnotationLastReloadedFrom) + Expect(err).NotTo(HaveOccurred()) + By("Updating the Vault secret (should trigger reload)") err = utils.UpdateVaultSecret( ctx, kubeClient, restConfig, vaultSecretPath, map[string]string{"token": "token-v2"}) @@ -240,9 +254,9 @@ var _ = Describe("CSI SecretProviderClass Tests", Label("csi"), Serial, func() { Expect(err).NotTo(HaveOccurred()) By("Verifying Deployment WAS reloaded for Vault secret change") - reloaded, err = adapter.WaitReloaded( + reloaded, err = adapter.WaitReloadedFrom( ctx, testNamespace, deploymentName, - utils.AnnotationLastReloadedFrom, utils.ReloadTimeout, + utils.AnnotationLastReloadedFrom, priorReload, utils.ReloadTimeout, ) Expect(err).NotTo(HaveOccurred()) Expect(reloaded).To(BeTrue(), "SPC auto annotation should trigger reload for Vault secret changes") diff --git a/test/e2e/utils/workload_adapter.go b/test/e2e/utils/workload_adapter.go index bc7f80ce..ce632625 100644 --- a/test/e2e/utils/workload_adapter.go +++ b/test/e2e/utils/workload_adapter.go @@ -69,12 +69,27 @@ type WorkloadAdapter interface { // WaitReloaded waits for the workload to have the reload annotation. // Returns true if the annotation was found, false if timeout occurred. + // It captures the baseline annotation value at call time, which races with Reloader + // if the reload trigger happened earlier — prefer WaitReloadedFrom in that case. WaitReloaded(ctx context.Context, namespace, name, annotationKey string, timeout time.Duration) (bool, error) + // WaitReloadedFrom waits for the reload annotation to be present with a value different + // from priorValue. Capture priorValue (via GetPodTemplateAnnotation) BEFORE performing + // the change that triggers the reload; capturing it afterwards can observe the already + // reloaded value and then wait for a further change that never comes. + WaitReloadedFrom(ctx context.Context, namespace, name, annotationKey, priorValue string, timeout time.Duration) (bool, error) + // WaitEnvVar waits for the workload to have a STAKATER_ env var (for envvars strategy). // Returns true if the env var was found, false if timeout occurred. + // It captures the baseline env var value at call time, which races with Reloader + // if the reload trigger happened earlier — prefer WaitEnvVarFrom in that case. WaitEnvVar(ctx context.Context, namespace, name, prefix string, timeout time.Duration) (bool, error) + // WaitEnvVarFrom waits for a STAKATER_ env var whose value differs from priorValue. + // Capture priorValue BEFORE performing the change that triggers the reload + // (an empty priorValue means the env var is expected to appear). + WaitEnvVarFrom(ctx context.Context, namespace, name, prefix, priorValue string, timeout time.Duration) (bool, error) + // SupportsEnvVarStrategy returns true if the workload supports env var reload strategy. // CronJob does not support this as it uses job creation instead. SupportsEnvVarStrategy() bool diff --git a/test/e2e/utils/workload_argo.go b/test/e2e/utils/workload_argo.go index 69d5163e..c976d746 100644 --- a/test/e2e/utils/workload_argo.go +++ b/test/e2e/utils/workload_argo.go @@ -58,6 +58,12 @@ func (a *ArgoRolloutAdapter) WaitReady(ctx context.Context, namespace, name stri // Captures the current annotation value first to avoid false positives from prior reloads. func (a *ArgoRolloutAdapter) WaitReloaded(ctx context.Context, namespace, name, annotationKey string, timeout time.Duration) (bool, error) { priorValue, _ := a.GetPodTemplateAnnotation(ctx, namespace, name, annotationKey) + return a.WaitReloadedFrom(ctx, namespace, name, annotationKey, priorValue, timeout) +} + +// WaitReloadedFrom waits for the reload annotation to be present with a value different from +// priorValue, which the caller captured before triggering the reload. +func (a *ArgoRolloutAdapter) WaitReloadedFrom(ctx context.Context, namespace, name, annotationKey, priorValue string, timeout time.Duration) (bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { return a.rolloutsClient.ArgoprojV1alpha1().Rollouts(namespace).Watch(ctx, opts) } @@ -72,6 +78,12 @@ func (a *ArgoRolloutAdapter) WaitEnvVar(ctx context.Context, namespace, name, pr if r, err := a.rolloutsClient.ArgoprojV1alpha1().Rollouts(namespace).Get(ctx, name, metav1.GetOptions{}); err == nil { priorValue = GetEnvVarValueByPrefix(r.Spec.Template.Spec.Containers, prefix) } + return a.WaitEnvVarFrom(ctx, namespace, name, prefix, priorValue, timeout) +} + +// WaitEnvVarFrom waits for a STAKATER_ env var whose value differs from priorValue, which the +// caller captured before triggering the reload. +func (a *ArgoRolloutAdapter) WaitEnvVarFrom(ctx context.Context, namespace, name, prefix, priorValue string, timeout time.Duration) (bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { return a.rolloutsClient.ArgoprojV1alpha1().Rollouts(namespace).Watch(ctx, opts) } diff --git a/test/e2e/utils/workload_cronjob.go b/test/e2e/utils/workload_cronjob.go index c681cce9..151fe9ae 100644 --- a/test/e2e/utils/workload_cronjob.go +++ b/test/e2e/utils/workload_cronjob.go @@ -50,6 +50,12 @@ func (a *CronJobAdapter) WaitReady(ctx context.Context, namespace, name string, // Captures the current annotation value first to avoid false positives from prior reloads. func (a *CronJobAdapter) WaitReloaded(ctx context.Context, namespace, name, annotationKey string, timeout time.Duration) (bool, error) { priorValue, _ := a.GetPodTemplateAnnotation(ctx, namespace, name, annotationKey) + return a.WaitReloadedFrom(ctx, namespace, name, annotationKey, priorValue, timeout) +} + +// WaitReloadedFrom waits for the reload annotation to be present with a value different from +// priorValue, which the caller captured before triggering the reload. +func (a *CronJobAdapter) WaitReloadedFrom(ctx context.Context, namespace, name, annotationKey, priorValue string, timeout time.Duration) (bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { return a.client.BatchV1().CronJobs(namespace).Watch(ctx, opts) } @@ -62,6 +68,11 @@ func (a *CronJobAdapter) WaitEnvVar(ctx context.Context, namespace, name, prefix return false, ErrUnsupportedOperation } +// WaitEnvVarFrom returns an error because CronJobs don't support env var reload strategy. +func (a *CronJobAdapter) WaitEnvVarFrom(ctx context.Context, namespace, name, prefix, priorValue string, timeout time.Duration) (bool, error) { + return false, ErrUnsupportedOperation +} + // SupportsEnvVarStrategy returns false as CronJobs don't support env var reload strategy. func (a *CronJobAdapter) SupportsEnvVarStrategy() bool { return false diff --git a/test/e2e/utils/workload_daemonset.go b/test/e2e/utils/workload_daemonset.go index 4a7a2b14..74aada0e 100644 --- a/test/e2e/utils/workload_daemonset.go +++ b/test/e2e/utils/workload_daemonset.go @@ -50,6 +50,12 @@ func (a *DaemonSetAdapter) WaitReady(ctx context.Context, namespace, name string // Captures the current annotation value first to avoid false positives from prior reloads. func (a *DaemonSetAdapter) WaitReloaded(ctx context.Context, namespace, name, annotationKey string, timeout time.Duration) (bool, error) { priorValue, _ := a.GetPodTemplateAnnotation(ctx, namespace, name, annotationKey) + return a.WaitReloadedFrom(ctx, namespace, name, annotationKey, priorValue, timeout) +} + +// WaitReloadedFrom waits for the reload annotation to be present with a value different from +// priorValue, which the caller captured before triggering the reload. +func (a *DaemonSetAdapter) WaitReloadedFrom(ctx context.Context, namespace, name, annotationKey, priorValue string, timeout time.Duration) (bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { return a.client.AppsV1().DaemonSets(namespace).Watch(ctx, opts) } @@ -64,6 +70,12 @@ func (a *DaemonSetAdapter) WaitEnvVar(ctx context.Context, namespace, name, pref if ds, err := a.client.AppsV1().DaemonSets(namespace).Get(ctx, name, metav1.GetOptions{}); err == nil { priorValue = GetEnvVarValueByPrefix(ds.Spec.Template.Spec.Containers, prefix) } + return a.WaitEnvVarFrom(ctx, namespace, name, prefix, priorValue, timeout) +} + +// WaitEnvVarFrom waits for a STAKATER_ env var whose value differs from priorValue, which the +// caller captured before triggering the reload. +func (a *DaemonSetAdapter) WaitEnvVarFrom(ctx context.Context, namespace, name, prefix, priorValue string, timeout time.Duration) (bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { return a.client.AppsV1().DaemonSets(namespace).Watch(ctx, opts) } diff --git a/test/e2e/utils/workload_deployment.go b/test/e2e/utils/workload_deployment.go index f7ef5e37..d64e194a 100644 --- a/test/e2e/utils/workload_deployment.go +++ b/test/e2e/utils/workload_deployment.go @@ -51,6 +51,12 @@ func (a *DeploymentAdapter) WaitReady(ctx context.Context, namespace, name strin // does not cause a false positive — the condition triggers only when the value changes. func (a *DeploymentAdapter) WaitReloaded(ctx context.Context, namespace, name, annotationKey string, timeout time.Duration) (bool, error) { priorValue, _ := a.GetPodTemplateAnnotation(ctx, namespace, name, annotationKey) + return a.WaitReloadedFrom(ctx, namespace, name, annotationKey, priorValue, timeout) +} + +// WaitReloadedFrom waits for the reload annotation to be present with a value different from +// priorValue, which the caller captured before triggering the reload. +func (a *DeploymentAdapter) WaitReloadedFrom(ctx context.Context, namespace, name, annotationKey, priorValue string, timeout time.Duration) (bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { return a.client.AppsV1().Deployments(namespace).Watch(ctx, opts) } @@ -66,6 +72,12 @@ func (a *DeploymentAdapter) WaitEnvVar(ctx context.Context, namespace, name, pre if d, err := a.client.AppsV1().Deployments(namespace).Get(ctx, name, metav1.GetOptions{}); err == nil { priorValue = GetEnvVarValueByPrefix(d.Spec.Template.Spec.Containers, prefix) } + return a.WaitEnvVarFrom(ctx, namespace, name, prefix, priorValue, timeout) +} + +// WaitEnvVarFrom waits for a STAKATER_ env var whose value differs from priorValue, which the +// caller captured before triggering the reload. +func (a *DeploymentAdapter) WaitEnvVarFrom(ctx context.Context, namespace, name, prefix, priorValue string, timeout time.Duration) (bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { return a.client.AppsV1().Deployments(namespace).Watch(ctx, opts) } diff --git a/test/e2e/utils/workload_job.go b/test/e2e/utils/workload_job.go index e71c86c2..5f8ad506 100644 --- a/test/e2e/utils/workload_job.go +++ b/test/e2e/utils/workload_job.go @@ -54,11 +54,22 @@ func (a *JobAdapter) WaitReloaded(ctx context.Context, namespace, name, annotati return false, ErrUnsupportedOperation } +// WaitReloadedFrom returns an error because Jobs are recreated, not updated. +// Use the Recreatable interface (GetOriginalUID + WaitRecreated) instead. +func (a *JobAdapter) WaitReloadedFrom(ctx context.Context, namespace, name, annotationKey, priorValue string, timeout time.Duration) (bool, error) { + return false, ErrUnsupportedOperation +} + // WaitEnvVar returns an error because Jobs don't support env var reload strategy. func (a *JobAdapter) WaitEnvVar(ctx context.Context, namespace, name, prefix string, timeout time.Duration) (bool, error) { return false, ErrUnsupportedOperation } +// WaitEnvVarFrom returns an error because Jobs don't support env var reload strategy. +func (a *JobAdapter) WaitEnvVarFrom(ctx context.Context, namespace, name, prefix, priorValue string, timeout time.Duration) (bool, error) { + return false, ErrUnsupportedOperation +} + // WaitRecreated waits for the Job to be recreated with a different UID using watches. func (a *JobAdapter) WaitRecreated(ctx context.Context, namespace, name, originalUID string, timeout time.Duration) (string, bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { diff --git a/test/e2e/utils/workload_openshift.go b/test/e2e/utils/workload_openshift.go index 6f758bf4..2fb11c6c 100644 --- a/test/e2e/utils/workload_openshift.go +++ b/test/e2e/utils/workload_openshift.go @@ -60,6 +60,12 @@ func (a *DeploymentConfigAdapter) WaitReady(ctx context.Context, namespace, name // Captures the current annotation value first to avoid false positives from prior reloads. func (a *DeploymentConfigAdapter) WaitReloaded(ctx context.Context, namespace, name, annotationKey string, timeout time.Duration) (bool, error) { priorValue, _ := a.GetPodTemplateAnnotation(ctx, namespace, name, annotationKey) + return a.WaitReloadedFrom(ctx, namespace, name, annotationKey, priorValue, timeout) +} + +// WaitReloadedFrom waits for the reload annotation to be present with a value different from +// priorValue, which the caller captured before triggering the reload. +func (a *DeploymentConfigAdapter) WaitReloadedFrom(ctx context.Context, namespace, name, annotationKey, priorValue string, timeout time.Duration) (bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { return a.openshiftClient.AppsV1().DeploymentConfigs(namespace).Watch(ctx, opts) } @@ -74,6 +80,12 @@ func (a *DeploymentConfigAdapter) WaitEnvVar(ctx context.Context, namespace, nam if dc, err := a.openshiftClient.AppsV1().DeploymentConfigs(namespace).Get(ctx, name, metav1.GetOptions{}); err == nil && dc.Spec.Template != nil { priorValue = GetEnvVarValueByPrefix(dc.Spec.Template.Spec.Containers, prefix) } + return a.WaitEnvVarFrom(ctx, namespace, name, prefix, priorValue, timeout) +} + +// WaitEnvVarFrom waits for a STAKATER_ env var whose value differs from priorValue, which the +// caller captured before triggering the reload. +func (a *DeploymentConfigAdapter) WaitEnvVarFrom(ctx context.Context, namespace, name, prefix, priorValue string, timeout time.Duration) (bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { return a.openshiftClient.AppsV1().DeploymentConfigs(namespace).Watch(ctx, opts) } diff --git a/test/e2e/utils/workload_statefulset.go b/test/e2e/utils/workload_statefulset.go index d071678a..28052dce 100644 --- a/test/e2e/utils/workload_statefulset.go +++ b/test/e2e/utils/workload_statefulset.go @@ -50,6 +50,12 @@ func (a *StatefulSetAdapter) WaitReady(ctx context.Context, namespace, name stri // Captures the current annotation value first to avoid false positives from prior reloads. func (a *StatefulSetAdapter) WaitReloaded(ctx context.Context, namespace, name, annotationKey string, timeout time.Duration) (bool, error) { priorValue, _ := a.GetPodTemplateAnnotation(ctx, namespace, name, annotationKey) + return a.WaitReloadedFrom(ctx, namespace, name, annotationKey, priorValue, timeout) +} + +// WaitReloadedFrom waits for the reload annotation to be present with a value different from +// priorValue, which the caller captured before triggering the reload. +func (a *StatefulSetAdapter) WaitReloadedFrom(ctx context.Context, namespace, name, annotationKey, priorValue string, timeout time.Duration) (bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { return a.client.AppsV1().StatefulSets(namespace).Watch(ctx, opts) } @@ -64,6 +70,12 @@ func (a *StatefulSetAdapter) WaitEnvVar(ctx context.Context, namespace, name, pr if sts, err := a.client.AppsV1().StatefulSets(namespace).Get(ctx, name, metav1.GetOptions{}); err == nil { priorValue = GetEnvVarValueByPrefix(sts.Spec.Template.Spec.Containers, prefix) } + return a.WaitEnvVarFrom(ctx, namespace, name, prefix, priorValue, timeout) +} + +// WaitEnvVarFrom waits for a STAKATER_ env var whose value differs from priorValue, which the +// caller captured before triggering the reload. +func (a *StatefulSetAdapter) WaitEnvVarFrom(ctx context.Context, namespace, name, prefix, priorValue string, timeout time.Duration) (bool, error) { watchFunc := func(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { return a.client.AppsV1().StatefulSets(namespace).Watch(ctx, opts) }