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 <noreply@anthropic.com>
This commit is contained in:
Michał Marszałek
2026-07-06 14:10:38 +02:00
co-authored by Claude Fable 5
parent 53ae379242
commit 0e45c6b24d
13 changed files with 201 additions and 34 deletions
+13 -4
View File
@@ -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")
})
+21 -6
View File
@@ -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")
})
+14 -4
View File
@@ -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")
})
+36 -14
View File
@@ -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")
})
+20 -6
View File
@@ -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")
+15
View File
@@ -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
+12
View File
@@ -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)
}
+11
View File
@@ -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
+12
View File
@@ -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)
}
+12
View File
@@ -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)
}
+11
View File
@@ -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) {
+12
View File
@@ -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)
}
+12
View File
@@ -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)
}