This is an automated email from the ASF dual-hosted git repository.

gyfora pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-kubernetes-operator.git


The following commit(s) were added to refs/heads/main by this push:
     new 207b149f [FLINK-35292] Set dummy savepoint path during last-state 
upgrade
207b149f is described below

commit 207b149f9b556d68c4ab98d16cdde2f7820659c0
Author: Gyula Fora <[email protected]>
AuthorDate: Mon May 6 10:14:29 2024 +0200

    [FLINK-35292] Set dummy savepoint path during last-state upgrade
---
 .../deployment/AbstractJobReconciler.java          | 11 ++++-
 .../deployment/ApplicationReconciler.java          | 12 +++++
 .../operator/service/AbstractFlinkService.java     |  9 ++++
 .../kubernetes/operator/service/FlinkService.java  |  2 +
 .../kubernetes/operator/utils/FlinkUtils.java      | 24 ++++++++-
 .../kubernetes/operator/utils/SnapshotUtils.java   | 21 ++++++++
 .../kubernetes/operator/TestingFlinkService.java   |  6 +++
 .../deployment/ApplicationReconcilerTest.java      | 23 ++++++++-
 .../ApplicationReconcilerUpgradeModeTest.java      | 54 +++++++++++++++++++-
 .../kubernetes/operator/utils/FlinkUtilsTest.java  | 57 ++++++++++++++++++++++
 10 files changed, 214 insertions(+), 5 deletions(-)

diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java
index 1e256b7a..bdaec55f 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/AbstractJobReconciler.java
@@ -33,6 +33,7 @@ import 
org.apache.flink.kubernetes.operator.api.status.SnapshotTriggerType;
 import 
org.apache.flink.kubernetes.operator.autoscaler.KubernetesJobAutoScalerContext;
 import 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
 import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext;
+import org.apache.flink.kubernetes.operator.exception.RecoveryFailureException;
 import org.apache.flink.kubernetes.operator.reconciler.ReconciliationUtils;
 import org.apache.flink.kubernetes.operator.reconciler.SnapshotType;
 import org.apache.flink.kubernetes.operator.service.CheckpointHistoryWrapper;
@@ -64,6 +65,8 @@ public abstract class AbstractJobReconciler<
 
     private static final Logger LOG = 
LoggerFactory.getLogger(AbstractJobReconciler.class);
 
