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

mbalassi 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 e73363f3 [FLINK-35278] Fix occasional NPE on getting latest resource 
for status replace
e73363f3 is described below

commit e73363f3486ed9e1df5cc05c9d0baec7c8c3a37f
Author: Ferenc Csaky <[email protected]>
AuthorDate: Sat May 4 11:12:40 2024 +0200

    [FLINK-35278] Fix occasional NPE on getting latest resource for status 
replace
---
 .../kubernetes/operator/utils/StatusRecorder.java  | 87 +++++++++++++---------
 .../operator/utils/StatusRecorderTest.java         | 27 ++++++-
 2 files changed, 76 insertions(+), 38 deletions(-)

diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/StatusRecorder.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/StatusRecorder.java
index 0197dadd..1b7e3fc2 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/StatusRecorder.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/StatusRecorder.java
@@ -18,6 +18,7 @@
 
 package org.apache.flink.kubernetes.operator.utils;
 
+import org.apache.flink.annotation.VisibleForTesting;
 import org.apache.flink.kubernetes.operator.api.AbstractFlinkResource;
 import org.apache.flink.kubernetes.operator.api.FlinkDeployment;
 import 
org.apache.flink.kubernetes.operator.api.lifecycle.ResourceLifecycleState;
@@ -127,47 +128,65 @@ public class StatusRecorder<
             } catch (KubernetesClientException kce) {
                 // 409 is the error code for conflicts resulting from the 
locking
                 if (kce.getCode() == 409) {
-                    var currentVersion = 
resource.getMetadata().getResourceVersion();
-                    LOG.debug(
-                            "Could not apply status update for resource 
version {}",
-                            currentVersion);
-
-                    var latest = client.resource(resource).get();
-                    var latestVersion = 
latest.getMetadata().getResourceVersion();
-
-                    if (latestVersion.equals(currentVersion)) {
-                        // This should not happen as long as the client works 
consistently
-                        LOG.error("Unable to fetch latest resource version");
-                        throw kce;
-                    }
-
-                    if (latest.getStatus().equals(prevStatus)) {
-                        if (retries++ < 3) {
-                            LOG.debug(
-                                    "Retrying status update for latest version 
{}", latestVersion);
-                            
resource.getMetadata().setResourceVersion(latestVersion);
-                        } else {
-                            // If we cannot get the latest version in 3 tries 
we throw the error to
-                            // retry with delay
-                            throw kce;
-                        }
-                    } else {
-                        throw new StatusConflictException(
-                                "Status have been modified externally in 
version "
-                                        + latestVersion
-                                        + " Previous: "
-                                        + 
objectMapper.writeValueAsString(prevStatus)
-                                        + " Latest: "
-                                        + 
objectMapper.writeValueAsString(latest.getStatus()));
-                    }
+                    handleLockingError(resource, prevStatus, client, retries, 
kce);
+                    ++retries;
                 } else {
-                    // We simply throw non conflict errors, to trigger retry 
with delay
+                    // We simply throw non-conflict errors, to trigger retry 
with delay
                     throw kce;
                 }
             }
         }
     }
 
+    @VisibleForTesting
+    void handleLockingError(
+            CR resource,
+            STATUS prevStatus,
+            KubernetesClient client,
+            int retries,
+            KubernetesClientException kce)
+            throws JsonProcessingException {
+
+        var currentVersion = resource.getMetadata().getResourceVersion();
+        LOG.debug("Could not apply status update for resource version {}", 
currentVersion);
+
+        var latest = client.resource(resource).get();
+        if (latest == null || latest.getMetadata() == null) {
+            // This can happen occasionally, we throw the error to retry with 
delay.
+            throw new KubernetesClientException(
+                    String.format(
+                            "Failed to retrieve latest %s",
+                            latest == null ? "resource" : "metadata"),
+                    kce);
+        }
+
+        var latestVersion = latest.getMetadata().getResourceVersion();
+        if (currentVersion.equals(latestVersion)) {
+            // This should not happen as long as the client works consistently
+            LOG.error("Unable to fetch latest resource version");
+            throw kce;
+        }
+
+        if (latest.getStatus().equals(prevStatus)) {
+            if (retries < 3) {
+                LOG.debug("Retrying status update for latest version {}", 
latestVersion);
+                resource.getMetadata().setResourceVersion(latestVersion);
+            } else {
+                // If we cannot get the latest version in 3 tries we throw the 
error to
+                // retry with delay
+                throw kce;
+            }
+        } else {
+            throw new StatusConflictException(
+                    "Status have been modified externally in version "
+                            + latestVersion
+                            + " Previous: "
+                            + objectMapper.writeValueAsString(prevStatus)
+                            + " Latest: "
+                            + 
objectMapper.writeValueAsString(latest.getStatus()));
+        }
+    }
+
     /**
      * Update the custom resource status based on the in-memory cached to 
ensure that any status
      * updates that we made previously are always visible in the 
reconciliation loop. This is
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/StatusRecorderTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/StatusRecorderTest.java
index aab8eb80..6f9b2e52 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/StatusRecorderTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/StatusRecorderTest.java
@@ -24,11 +24,13 @@ import 
org.apache.flink.kubernetes.operator.api.status.ReconciliationState;
 import org.apache.flink.kubernetes.operator.metrics.MetricManager;
 
 import io.fabric8.kubernetes.client.KubernetesClient;
+import io.fabric8.kubernetes.client.KubernetesClientException;
 import io.fabric8.kubernetes.client.server.mock.EnableKubernetesMockClient;
 import io.fabric8.kubernetes.client.server.mock.KubernetesMockServer;
 import org.junit.jupiter.api.Test;
 
-import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
 
 /** Test for {@link StatusRecorder}. */
 @EnableKubernetesMockClient(crud = true)
@@ -47,17 +49,34 @@ public class StatusRecorderTest {
         var lastRequest = mockServer.getLastRequest();
 
         helper.patchAndCacheStatus(deployment, kubernetesClient);
-        assertTrue(mockServer.getLastRequest() != lastRequest);
+        assertThat(lastRequest).isNotSameAs(mockServer.getLastRequest());
         lastRequest = mockServer.getLastRequest();
         
deployment.getStatus().getReconciliationStatus().setState(ReconciliationState.ROLLING_BACK);
         helper.patchAndCacheStatus(deployment, kubernetesClient);
 
         // We intentionally compare references
-        assertTrue(mockServer.getLastRequest() != lastRequest);
+        assertThat(lastRequest).isNotSameAs(mockServer.getLastRequest());
         lastRequest = mockServer.getLastRequest();
 
         // No update
         helper.patchAndCacheStatus(deployment, kubernetesClient);
-        assertTrue(mockServer.getLastRequest() == lastRequest);
+        assertThat(lastRequest).isSameAs(mockServer.getLastRequest());
+    }
+
+    @Test
+    public void testNullLatestResource() {
+        var statusRecorder =
+                new StatusRecorder<FlinkDeployment, FlinkDeploymentStatus>(
+                        new MetricManager<>(), (e, s) -> {});
+
+        var resource = TestUtils.buildApplicationCluster();
+        var cause = new KubernetesClientException("dummy");
+        assertThatThrownBy(
+                        () ->
+                                statusRecorder.handleLockingError(
+                                        resource, null, kubernetesClient, 0, 
cause))
+                .isInstanceOf(KubernetesClientException.class)
+                .hasMessage("Failed to retrieve latest resource")
+                .hasCause(cause);
     }
 }

Reply via email to