diff --git a/.github/actions/run-systemtests/action.yml b/.github/actions/run-systemtests/action.yml index 351c4509..1410dc6c 100644 --- a/.github/actions/run-systemtests/action.yml +++ b/.github/actions/run-systemtests/action.yml @@ -35,4 +35,12 @@ runs: ./mvnw verify -Psystemtests -pl systemtests - -Dstyle.color=always + -Dstyle.color=always + + - name: Upload test logs + if: failure() + uses: actions/upload-artifact@v4 + with: + name: test-logs-${{ github.job }} + path: | + operator/systemtests/target/logs/**/* diff --git a/systemtests/pom.xml b/systemtests/pom.xml index 04343285..27007261 100644 --- a/systemtests/pom.xml +++ b/systemtests/pom.xml @@ -15,7 +15,7 @@ 1.10.2 2.24.2 3.27.3 - 0.2.1 + 0.13.0 0.0.1-alpha1 5.0.0-alpha.12 4.2.1 @@ -97,6 +97,11 @@ test-frame-kubernetes ${testframe.version} + + io.skodjob + test-frame-log-collector + ${testframe.version} + io.debezium debezium-operator-api diff --git a/systemtests/src/main/java/io/debezium/operator/systemtests/ConfigProperties.java b/systemtests/src/main/java/io/debezium/operator/systemtests/ConfigProperties.java index 064b8115..b573962f 100644 --- a/systemtests/src/main/java/io/debezium/operator/systemtests/ConfigProperties.java +++ b/systemtests/src/main/java/io/debezium/operator/systemtests/ConfigProperties.java @@ -11,4 +11,5 @@ public final class ConfigProperties { public static final Integer HTTP_POLL_INTERVAL = Integer.valueOf(System.getProperty("test.http.poll.interval", "200")); public static final Integer FABRIC8_POLL_INTERVAL = Integer.valueOf(System.getProperty("test.fabric8.poll.interval", "2")); public static final Integer FABRIC8_POLL_TIMEOUT = Integer.valueOf(System.getProperty("test.fabric8.poll.timeout", "60")); + public static final String LOG_COLLECT_LABEL = System.getProperty("test.log.collector.label", "collect-logs"); } diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/ConfigMapOffsetStorageTest.java b/systemtests/src/test/java/io/debezium/operator/systemtests/ConfigMapOffsetStorageTest.java index 6166b5fe..34444003 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/ConfigMapOffsetStorageTest.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/ConfigMapOffsetStorageTest.java @@ -44,10 +44,10 @@ void testConfigMapOffsetStorage() { .build(); server.getSpec().getSource().setOffset(offset); - KubeResourceManager.getInstance().createResourceWithWait(server); + KubeResourceManager.get().createResourceWithWait(server); assertStreamingWorks(); - ConfigMap configMap = KubeResourceManager.getKubeClient().getClient() + ConfigMap configMap = KubeResourceManager.get().kubeClient().getClient() .configMaps() .inNamespace(namespace) .withName("my-debezium-offsets") @@ -68,7 +68,7 @@ void configMapOffsetStorageMustNotBeCreatedIfAlreadyExists() { logger.info("Deploying Operator"); operatorBundleResource.deploy(); - KubeResourceManager.getInstance().createResourceWithWait(new ConfigMapBuilder() + KubeResourceManager.get().createResourceWithWait(new ConfigMapBuilder() .withMetadata(new ObjectMetaBuilder() .withNamespace(namespace) .withName("debezium-offsets") @@ -85,10 +85,10 @@ void configMapOffsetStorageMustNotBeCreatedIfAlreadyExists() { .build(); server.getSpec().getSource().setOffset(offset); - KubeResourceManager.getInstance().createResourceWithWait(server); + KubeResourceManager.get().createResourceWithWait(server); assertStreamingWorks(); - ConfigMap configMap = KubeResourceManager.getKubeClient().getClient() + ConfigMap configMap = KubeResourceManager.get().kubeClient().getClient() .configMaps() .inNamespace(namespace) .withName("debezium-offsets") diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/RedisOffsetStorageTest.java b/systemtests/src/test/java/io/debezium/operator/systemtests/RedisOffsetStorageTest.java index ee2afc0b..29ea0ac3 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/RedisOffsetStorageTest.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/RedisOffsetStorageTest.java @@ -28,7 +28,7 @@ public class RedisOffsetStorageTest extends TestBase { private final Logger logger = LoggerFactory.getLogger(this.getClass()); @Test - void testRedisOffsetStorage() { + void testRedisOffsetStorage() throws InterruptedException { String namespace = NamespaceHolder.INSTANCE.getCurrentNamespace(); DebeziumOperatorBundleResource operatorBundleResource = new DebeziumOperatorBundleResource(); operatorBundleResource.configureAsDefault(namespace); @@ -45,7 +45,7 @@ void testRedisOffsetStorage() { .build(); server.getSpec().getSource().setOffset(offset); - KubeResourceManager.getInstance().createResourceWithWait(server); + KubeResourceManager.get().createResourceWithWait(server); assertStreamingWorks(); try (LocalPortForward lcp = dmtResource.portForward(portForwardPort, namespace)) { @@ -58,7 +58,7 @@ void testRedisOffsetStorage() { } server.getSpec().getSource().getOffset().getRedis().setKey("metadata:debezium_n:offsets"); - KubeResourceManager.getInstance().createOrUpdateResourceWithWait(server); + KubeResourceManager.get().createOrUpdateResourceWithWait(server); assertStreamingWorks(10, 20); try (LocalPortForward lcp = dmtResource.portForward(portForwardPort, namespace)) { diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/SmokeTest.java b/systemtests/src/test/java/io/debezium/operator/systemtests/SmokeTest.java index 5d960579..002ad28e 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/SmokeTest.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/SmokeTest.java @@ -27,7 +27,7 @@ void firstInstanceIT() { operatorBundleResource.deploy(); logger.info("Deploying Debezium Server"); DebeziumServer server = DebeziumServerGenerator.generateDefaultMysqlToRedis(namespace); - KubeResourceManager.getInstance().createResourceWithWait(server); + KubeResourceManager.get().createResourceWithWait(server); assertStreamingWorks(); } @@ -40,7 +40,7 @@ void secondInstanceIT() { operatorBundleResource.deploy(); logger.info("Deploying Debezium Server"); DebeziumServer server = DebeziumServerGenerator.generateDefaultMysqlToRedis(namespace); - KubeResourceManager.getInstance().createResourceWithWait(server); + KubeResourceManager.get().createResourceWithWait(server); assertStreamingWorks(); } } diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/TestBase.java b/systemtests/src/test/java/io/debezium/operator/systemtests/TestBase.java index 94527659..a8526835 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/TestBase.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/TestBase.java @@ -10,8 +10,13 @@ import static org.awaitility.Awaitility.await; import java.io.IOException; +import java.nio.file.Path; +import java.nio.file.Paths; import java.time.Duration; +import java.time.LocalDateTime; +import java.time.format.DateTimeFormatter; +import org.awaitility.core.ConditionTimeoutException; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeAll; @@ -24,12 +29,15 @@ import io.debezium.operator.systemtests.resources.databases.MysqlResource; import io.debezium.operator.systemtests.resources.dmt.DmtClient; import io.debezium.operator.systemtests.resources.dmt.DmtResource; +import io.debezium.operator.systemtests.resources.logs.MustGatherImpl; import io.debezium.operator.systemtests.resources.sinks.RedisResource; import io.fabric8.kubernetes.client.LocalPortForward; +import io.skodjob.testframe.annotations.MustGather; import io.skodjob.testframe.annotations.ResourceManager; import io.skodjob.testframe.annotations.TestVisualSeparator; -@ResourceManager +@ResourceManager(asyncDeletion = false) +@MustGather(config = MustGatherImpl.class) @TestVisualSeparator @DebeziumResourceTypes @TestInstance(TestInstance.Lifecycle.PER_CLASS) @@ -38,6 +46,10 @@ public class TestBase { protected final DmtResource dmtResource = new DmtResource(); protected final String portForwardHost = "127.0.0.1"; protected int portForwardPort = 8080; + private static final DateTimeFormatter DATE_FORMAT = DateTimeFormatter.ofPattern("yyyy-MM-dd_HH-mm"); + private static final String USER_PATH = System.getProperty("user.dir"); + public static final Path LOG_DIR = Paths.get(USER_PATH, "target", "logs") + .resolve("test-run-" + DATE_FORMAT.format(LocalDateTime.now())); @BeforeAll void initDefault() { @@ -86,7 +98,15 @@ public void assertStreamingWorks(int messagesToDatabase, int expectedMessages) { DmtClient.waitForFilledRedis(portForwardHost, portForwardPort, Duration.ofSeconds(60), "inventory.inventory.operator_test"); await().atMost(Duration.ofMinutes(HTTP_POLL_TIMEOUT)) .pollInterval(Duration.ofMillis(HTTP_POLL_INTERVAL)) - .until(() -> DmtClient.digStreamedData(portForwardHost, portForwardPort, expectedMessages) == expectedMessages); + .until(() -> { + try { + return DmtClient.digStreamedData(portForwardHost, portForwardPort, expectedMessages) == expectedMessages; + } + catch (ConditionTimeoutException ex) { + logger.error(ex.getMessage()); + return false; + } + }); } catch (IOException e) { throw new RuntimeException(e); diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/NamespaceHolder.java b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/NamespaceHolder.java index bde9068d..46c3b653 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/NamespaceHolder.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/NamespaceHolder.java @@ -30,7 +30,7 @@ public void createNewNamespace() { .endMetadata() .build(); this.currentNamespace = name; - KubeResourceManager.getInstance().createResourceWithWait(namespace); + KubeResourceManager.get().createResourceWithWait(namespace); } public String getCurrentNamespace() { diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/annotations/extensions/DebeziumResourceTypesExtension.java b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/annotations/extensions/DebeziumResourceTypesExtension.java index 3a335874..04348151 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/annotations/extensions/DebeziumResourceTypesExtension.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/annotations/extensions/DebeziumResourceTypesExtension.java @@ -8,14 +8,31 @@ import org.junit.jupiter.api.extension.BeforeAllCallback; import org.junit.jupiter.api.extension.ExtensionContext; +import io.debezium.operator.systemtests.ConfigProperties; import io.debezium.operator.systemtests.resources.server.DebeziumServerResource; +import io.skodjob.testframe.resources.ConfigMapType; import io.skodjob.testframe.resources.CustomResourceDefinitionType; +import io.skodjob.testframe.resources.DeploymentType; import io.skodjob.testframe.resources.KubeResourceManager; import io.skodjob.testframe.resources.NamespaceType; +import io.skodjob.testframe.resources.ServiceType; +import io.skodjob.testframe.utils.KubeUtils; public class DebeziumResourceTypesExtension implements BeforeAllCallback { @Override public void beforeAll(ExtensionContext extensionContext) { - KubeResourceManager.getInstance().setResourceTypes(new NamespaceType(), new CustomResourceDefinitionType(), new DebeziumServerResource()); + KubeResourceManager.get().setResourceTypes( + new NamespaceType(), + new CustomResourceDefinitionType(), + new DebeziumServerResource(), + new DeploymentType(), + new ServiceType(), + new ConfigMapType()); + + KubeResourceManager.get().addCreateCallback(r -> { + if (r.getKind().equals("Namespace")) { + KubeUtils.labelNamespace(r.getMetadata().getName(), ConfigProperties.LOG_COLLECT_LABEL, "true"); + } + }); } } diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/databases/MysqlResource.java b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/databases/MysqlResource.java index 7cfae309..13aa1f39 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/databases/MysqlResource.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/databases/MysqlResource.java @@ -202,7 +202,7 @@ public void configureAsDefault(String namespace) { } public void deploy() { - KubeResourceManager.getInstance().createResourceWithoutWait(this.persistentVolumeClaim, this.credentials, this.service); - KubeResourceManager.getInstance().createResourceWithWait(this.deployment); + KubeResourceManager.get().createResourceWithoutWait(this.persistentVolumeClaim, this.credentials, this.service); + KubeResourceManager.get().createResourceWithWait(this.deployment); } } diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/dmt/DmtClient.java b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/dmt/DmtClient.java index 53e27020..70d16e92 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/dmt/DmtClient.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/dmt/DmtClient.java @@ -199,16 +199,16 @@ public static Response insertDataToDatabase(String host, int port, DatabaseEntry public static Response sendPostRequest(String host, int port, String command, String body) { OkHttpClient client = defaultClient(); - Request request = new Request.Builder() - .url("http://" + host + ":" + port + command) - .post(RequestBody.create(body, MEDIATYPE_JSON)) - .build(); - Call call = client.newCall(request); AtomicReference responseAtomicReference = new AtomicReference<>(); await().atMost(Duration.ofSeconds(HTTP_POLL_TIMEOUT)) .pollInterval(Duration.ofMillis(HTTP_POLL_INTERVAL)) .until(() -> { + Request request = new Request.Builder() + .url("http://" + host + ":" + port + command) + .post(RequestBody.create(body, MEDIATYPE_JSON)) + .build(); + Call call = client.newCall(request); try (Response response = call.execute()) { if (response.isSuccessful()) { responseAtomicReference.set(response); @@ -219,6 +219,7 @@ public static Response sendPostRequest(String host, int port, String command, St } } catch (Exception e) { + LOGGER.error("Cannot send POST request to DMT: {}", e.getMessage()); return false; } }); @@ -243,15 +244,15 @@ public static String sendPostRequest(String host, int port, String command, Map< } HttpUrl url = builder.build(); RequestBody requestBody = RequestBody.create(body, MEDIATYPE_JSON); - Request request = new Request.Builder() - .url(url) - .method("POST", requestBody) - .build(); - Call call = client.newCall(request); AtomicReference responseAtomicReference = new AtomicReference<>(); await().atMost(Duration.ofSeconds(HTTP_POLL_TIMEOUT)) .pollInterval(Duration.ofMillis(HTTP_POLL_INTERVAL)) .until(() -> { + Request request = new Request.Builder() + .url(url) + .method("POST", requestBody) + .build(); + Call call = client.newCall(request); try (Response response = call.execute()) { if (response.isSuccessful()) { responseAtomicReference.set(response.body().string()); @@ -262,6 +263,7 @@ public static String sendPostRequest(String host, int port, String command, Map< } } catch (Exception e) { + LOGGER.error("Cannot send POST request to DMT: {}", e.getMessage()); return false; } }); diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/dmt/DmtResource.java b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/dmt/DmtResource.java index 6b79ab78..85230cc3 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/dmt/DmtResource.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/dmt/DmtResource.java @@ -40,25 +40,25 @@ public void configureAsDefault(String namespace) { @Override public void deploy() { - KubeResourceManager.getInstance().createResourceWithoutWait(this.configMap, this.service); - KubeResourceManager.getInstance().createResourceWithWait(this.deployment); + KubeResourceManager.get().createResourceWithoutWait(this.configMap, this.service); + KubeResourceManager.get().createResourceWithWait(this.deployment); } // TODO: The port should be configurable public LocalPortForward portForward(int port, String namespace) { Map labels = new HashMap<>(); labels.put("app.kubernetes.io/name", "database-manipulation-tool"); - List names = KubeResourceManager.getKubeClient().getClient().pods().inNamespace(namespace) + List names = KubeResourceManager.get().kubeClient().getClient().pods().inNamespace(namespace) .withLabels(labels).list().getItems() .stream().map(p -> p.getMetadata().getName()).toList(); if (names.size() == 1) { try { - int cp = KubeResourceManager.getKubeClient().getClient() + int cp = KubeResourceManager.get().kubeClient().getClient() .pods().inNamespace(namespace) .withName(names.get(0)).get().getSpec().getContainers().get(0).getPorts().get(0).getContainerPort(); - return KubeResourceManager.getKubeClient().getClient() + return KubeResourceManager.get().kubeClient().getClient() .pods().inNamespace(namespace) .withName(names.get(0)).portForward(cp, InetAddress.getByAddress(new byte[]{ 127, 0, 0, 1 }), port); } diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/logs/MustGatherImpl.java b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/logs/MustGatherImpl.java new file mode 100644 index 00000000..44d45a2c --- /dev/null +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/logs/MustGatherImpl.java @@ -0,0 +1,73 @@ +/* + * Copyright Debezium Authors. + * + * Licensed under the Apache Software License version 2.0, available at http://www.apache.org/licenses/LICENSE-2.0 + */ +package io.debezium.operator.systemtests.resources.logs; + +import java.nio.file.Path; +import java.nio.file.Paths; +import java.util.Collections; + +import org.junit.jupiter.api.extension.ExtensionContext; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import io.debezium.operator.systemtests.ConfigProperties; +import io.debezium.operator.systemtests.TestBase; +import io.fabric8.kubernetes.api.model.LabelSelectorBuilder; +import io.skodjob.testframe.LogCollector; +import io.skodjob.testframe.LogCollectorBuilder; +import io.skodjob.testframe.interfaces.MustGatherSupplier; +import io.skodjob.testframe.resources.KubeResourceManager; + +public class MustGatherImpl implements MustGatherSupplier { + private static final Logger logger = LoggerFactory.getLogger(MustGatherImpl.class); + + @Override + public void saveKubernetesState(ExtensionContext extensionContext) { + LogCollector logCollector = new LogCollectorBuilder() + .withNamespacedResources( + "debeziumserver", + "deployment", + "secret", + "configmap", + "role", + "rolebinding", + "serviceaccount", + "pvc", + "statefulset", + "replicaset", + "service", + "route", + "ingress", + "networkpolicy") + .withClusterWideResources( + "node", + "pv") + .withKubeClient(KubeResourceManager.get().kubeClient()) + .withKubeCmdClient(KubeResourceManager.get().kubeCmdClient()) + .withRootFolderPath(getLogPath( + TestBase.LOG_DIR.resolve("failedTest").toString(), extensionContext).toString()) + .build(); + try { + logCollector.collectFromNamespacesWithLabels(new LabelSelectorBuilder() + .withMatchLabels(Collections.singletonMap(ConfigProperties.LOG_COLLECT_LABEL, "true")) + .build()); + } + catch (Exception ignored) { + logger.warn("Failed to collect"); + } + logCollector.collectClusterWideResources(); + } + + private Path getLogPath(String folderName, ExtensionContext context) { + String testMethod = context.getDisplayName(); + String testClassName = context.getTestClass().map(Class::getName).orElse("NOCLASS"); + Path path = TestBase.LOG_DIR.resolve(Paths.get(folderName, testClassName)); + if (testMethod != null) { + path = path.resolve(testMethod.replace("(", "").replace(")", "")); + } + return path; + } +} diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/operator/DebeziumOperatorBundleResource.java b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/operator/DebeziumOperatorBundleResource.java index e71cf27a..217ce7fd 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/operator/DebeziumOperatorBundleResource.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/operator/DebeziumOperatorBundleResource.java @@ -33,7 +33,7 @@ public class DebeziumOperatorBundleResource implements DeployableResourceGroup { @Override public void configureAsDefault(String namespace) { try { - List res = KubeResourceManager.getKubeClient().readResourcesFromFile(Paths.get(BUNDLE_PATH + "kubernetes.yml")); + List res = KubeResourceManager.get().kubeClient().readResourcesFromFile(Paths.get(BUNDLE_PATH + "kubernetes.yml")); for (HasMetadata object : res) { object.getMetadata().setNamespace(namespace); switch (object.getKind()) { @@ -61,7 +61,7 @@ public void configureAsDefault(String namespace) { break; } } - res = KubeResourceManager.getKubeClient().readResourcesFromFile(Paths.get(BUNDLE_PATH + "/debeziumservers.debezium.io-v1.yml")); + res = KubeResourceManager.get().kubeClient().readResourcesFromFile(Paths.get(BUNDLE_PATH + "/debeziumservers.debezium.io-v1.yml")); if (res.size() != 1) { throw new IOException("Specified file cannot be found or is in wrong format!"); } @@ -74,8 +74,8 @@ public void configureAsDefault(String namespace) { @Override public void deploy() { - KubeResourceManager.getInstance().createOrUpdateResourceWithoutWait(crd, serviceAccount, + KubeResourceManager.get().createOrUpdateResourceWithoutWait(crd, serviceAccount, clusterRole, debeziumClusterRoleBinding, viewRoleBinding, service); - KubeResourceManager.getInstance().createOrUpdateResourceWithWait(deployment); + KubeResourceManager.get().createOrUpdateResourceWithWait(deployment); } } diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/server/DebeziumServerResource.java b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/server/DebeziumServerResource.java index e87f2789..7ed4f6d9 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/server/DebeziumServerResource.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/server/DebeziumServerResource.java @@ -5,12 +5,6 @@ */ package io.debezium.operator.systemtests.resources.server; -import static io.debezium.operator.systemtests.ConfigProperties.FABRIC8_POLL_INTERVAL; -import static io.debezium.operator.systemtests.ConfigProperties.FABRIC8_POLL_TIMEOUT; -import static org.awaitility.Awaitility.await; - -import java.io.InputStream; -import java.time.Duration; import java.util.function.Consumer; import org.slf4j.Logger; @@ -20,7 +14,6 @@ import io.fabric8.kubernetes.client.dsl.MixedOperation; import io.fabric8.kubernetes.client.dsl.NonNamespaceOperation; import io.fabric8.kubernetes.client.dsl.Resource; -import io.fabric8.kubernetes.client.dsl.internal.HasMetadataOperationsImpl; import io.skodjob.testframe.interfaces.ResourceType; import io.skodjob.testframe.resources.KubeResourceManager; @@ -30,7 +23,7 @@ public class DebeziumServerResource implements ResourceType { private final Logger logger = LoggerFactory.getLogger(this.getClass().getName()); public DebeziumServerResource() { - this.client = KubeResourceManager.getKubeClient().getClient().resources(DebeziumServer.class, DebeziumServerList.class); + this.client = KubeResourceManager.get().kubeClient().getClient().resources(DebeziumServer.class, DebeziumServerList.class); } public DebeziumServer get(String namespace, String name) { @@ -58,44 +51,27 @@ public void update(DebeziumServer debeziumServer) { } @Override - public void delete(String name) { - client.list().getItems().stream() - .filter(n -> n.getMetadata().getName().equals(name)).findFirst().ifPresent(client::delete); + public void delete(DebeziumServer debeziumServer) { + client.inNamespace(debeziumServer.getMetadata().getNamespace()) + .withName(debeziumServer.getMetadata().getName()).delete(); } @Override - public void replace(String s, Consumer editor) { - DebeziumServer toBeReplaced = client.withName(s).get(); - editor.accept(toBeReplaced); - update(toBeReplaced); + public void replace(DebeziumServer debeziumServer, Consumer consumer) { + DebeziumServer toBeUpdated = client.inNamespace(debeziumServer.getMetadata().getNamespace()) + .withName(debeziumServer.getMetadata().getName()).get(); + consumer.accept(toBeUpdated); + update(toBeUpdated); } @Override - public boolean waitForReadiness(DebeziumServer debeziumServer) { - await().atMost(Duration.ofSeconds(FABRIC8_POLL_TIMEOUT)).pollInterval(Duration.ofSeconds(FABRIC8_POLL_INTERVAL)) - .until(() -> { - DebeziumServer dbzServer = client.inNamespace(debeziumServer.getMetadata().getNamespace()) - .withName(debeziumServer.getMetadata().getName()).get(); - - boolean ready = dbzServer.getStatus().getConditions().stream() - .anyMatch(condition -> condition.getType().equals("Ready") && condition.getStatus().equals("True")); - if (ready) { - return true; - } - else { - logger.info("Waiting for readiness of Debezium Server..."); - return false; - } - }); - return true; + public boolean isReady(DebeziumServer debeziumServer) { + return debeziumServer.getStatus().getConditions().stream() + .anyMatch(condition -> condition.getType().equals("Ready") && condition.getStatus().equals("True")); } @Override - public boolean waitForDeletion(DebeziumServer debeziumServer) { + public boolean isDeleted(DebeziumServer debeziumServer) { return debeziumServer == null; } - - public DebeziumServer loadResource(InputStream is) { - return (DebeziumServer) ((HasMetadataOperationsImpl) this.client.load(is)).getItem(); - } } diff --git a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/sinks/RedisResource.java b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/sinks/RedisResource.java index 323ce7e1..72a4704d 100644 --- a/systemtests/src/test/java/io/debezium/operator/systemtests/resources/sinks/RedisResource.java +++ b/systemtests/src/test/java/io/debezium/operator/systemtests/resources/sinks/RedisResource.java @@ -65,16 +65,16 @@ public void setConfigMap(ConfigMap configMap) { @Override public void configureAsDefault(String namespace) { try { - this.pod = (Pod) KubeResourceManager.getKubeClient().readResourcesFromFile(this.getClass().getClassLoader().getResourceAsStream("redis/redis-pod.yaml")) + this.pod = (Pod) KubeResourceManager.get().kubeClient().readResourcesFromFile(this.getClass().getClassLoader().getResourceAsStream("redis/redis-pod.yaml")) .get(0); pod.getMetadata().setNamespace(namespace); - this.persistentVolumeClaim = (PersistentVolumeClaim) KubeResourceManager.getKubeClient() + this.persistentVolumeClaim = (PersistentVolumeClaim) KubeResourceManager.get().kubeClient() .readResourcesFromFile(this.getClass().getClassLoader().getResourceAsStream("redis/redis-pvc.yaml")).get(0); persistentVolumeClaim.getMetadata().setNamespace(namespace); - this.configMap = (ConfigMap) KubeResourceManager.getKubeClient() + this.configMap = (ConfigMap) KubeResourceManager.get().kubeClient() .readResourcesFromFile(this.getClass().getClassLoader().getResourceAsStream("redis/redis-cfg.yaml")).get(0); configMap.getMetadata().setNamespace(namespace); - this.service = (Service) KubeResourceManager.getKubeClient() + this.service = (Service) KubeResourceManager.get().kubeClient() .readResourcesFromFile(this.getClass().getClassLoader().getResourceAsStream("redis/redis-service.yaml")).get(0); service.getMetadata().setNamespace(namespace); } @@ -89,8 +89,8 @@ public static String getDefaultRedisAddress() { @Override public void deploy() { - KubeResourceManager.getInstance().createResourceWithoutWait(configMap, service, persistentVolumeClaim); - KubeResourceManager.getInstance().createResourceWithWait(pod); + KubeResourceManager.get().createResourceWithoutWait(configMap, service, persistentVolumeClaim); + KubeResourceManager.get().createResourceWithWait(pod); } } diff --git a/systemtests/src/test/resources/redis/redis-pod.yaml b/systemtests/src/test/resources/redis/redis-pod.yaml index 16c94a26..9f7611af 100644 --- a/systemtests/src/test/resources/redis/redis-pod.yaml +++ b/systemtests/src/test/resources/redis/redis-pod.yaml @@ -8,7 +8,7 @@ metadata: spec: containers: - name: redis - image: mirror.gcr.io/library/redis:7.2.4 + image: mirror.gcr.io/library/redis:7.2.8 command: - redis-server - "/redis-master/redis.conf"