+    public static final String LAST_STATE_DUMMY_SP_PATH = 
"KUBERNETES_OPERATOR_LAST_STATE";
+
     public AbstractJobReconciler(
             EventRecorder eventRecorder,
             StatusRecorder<CR, STATUS> statusRecorder,
@@ -179,6 +182,12 @@ public abstract class AbstractJobReconciler<
         var flinkService = ctx.getFlinkService();
         if (ReconciliationUtils.isJobInTerminalState(status)
                 && 
!flinkService.isHaMetadataAvailable(ctx.getObserveConfig())) {
+
+            if (!SnapshotUtils.lastSavepointKnown(status)) {
+                throw new RecoveryFailureException(
+                        "Job is in terminal state but last checkpoint is 
unknown, possibly due to an unrecoverable restore error. Manual restore 
required.",
+                        "UpgradeFailed");
+            }
             LOG.info(
                     "Job is in terminal state, ready for upgrade from observed 
latest checkpoint/savepoint");
             return AvailableUpgradeMode.of(UpgradeMode.SAVEPOINT);
@@ -265,7 +274,7 @@ public abstract class AbstractJobReconciler<
             throws Exception {
         Optional<String> savepointOpt = Optional.empty();
 
-        if (spec.getJob().getUpgradeMode() != UpgradeMode.STATELESS) {
+        if (spec.getJob().getUpgradeMode() == UpgradeMode.SAVEPOINT) {
             savepointOpt =
                     Optional.ofNullable(
                                     ctx.getResource()
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java
index 7992a382..72f3a51f 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconciler.java
@@ -29,6 +29,8 @@ import 
org.apache.flink.kubernetes.operator.api.spec.FlinkDeploymentSpec;
 import org.apache.flink.kubernetes.operator.api.spec.UpgradeMode;
 import org.apache.flink.kubernetes.operator.api.status.FlinkDeploymentStatus;
 import 
org.apache.flink.kubernetes.operator.api.status.JobManagerDeploymentStatus;
+import org.apache.flink.kubernetes.operator.api.status.Savepoint;
+import org.apache.flink.kubernetes.operator.api.status.SnapshotTriggerType;
 import 
org.apache.flink.kubernetes.operator.autoscaler.KubernetesJobAutoScalerContext;
 import 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
 import org.apache.flink.kubernetes.operator.controller.FlinkResourceContext;
@@ -159,8 +161,18 @@ public class ApplicationReconciler
                 relatedResource.getStatus().getClusterInfo());
 
         if (savepoint.isPresent()) {
+            // Savepoint deployment
             deployConfig.set(SavepointConfigOptions.SAVEPOINT_PATH, 
savepoint.get());
+        } else if (requireHaMetadata && 
flinkService.atLeastOneCheckpoint(deployConfig)) {
+            // Last state deployment, explicitly set a dummy savepoint path to 
avoid accidental
+            // incorrect state restore in case the HA metadata is deleted by 
the user
+            deployConfig.set(SavepointConfigOptions.SAVEPOINT_PATH, 
LAST_STATE_DUMMY_SP_PATH);
+            status.getJobStatus()
+                    .getSavepointInfo()
+                    .setLastSavepoint(
+                            Savepoint.of(LAST_STATE_DUMMY_SP_PATH, 
SnapshotTriggerType.UNKNOWN));
         } else {
+            // Stateless deployment, remove any user configured savepoint path
             deployConfig.removeConfig(SavepointConfigOptions.SAVEPOINT_PATH);
         }
 
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
index 48088435..001b9673 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
@@ -229,6 +229,15 @@ public abstract class AbstractFlinkService implements 
FlinkService {
         return false;
     }
 
+    @Override
+    public boolean atLeastOneCheckpoint(Configuration conf) {
+        if (FlinkUtils.isKubernetesHAActivated(conf)) {
+            return 
FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(conf, 
kubernetesClient);
+        } else {
+            return isHaMetadataAvailable(conf);
+        }
+    }
+
     @Override
     public JobID submitJobToSessionCluster(
             ObjectMeta meta,
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkService.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkService.java
index 0e7bfca0..1e2630d6 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkService.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/FlinkService.java
@@ -57,6 +57,8 @@ public interface FlinkService {
 
     boolean isHaMetadataAvailable(Configuration conf);
 
+    boolean atLeastOneCheckpoint(Configuration conf);
+
     void submitSessionCluster(Configuration conf) throws Exception;
 
     JobID submitJobToSessionCluster(
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/FlinkUtils.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/FlinkUtils.java
index 10288c70..3d81e509 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/FlinkUtils.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/FlinkUtils.java
@@ -58,6 +58,7 @@ import java.util.Iterator;
 import java.util.LinkedHashMap;
 import java.util.Map;
 import java.util.Optional;
+import java.util.function.Predicate;
 
 import static 
org.apache.flink.kubernetes.utils.Constants.LABEL_CONFIGMAP_TYPE_HIGH_AVAILABILITY;
 
@@ -287,6 +288,20 @@ public class FlinkUtils {
 
     public static boolean isKubernetesHaMetadataAvailable(
             Configuration conf, KubernetesClient kubernetesClient) {
+        return isKubernetesHaMetadataAvailable(
+                conf, kubernetesClient, FlinkUtils::isValidHaConfigMap);
+    }
+
+    public static boolean isKubernetesHaMetadataAvailableWithCheckpoint(
+            Configuration conf, KubernetesClient kubernetesClient) {
+        return isKubernetesHaMetadataAvailable(
+                conf, kubernetesClient, cm -> isValidHaConfigMap(cm) && 
checkpointExists(cm));
+    }
+
+    private static boolean isKubernetesHaMetadataAvailable(
+            Configuration conf,
+            KubernetesClient kubernetesClient,
+            Predicate<ConfigMap> cmPredicate) {
 
         String clusterId = conf.get(KubernetesConfigOptions.CLUSTER_ID);
         String namespace = conf.get(KubernetesConfigOptions.NAMESPACE);
@@ -303,7 +318,7 @@ public class FlinkUtils {
                         .list()
                         .getItems();
 
-        return configMaps.stream().anyMatch(FlinkUtils::isValidHaConfigMap);
+        return configMaps.stream().anyMatch(cmPredicate);
     }
 
     private static boolean isValidHaConfigMap(ConfigMap cm) {
@@ -319,6 +334,13 @@ public class FlinkUtils {
         return name.endsWith("-jobmanager-leader");
     }
 
+    private static boolean checkpointExists(ConfigMap cm) {
+        var data = cm.getData();
+        return data != null
+                && data.keySet().stream()
+                        .anyMatch(s -> 
s.startsWith(Constants.CHECKPOINT_ID_KEY_PREFIX));
+    }
+
     private static boolean isJobGraphKey(Map.Entry<String, String> entry) {
         return entry.getKey().startsWith(Constants.JOB_GRAPH_STORE_KEY_PREFIX);
     }
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/SnapshotUtils.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/SnapshotUtils.java
index 2322e484..ea9345ed 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/SnapshotUtils.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/SnapshotUtils.java
@@ -22,12 +22,14 @@ import org.apache.flink.configuration.Configuration;
 import org.apache.flink.configuration.ConfigurationUtils;
 import org.apache.flink.kubernetes.operator.api.AbstractFlinkResource;
 import org.apache.flink.kubernetes.operator.api.spec.FlinkVersion;
+import org.apache.flink.kubernetes.operator.api.status.CommonStatus;
 import org.apache.flink.kubernetes.operator.api.status.JobStatus;
 import org.apache.flink.kubernetes.operator.api.status.SnapshotInfo;
 import org.apache.flink.kubernetes.operator.api.status.SnapshotTriggerType;
 import 
org.apache.flink.kubernetes.operator.config.KubernetesOperatorConfigOptions;
 import org.apache.flink.kubernetes.operator.reconciler.ReconciliationUtils;
 import org.apache.flink.kubernetes.operator.reconciler.SnapshotType;
+import 
org.apache.flink.kubernetes.operator.reconciler.deployment.AbstractJobReconciler;
 import org.apache.flink.kubernetes.operator.service.FlinkService;
 
 import io.fabric8.kubernetes.client.KubernetesClient;
@@ -402,4 +404,23 @@ public class SnapshotUtils {
             }
         }
     }
+
+    /**
+     * Check if the last snapshot information is known. True if the snapshot 
location is known
+     * explicitly (not implicitly through a last-state upgrade) or if the 
savepoint is known to be
+     * empty.
+     *
+     * @param status Flink resource status
+     * @return True if last savepoint is known
+     */
+    public static boolean lastSavepointKnown(CommonStatus<?> status) {
+        var lastSavepoint = 
status.getJobStatus().getSavepointInfo().getLastSavepoint();
+
+        if (lastSavepoint == null) {
+            return true;
+        }
+        String location = lastSavepoint.getLocation();
+
+        return 
!location.equals(AbstractJobReconciler.LAST_STATE_DUMMY_SP_PATH);
+    }
 }
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java
index 3283191d..b7427b01 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/TestingFlinkService.java
@@ -125,6 +125,7 @@ public class TestingFlinkService extends 
AbstractFlinkService {
     @Setter private boolean isFlinkJobTerminatedWithoutCancellation = false;
     @Setter private boolean isPortReady = true;
     @Setter private boolean haDataAvailable = true;
+    @Setter private boolean checkpointAvailable = true;
     @Setter private boolean jobManagerReady = true;
     @Setter private boolean deployFailure = false;
     @Setter private Runnable sessionJobSubmittedCallback;
@@ -236,6 +237,11 @@ public class TestingFlinkService extends 
AbstractFlinkService {
         }
     }
 
+    @Override
+    public boolean atLeastOneCheckpoint(Configuration conf) {
+        return isHaMetadataAvailable(conf) && checkpointAvailable;
+    }
+
     @Override
     public boolean isHaMetadataAvailable(Configuration conf) {
         return HighAvailabilityMode.isHighAvailabilityModeActivated(conf) && 
haDataAvailable;
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerTest.java
index 6f9e17fd..2840dc32 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerTest.java
@@ -635,8 +635,15 @@ public class ApplicationReconcilerTest extends 
OperatorTestBase {
     private void verifyAndSetRunningJobsToStatus(
             FlinkDeployment deployment,
             List<Tuple3<String, JobStatusMessage, Configuration>> runningJobs) 
{
+        verifyAndSetRunningJobsToStatus(deployment, runningJobs, null);
+    }
+
+    private void verifyAndSetRunningJobsToStatus(
+            FlinkDeployment deployment,
+            List<Tuple3<String, JobStatusMessage, Configuration>> runningJobs,
+            String savepoint) {
         assertEquals(1, runningJobs.size());
-        assertNull(runningJobs.get(0).f0);
+        assertEquals(savepoint, runningJobs.get(0).f0);
         deployment
                 .getStatus()
                 .setJobStatus(
@@ -1124,7 +1131,10 @@ public class ApplicationReconcilerTest extends 
OperatorTestBase {
                 deployment.getSpec().getRestartNonce(), 
lastReconciledSpec.getRestartNonce());
 
         // Set to running to let savepoint upgrade proceed
-        verifyAndSetRunningJobsToStatus(deployment, flinkService.listJobs());
+        verifyAndSetRunningJobsToStatus(
+                deployment,
+                flinkService.listJobs(),
+                ApplicationReconciler.LAST_STATE_DUMMY_SP_PATH);
 
         reconciler.reconcile(deployment, context);
         // Make sure upgrade is properly triggered now
@@ -1132,6 +1142,15 @@ public class ApplicationReconcilerTest extends 
OperatorTestBase {
                 
deployment.getStatus().getReconciliationStatus().deserializeLastReconciledSpec();
         assertEquals(deployment.getSpec().getRestartNonce(), 
lastReconciledSpec.getRestartNonce());
         assertEquals(JobState.SUSPENDED, 
lastReconciledSpec.getJob().getState());
+        assertEquals(UpgradeMode.SAVEPOINT, 
lastReconciledSpec.getJob().getUpgradeMode());
+        assertEquals(
+                "savepoint_0",
+                deployment
+                        .getStatus()
+                        .getJobStatus()
+                        .getSavepointInfo()
+                        .getLastSavepoint()
+                        .getLocation());
     }
 
     @Test
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java
index a5416a04..ae32f197 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/reconciler/deployment/ApplicationReconcilerUpgradeModeTest.java
@@ -52,6 +52,7 @@ import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.Arguments;
 import org.junit.jupiter.params.provider.MethodSource;
+import org.junit.jupiter.params.provider.ValueSource;
 
 import java.time.Duration;
 import java.util.ArrayList;
@@ -568,6 +569,48 @@ public class ApplicationReconcilerUpgradeModeTest extends 
OperatorTestBase {
         assertEquals(JobState.SUSPENDED, 
lastReconciledSpec.getJob().getState());
     }
 
+    @ParameterizedTest
+    @ValueSource(booleans = {true, false})
+    public void testLastStateDummySpPath(boolean checkpointAvailable) throws 
Exception {
+        // Bootstrap running deployment
+        var deployment = TestUtils.buildApplicationCluster();
+        deployment.getSpec().getJob().setUpgradeMode(UpgradeMode.LAST_STATE);
+
+        reconciler.reconcile(deployment, context);
+        verifyAndSetRunningJobsToStatus(deployment, flinkService.listJobs());
+
+        flinkService.setHaDataAvailable(true);
+        flinkService.setCheckpointAvailable(checkpointAvailable);
+
+        // Submit upgrade
+        deployment.getSpec().setRestartNonce(123L);
+        reconciler.reconcile(deployment, context);
+        reconciler.reconcile(deployment, context);
+
+        var lastReconciledSpec =
+                
deployment.getStatus().getReconciliationStatus().deserializeLastReconciledSpec();
+
+        // Make sure we correctly record upgrade mode to last state
+        assertEquals(UpgradeMode.LAST_STATE, 
lastReconciledSpec.getJob().getUpgradeMode());
+
+        if (checkpointAvailable) {
+            assertEquals(
+                    ApplicationReconciler.LAST_STATE_DUMMY_SP_PATH,
+                    deployment
+                            .getStatus()
+                            .getJobStatus()
+                            .getSavepointInfo()
+                            .getLastSavepoint()
+                            .getLocation());
+            assertEquals(
+                    ApplicationReconciler.LAST_STATE_DUMMY_SP_PATH,
+                    flinkService.listJobs().get(0).f0);
+        } else {
+            
assertNull(deployment.getStatus().getJobStatus().getSavepointInfo().getLastSavepoint());
+            assertNull(flinkService.listJobs().get(0).f0);
+        }
+    }
+
     @Test
     public void testUpgradeModeChangeFromSavepointToLastState() throws 
Exception {
         final String expectedSavepointPath = "savepoint_0";
@@ -682,7 +725,16 @@ public class ApplicationReconcilerUpgradeModeTest extends 
OperatorTestBase {
                         .getImage());
         // Upgrade mode changes from stateless to last-state while HA enabled 
previously should not
         // trigger a savepoint
-        assertNull(flinkService.listJobs().get(0).f0);
+        assertEquals(
+                ApplicationReconciler.LAST_STATE_DUMMY_SP_PATH,
+                deployment
+                        .getStatus()
+                        .getJobStatus()
+                        .getSavepointInfo()
+                        .getLastSavepoint()
+                        .getLocation());
+        assertEquals(
+                ApplicationReconciler.LAST_STATE_DUMMY_SP_PATH, 
flinkService.listJobs().get(0).f0);
     }
 
     public static FlinkDeployment buildApplicationCluster(
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/FlinkUtilsTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/FlinkUtilsTest.java
index a4c97f6d..34fdc11b 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/FlinkUtilsTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/FlinkUtilsTest.java
@@ -201,6 +201,63 @@ public class FlinkUtilsTest {
                         kubernetesClient));
     }
 
+    @Test
+    public void kubernetesHaMetaDataCheckpointCheckTest() {
+        var cr = TestUtils.buildApplicationCluster();
+        var confManager = new FlinkConfigManager(new Configuration());
+        assertFalse(
+                FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(
+                        confManager.getDeployConfig(cr.getMetadata(), 
cr.getSpec()),
+                        kubernetesClient));
+
+        var withCheckpoint = Map.of("checkpointID-2", "p");
+        var withoutCheckpoint = Map.of("counter", "2");
+
+        // Wrong CM name
+        createHAConfigMapWithData(
+                cr.getMetadata().getName() + "-wrong-name",
+                cr.getMetadata().getNamespace(),
+                cr.getMetadata().getName(),
+                withCheckpoint);
+        assertFalse(
+                FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(
+                        confManager.getDeployConfig(cr.getMetadata(), 
cr.getSpec()),
+                        kubernetesClient));
+
+        // Missing data
+        createHAConfigMapWithData(
+                cr.getMetadata().getName() + "-000000000000-config-map",
+                cr.getMetadata().getNamespace(),
+                cr.getMetadata().getName(),
+                null);
+        assertFalse(
+                FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(
+                        confManager.getDeployConfig(cr.getMetadata(), 
cr.getSpec()),
+                        kubernetesClient));
+
+        // CM data without CP
+        createHAConfigMapWithData(
+                cr.getMetadata().getName() + "-000000000000-config-map",
+                cr.getMetadata().getNamespace(),
+                cr.getMetadata().getName(),
+                withoutCheckpoint);
+        assertFalse(
+                FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(
+                        confManager.getDeployConfig(cr.getMetadata(), 
cr.getSpec()),
+                        kubernetesClient));
+
+        // CM data with CP
+        createHAConfigMapWithData(
+                cr.getMetadata().getName() + "-000000000000-config-map",
+                cr.getMetadata().getNamespace(),
+                cr.getMetadata().getName(),
+                withCheckpoint);
+        assertTrue(
+                FlinkUtils.isKubernetesHaMetadataAvailableWithCheckpoint(
+                        confManager.getDeployConfig(cr.getMetadata(), 
cr.getSpec()),
+                        kubernetesClient));
+    }
+
     @Test
     public void testJmNeverStartedDetection() {
         var jmDeployment = new Deployment();

Reply via email to