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 4293d583 [FLINK-31860] FlinkDeployments never finalize when namespace 
is deleted
4293d583 is described below

commit 4293d58329af562e9c50216c3005b4577a289b90
Author: zhou-jiang <[email protected]>
AuthorDate: Wed Apr 3 14:31:59 2024 -0700

    [FLINK-31860] FlinkDeployments never finalize when namespace is deleted
---
 .../kubernetes/operator/utils/EventUtils.java      | 28 +++++++++++--
 .../kubernetes/operator/utils/EventUtilsTest.java  | 47 ++++++++++++++++++++++
 2 files changed, 72 insertions(+), 3 deletions(-)

diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/EventUtils.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/EventUtils.java
index ef32cd59..b515cdda 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/EventUtils.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/utils/EventUtils.java
@@ -23,12 +23,17 @@ import io.fabric8.kubernetes.api.model.HasMetadata;
 import io.fabric8.kubernetes.api.model.ObjectMeta;
 import io.fabric8.kubernetes.api.model.ObjectReferenceBuilder;
 import io.fabric8.kubernetes.client.KubernetesClient;
+import io.fabric8.kubernetes.client.KubernetesClientException;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 import javax.annotation.Nullable;
 
+import java.net.HttpURLConnection;
 import java.time.Duration;
 import java.time.Instant;
 import java.util.Map;
+import java.util.Optional;
 import java.util.function.Consumer;
 import java.util.function.Predicate;
 
@@ -37,6 +42,7 @@ import java.util.function.Predicate;
  * 
https://github.com/EnMasseProject/enmasse/blob/master/k8s-api/src/main/java/io/enmasse/k8s/api/KubeEventLogger.java
  */
 public class EventUtils {
+    private static final Logger LOG = 
LoggerFactory.getLogger(EventUtils.class);
 
     public static String generateEventName(
             HasMetadata target,
@@ -108,7 +114,7 @@ public class EventUtils {
             return false;
         } else {
             Event event = buildEvent(target, type, reason, message, component, 
eventName);
-            eventListener.accept(client.resource(event).createOrReplace());
+            createOrReplaceEvent(client, event).ifPresent(eventListener);
             return true;
         }
     }
@@ -139,7 +145,7 @@ public class EventUtils {
         } else {
             Event event = buildEvent(target, type, reason, message, component, 
eventName);
             setLabels(event, labels);
-            eventListener.accept(client.resource(event).createOrReplace());
+            createOrReplaceEvent(client, event).ifPresent(eventListener);
             return true;
         }
     }
@@ -154,7 +160,7 @@ public class EventUtils {
         existing.setCount(existing.getCount() + 1);
         existing.setMessage(message);
         setLabels(existing, labels);
-        eventListener.accept(client.resource(existing).createOrReplace());
+        createOrReplaceEvent(client, existing).ifPresent(eventListener);
     }
 
     private static void setLabels(Event existing, @Nullable Map<String, 
String> labels) {
@@ -213,4 +219,20 @@ public class EventUtils {
                 || (existing.getMetadata() != null
                         && 
dedupePredicate.test(existing.getMetadata().getLabels()));
     }
+
+    private static Optional<Event> createOrReplaceEvent(KubernetesClient 
client, Event event) {
+        try {
+            Event createdEvent = client.resource(event).createOrReplace();
+            return Optional.of(createdEvent);
+        } catch (KubernetesClientException e) {
+            if (e.getCode() == HttpURLConnection.HTTP_FORBIDDEN) {
+                // fail the reconcile when events cannot be delivered, unless 
FORBIDDEN
+                // which can be a result of recoverable rbac issue, or 
namespace is terminating
+                LOG.warn("Cannot create or update events, proceeding.", e);
+            } else {
+                throw e;
+            }
+        }
+        return Optional.empty();
+    }
 }
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/EventUtilsTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/EventUtilsTest.java
index 24a4c999..c1e758f3 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/EventUtilsTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/utils/EventUtilsTest.java
@@ -20,6 +20,7 @@ package org.apache.flink.kubernetes.operator.utils;
 import org.apache.flink.kubernetes.operator.TestUtils;
 
 import io.fabric8.kubernetes.api.model.Event;
+import io.fabric8.kubernetes.api.model.EventBuilder;
 import io.fabric8.kubernetes.client.KubernetesClient;
 import io.fabric8.kubernetes.client.server.mock.EnableKubernetesMockClient;
 import io.fabric8.kubernetes.client.server.mock.KubernetesMockServer;
@@ -30,6 +31,7 @@ import org.junit.jupiter.params.provider.ValueSource;
 
 import javax.annotation.Nullable;
 
+import java.net.HttpURLConnection;
 import java.time.Duration;
 import java.util.Map;
 import java.util.function.Consumer;
@@ -525,4 +527,49 @@ public class EventUtilsTest {
         Assertions.assertEquals(1, event.getCount());
         Assertions.assertNull(eventConsumed);
     }
+
+    @Test
+    public void testCreateOrReplaceEventOnDeletedNamespace() {
+        var consumer =
+                new Consumer<Event>() {
+                    @Override
+                    public void accept(Event event) {
+                        eventConsumed = event;
+                    }
+                };
+        var flinkApp = TestUtils.buildApplicationCluster();
+        var reason = "Cleanup";
+        var message = "message";
+        var eventName =
+                EventUtils.generateEventName(
+                        flinkApp,
+                        EventRecorder.Type.Warning,
+                        reason,
+                        message,
+                        EventRecorder.Component.Operator);
+
+        var namespaceName = flinkApp.getMetadata().getNamespace();
+
+        String eventCreatePath = String.format("/api/v1/namespaces/%s/events", 
namespaceName);
+
+        mockServer
+                .expect()
+                .post()
+                .withPath(eventCreatePath)
+                .andReturn(HttpURLConnection.HTTP_FORBIDDEN, new 
EventBuilder().build())
+                .once();
+
+        Assertions.assertTrue(
+                EventUtils.createOrUpdateEventWithInterval(
+                        kubernetesClient,
+                        flinkApp,
+                        EventRecorder.Type.Warning,
+                        reason,
+                        message,
+                        EventRecorder.Component.Operator,
+                        consumer,
+                        null,
+                        null));
+        Assertions.assertNull(eventConsumed);
+    }
 }

Reply via email to