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);
+ }
}