Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -330,7 +330,7 @@ protected Future<Void> manualRollingUpdate(Reconciliation reconciliation, KafkaC
if (!podNamesToRoll.isEmpty()) {
// There are some pods to roll
KafkaConnectRoller roller = new KafkaConnectRoller(reconciliation, connect, operationTimeoutMs, podOperations);
return roller.maybeRoll(podNamesToRoll, pod -> RestartReasons.of(RestartReason.MANUAL_ROLLING_UPDATE))
return VertxUtil.toFuture(roller.maybeRoll(podNamesToRoll, pod -> RestartReasons.of(RestartReason.MANUAL_ROLLING_UPDATE)))
.recover(error -> {
LOGGER.warnCr(reconciliation, "Manual rolling update failed (reconciliation will be continued)", error);
return Future.succeededFuture();
Expand Down Expand Up @@ -373,7 +373,7 @@ protected Future<Void> reconcilePodSet(Reconciliation reconciliation,
return VertxUtil.toFuture(podSetOperations.reconcile(reconciliation, reconciliation.namespace(), connect.getComponentName(), connect.generatePodSet(connect.getReplicas(), podSetAnnotations, podAnnotations, imagePullPolicy, imagePullSecrets, customContainerImage)))
.compose(reconciliationResult -> {
KafkaConnectRoller roller = new KafkaConnectRoller(reconciliation, connect, operationTimeoutMs, podOperations);
return roller.maybeRoll(PodSetUtils.podNames(reconciliationResult.resource()), pod -> KafkaConnectRoller.needsRollingRestart(reconciliation, reconciliationResult.resource(), pod));
return VertxUtil.toFuture(roller.maybeRoll(PodSetUtils.podNames(reconciliationResult.resource()), pod -> KafkaConnectRoller.needsRollingRestart(reconciliation, reconciliationResult.resource(), pod)));
})
.compose(i -> VertxUtil.toFuture(podSetOperations.readiness(reconciliation, reconciliation.namespace(), connect.getComponentName(), 1_000, operationTimeoutMs)));
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -12,16 +12,16 @@
import io.strimzi.operator.cluster.model.PodRevision;
import io.strimzi.operator.cluster.model.RestartReason;
import io.strimzi.operator.cluster.model.RestartReasons;
import io.strimzi.operator.cluster.operator.VertxUtil;
import io.strimzi.operator.cluster.operator.resource.kubernetes.PodOperator;
import io.strimzi.operator.common.Reconciliation;
import io.strimzi.operator.common.ReconciliationLogger;
import io.vertx.core.Future;

import java.util.ArrayDeque;
import java.util.Deque;
import java.util.List;
import java.util.Queue;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionStage;
import java.util.function.Function;

/**
Expand Down Expand Up @@ -65,12 +65,11 @@ public KafkaConnectRoller(
* @param podNeedsRestart Function that evaluates the PodSet and Pods and decides if restart of the Pod is
* needed or not
*
* @return Future which completes when the rolling update is done
* @return CompletionStage which completes when the rolling update is done
*/
public Future<Void> maybeRoll(List<String> podNamesToConsider, Function<Pod, RestartReasons> podNeedsRestart) {
return VertxUtil.toFuture(podOperator.listAsync(reconciliation.namespace(), connect.getSelectorLabels()))
.compose(pods -> Future.succeededFuture(prepareRollingOrder(podNamesToConsider, pods)))
.compose(rollingOrder -> maybeRollPods(podNeedsRestart, rollingOrder));
public CompletionStage<Void> maybeRoll(List<String> podNamesToConsider, Function<Pod, RestartReasons> podNeedsRestart) {
return podOperator.listAsync(reconciliation.namespace(), connect.getSelectorLabels())
.thenCompose(pods -> maybeRollPods(podNeedsRestart, prepareRollingOrder(podNamesToConsider, pods)));
}

/* test */ Queue<String> prepareRollingOrder(List<String> podNamesToConsider, List<Pod> pods) {
Expand Down Expand Up @@ -100,19 +99,18 @@ public Future<Void> maybeRoll(List<String> podNamesToConsider, Function<Pod, Res
* or not
* @param rollingOrder Queue with the pod names in the order of their rolling
*
* @return Future which completes when all pods were rolled / considered for rolling
* @return CompletionStage which completes when all pods were rolled / considered for rolling
*/
private Future<Void> maybeRollPods(Function<Pod, RestartReasons> podNeedsRestart,
Queue<String> rollingOrder) {
private CompletionStage<Void> maybeRollPods(Function<Pod, RestartReasons> podNeedsRestart, Queue<String> rollingOrder) {
String podName = rollingOrder.poll();

if (podName != null) {
// The queue is not empty. We consider rolling of this pod and call this method again to handle the next pod
return maybeRollPod(podNeedsRestart, podName)
.compose(i -> maybeRollPods(podNeedsRestart, rollingOrder));
.thenCompose(i -> maybeRollPods(podNeedsRestart, rollingOrder));
} else {
// Queue is empty => we return completely
return Future.succeededFuture();
return CompletableFuture.completedFuture(null);
}
}

Expand All @@ -125,32 +123,31 @@ private Future<Void> maybeRollPods(Function<Pod, RestartReasons> podNeedsRestart
* or not
* @param podName Name of the pod which should be considered
*
* @return Future which completes when the pod is maybe rolled and ready
* @return CompletionStage which completes when the pod is maybe rolled and ready
*/
/* test */ Future<Void> maybeRollPod(Function<Pod, RestartReasons> podNeedsRestart,
String podName) {
return VertxUtil.toFuture(podOperator.getAsync(reconciliation.namespace(), podName))
.compose(pod -> {
/* test */ CompletionStage<Void> maybeRollPod(Function<Pod, RestartReasons> podNeedsRestart, String podName) {
return podOperator.getAsync(reconciliation.namespace(), podName)
.thenCompose(pod -> {
if (pod == null) {
LOGGER.debugCr(reconciliation, "Pod {} does not exist => waiting for its creation", podName);
return Future.succeededFuture();
return CompletableFuture.completedFuture(null);
} else {
RestartReasons restartReasons = podNeedsRestart.apply(pod);

if (restartReasons.shouldRestart()) {
// Pods changed and needs rolling
LOGGER.infoCr(reconciliation, "Rolling pod {}: {}", podName, restartReasons.getAllReasonNotes());
return VertxUtil.toFuture(podOperator.deleteAsync(reconciliation, reconciliation.namespace(), podName, false));
return podOperator.deleteAsync(reconciliation, reconciliation.namespace(), podName, false);
} else {
// Pod exists and does not need to be rolled
LOGGER.debugCr(reconciliation, "Pod {} does not need to be rolled", podName);
return Future.succeededFuture();
return CompletableFuture.completedFuture(null);
}
}
})
.compose(i -> {
.thenCompose(i -> {
LOGGER.debugCr(reconciliation, "Waiting for pod {} to become ready", podName);
return VertxUtil.toFuture(podOperator.readiness(reconciliation, reconciliation.namespace(), podName, 1_000, operationTimeoutMs));
return podOperator.readiness(reconciliation, reconciliation.namespace(), podName, 1_000, operationTimeoutMs);
});
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,21 +26,21 @@
import io.strimzi.operator.common.Annotations;
import io.strimzi.operator.common.Reconciliation;
import io.strimzi.operator.common.model.Labels;
import io.vertx.junit5.Checkpoint;
import io.vertx.junit5.VertxExtension;
import io.vertx.junit5.VertxTestContext;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.api.Timeout;
import org.mockito.ArgumentCaptor;

import java.util.List;
import java.util.Map;
import java.util.Queue;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
import java.util.concurrent.TimeUnit;

import static org.hamcrest.CoreMatchers.is;
import static org.hamcrest.MatcherAssert.assertThat;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.ArgumentMatchers.anyLong;
Expand All @@ -51,7 +51,7 @@
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;

@ExtendWith(VertxExtension.class)
@Timeout(value = 30, unit = TimeUnit.SECONDS)
@SuppressWarnings("unchecked")
public class KafkaConnectRollerTest {
private final static String NAME = "my-connect";
Expand Down Expand Up @@ -177,7 +177,7 @@ public void testRollingOrderWithUnreadyAndMissingPod() {
}

@Test
public void testMaybeRollPodNoChange(VertxTestContext context) {
public void testMaybeRollPodNoChange() {
StrimziPodSet podSet = new StrimziPodSetBuilder()
.withNewMetadata()
.withName("my-connect-connect")
Expand All @@ -193,17 +193,14 @@ public void testMaybeRollPodNoChange(VertxTestContext context) {

KafkaConnectRoller roller = new KafkaConnectRoller(RECONCILIATION, CLUSTER, 1_000L, mockPodOps);

Checkpoint async = context.checkpoint();
roller.maybeRollPod(pod -> KafkaConnectRoller.needsRollingRestart(RECONCILIATION, podSet, pod), "my-connect-connect-0")
.onComplete(context.succeeding(v -> context.verify(() -> {
verify(mockPodOps, never()).deleteAsync(any(), any(), any(), anyBoolean());
.toCompletableFuture().join();

async.flag();
})));
verify(mockPodOps, never()).deleteAsync(any(), any(), any(), anyBoolean());
}

@Test
public void testMaybeRollPodMissingPod(VertxTestContext context) {
public void testMaybeRollPodMissingPod() {
StrimziPodSet podSet = new StrimziPodSetBuilder()
.withNewMetadata()
.withName("my-connect-connect")
Expand All @@ -219,17 +216,14 @@ public void testMaybeRollPodMissingPod(VertxTestContext context) {

KafkaConnectRoller roller = new KafkaConnectRoller(RECONCILIATION, CLUSTER, 1_000L, mockPodOps);

Checkpoint async = context.checkpoint();
roller.maybeRollPod(pod -> KafkaConnectRoller.needsRollingRestart(RECONCILIATION, podSet, pod), "my-connect-connect-0")
.onComplete(context.succeeding(v -> context.verify(() -> {
verify(mockPodOps, never()).deleteAsync(any(), any(), any(), anyBoolean());
.toCompletableFuture().join();

async.flag();
})));
verify(mockPodOps, never()).deleteAsync(any(), any(), any(), anyBoolean());
}

@Test
public void testMaybeRollPodNeedRolling(VertxTestContext context) {
public void testMaybeRollPodNeedRolling() {
StrimziPodSet podSet = new StrimziPodSetBuilder()
.withNewMetadata()
.withName("my-connect-connect")
Expand All @@ -246,17 +240,14 @@ public void testMaybeRollPodNeedRolling(VertxTestContext context) {

KafkaConnectRoller roller = new KafkaConnectRoller(RECONCILIATION, CLUSTER, 1_000L, mockPodOps);

Checkpoint async = context.checkpoint();
roller.maybeRollPod(pod -> KafkaConnectRoller.needsRollingRestart(RECONCILIATION, podSet, pod), "my-connect-connect-0")
.onComplete(context.succeeding(v -> context.verify(() -> {
verify(mockPodOps, times(1)).deleteAsync(any(), eq(NAMESPACE), eq("my-connect-connect-0"), eq(false));
.toCompletableFuture().join();

async.flag();
})));
verify(mockPodOps, times(1)).deleteAsync(any(), eq(NAMESPACE), eq("my-connect-connect-0"), eq(false));
}

@Test
public void testMaybeRollPodFailsWhenNotReady(VertxTestContext context) {
public void testMaybeRollPodFailsWhenNotReady() {
StrimziPodSet podSet = new StrimziPodSetBuilder()
.withNewMetadata()
.withName("my-connect-connect")
Expand All @@ -272,19 +263,16 @@ public void testMaybeRollPodFailsWhenNotReady(VertxTestContext context) {

KafkaConnectRoller roller = new KafkaConnectRoller(RECONCILIATION, CLUSTER, 1_000L, mockPodOps);

Checkpoint async = context.checkpoint();
roller.maybeRollPod(pod -> KafkaConnectRoller.needsRollingRestart(RECONCILIATION, podSet, pod), "my-connect-connect-0")
.onComplete(context.failing(v -> context.verify(() -> {
verify(mockPodOps, never()).deleteAsync(any(), any(), any(), anyBoolean());

assertThat(v.getMessage(), is("Timeout"));
CompletionException ex = assertThrows(CompletionException.class, () ->
roller.maybeRollPod(pod -> KafkaConnectRoller.needsRollingRestart(RECONCILIATION, podSet, pod), "my-connect-connect-0")
.toCompletableFuture().join());

async.flag();
})));
verify(mockPodOps, never()).deleteAsync(any(), any(), any(), anyBoolean());
assertThat(ex.getCause().getMessage(), is("Timeout"));
}

@Test
public void testMaybeRollInOrder(VertxTestContext context) {
public void testMaybeRollInOrder() {
StrimziPodSet podSet = new StrimziPodSetBuilder()
.withNewMetadata()
.withName("my-connect-connect")
Expand Down Expand Up @@ -314,29 +302,26 @@ public void testMaybeRollInOrder(VertxTestContext context) {

KafkaConnectRoller roller = new KafkaConnectRoller(RECONCILIATION, CLUSTER, 1_000L, mockPodOps);

Checkpoint async = context.checkpoint();
roller.maybeRoll(POD_NAMES, pod -> KafkaConnectRoller.needsRollingRestart(RECONCILIATION, podSet, pod))
.onComplete(context.succeeding(v -> context.verify(() -> {
verify(mockPodOps, never()).deleteAsync(any(), any(), any(), anyBoolean());

List<String> getAsync = getAsyncCaptor.getAllValues();
assertThat(getAsync.size(), is(3));
assertThat(getAsync.get(0), is("my-connect-connect-1"));
assertThat(getAsync.get(1), is("my-connect-connect-0"));
assertThat(getAsync.get(2), is("my-connect-connect-2"));

List<String> readiness = readinessCaptor.getAllValues();
assertThat(readiness.size(), is(3));
assertThat(readiness.get(0), is("my-connect-connect-1"));
assertThat(readiness.get(1), is("my-connect-connect-0"));
assertThat(readiness.get(2), is("my-connect-connect-2"));

async.flag();
})));
.toCompletableFuture().join();

verify(mockPodOps, never()).deleteAsync(any(), any(), any(), anyBoolean());

List<String> getAsync = getAsyncCaptor.getAllValues();
assertThat(getAsync.size(), is(3));
assertThat(getAsync.get(0), is("my-connect-connect-1"));
assertThat(getAsync.get(1), is("my-connect-connect-0"));
assertThat(getAsync.get(2), is("my-connect-connect-2"));

List<String> readiness = readinessCaptor.getAllValues();
assertThat(readiness.size(), is(3));
assertThat(readiness.get(0), is("my-connect-connect-1"));
assertThat(readiness.get(1), is("my-connect-connect-0"));
assertThat(readiness.get(2), is("my-connect-connect-2"));
}

@Test
public void testMaybeRollNotReady(VertxTestContext context) {
public void testMaybeRollNotReady() {
StrimziPodSet podSet = new StrimziPodSetBuilder()
.withNewMetadata()
.withName("my-connect-connect")
Expand All @@ -358,12 +343,11 @@ public void testMaybeRollNotReady(VertxTestContext context) {

KafkaConnectRoller roller = new KafkaConnectRoller(RECONCILIATION, CLUSTER, 1_000L, mockPodOps);

Checkpoint async = context.checkpoint();
roller.maybeRoll(POD_NAMES, pod -> KafkaConnectRoller.needsRollingRestart(RECONCILIATION, podSet, pod))
.onComplete(context.failing(v -> context.verify(() -> {
assertThat(v.getMessage(), is("Timeout"));
async.flag();
})));
CompletionException ex = assertThrows(CompletionException.class, () ->
roller.maybeRoll(POD_NAMES, pod -> KafkaConnectRoller.needsRollingRestart(RECONCILIATION, podSet, pod))
.toCompletableFuture().join());

assertThat(ex.getCause().getMessage(), is("Timeout"));
}

@Test
Expand Down
Loading