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
10 changes: 9 additions & 1 deletion .github/actions/run-systemtests/action.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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/**/*
7 changes: 6 additions & 1 deletion systemtests/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
<junit.platform.launcher.version>1.10.2</junit.platform.launcher.version>
<log4j2.version>2.24.2</log4j2.version>
<assertj.version>3.27.3</assertj.version>
<testframe.version>0.2.1</testframe.version>
<testframe.version>0.13.0</testframe.version>
<dmt.version>0.0.1-alpha1</dmt.version>
<okhttp.version>5.0.0-alpha.12</okhttp.version>
<awaitility.version>4.2.1</awaitility.version>
Expand Down Expand Up @@ -97,6 +97,11 @@
<artifactId>test-frame-kubernetes</artifactId>
<version>${testframe.version}</version>
</dependency>
<dependency>
<groupId>io.skodjob</groupId>
<artifactId>test-frame-log-collector</artifactId>
<version>${testframe.version}</version>
</dependency>
<dependency>
<groupId>io.debezium</groupId>
<artifactId>debezium-operator-api</artifactId>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand All @@ -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")
Expand All @@ -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")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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)) {
Expand All @@ -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)) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}

Expand All @@ -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();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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)
Expand All @@ -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() {
Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ public void createNewNamespace() {
.endMetadata()
.build();
this.currentNamespace = name;
KubeResourceManager.getInstance().createResourceWithWait(namespace);
KubeResourceManager.get().createResourceWithWait(namespace);
}

public String getCurrentNamespace() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
});
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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<Response> 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);
Expand All @@ -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;
}
});
Expand All @@ -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<String> 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());
Expand All @@ -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;
}
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<String, String> labels = new HashMap<>();
labels.put("app.kubernetes.io/name", "database-manipulation-tool");
List<String> names = KubeResourceManager.getKubeClient().getClient().pods().inNamespace(namespace)
List<String> 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);
}
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
}
}
Loading