[Feature]Add zookeeper e2e code (#3030)

Co-authored-by: Calvin <naruse_shinji@163.com>
Co-authored-by: shown <yuluo08290126@gmail.com>
This commit is contained in:
Jast
2025-02-07 23:22:10 +08:00
committed by GitHub
co-authored by Calvin shown
parent 05e42182d9
commit 7effbc1f0d
5 changed files with 188 additions and 11 deletions
@@ -53,5 +53,10 @@
<version>${hertzbeat.version}</version> <version>${hertzbeat.version}</version>
<scope>test</scope> <scope>test</scope>
</dependency> </dependency>
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>testcontainers</artifactId>
<scope>test</scope>
</dependency>
</dependencies> </dependencies>
</project> </project>
@@ -56,6 +56,7 @@ public class HttpMonitorE2eTest extends AbstractCollectE2eTest {
private static final int MOCK_SERVER_PORT = 52376; private static final int MOCK_SERVER_PORT = 52376;
private static final String LOCALHOST = "127.0.0.1"; private static final String LOCALHOST = "127.0.0.1";
private static final String RELATIVE_PATH = "/"; private static final String RELATIVE_PATH = "/";
private static final List<String> ALLOW_EMPTY_WHITE_LIST = List.of("header");
private static HttpServer mockServer; private static HttpServer mockServer;
@AfterAll @AfterAll
@@ -96,7 +97,12 @@ public class HttpMonitorE2eTest extends AbstractCollectE2eTest {
List<Map<String, Configmap>> configmapFromPreCollectData = new LinkedList<>(); List<Map<String, Configmap>> configmapFromPreCollectData = new LinkedList<>();
for (Metrics metricsDef : dockerJob.getMetrics()) { for (Metrics metricsDef : dockerJob.getMetrics()) {
metricsDef = CollectUtil.replaceCryPlaceholderToMetrics(metricsDef, configmapFromPreCollectData.size() > 0 ? configmapFromPreCollectData.get(0) : new HashMap<>()); metricsDef = CollectUtil.replaceCryPlaceholderToMetrics(metricsDef, configmapFromPreCollectData.size() > 0 ? configmapFromPreCollectData.get(0) : new HashMap<>());
CollectRep.MetricsData metricsData = validateMetricsCollection(metricsDef, metricsDef.getName()); CollectRep.MetricsData metricsData;
if (ALLOW_EMPTY_WHITE_LIST.contains(metricsDef.getName())) {
metricsData = validateMetricsCollection(metricsDef, metricsDef.getName(), true);
} else {
metricsData = validateMetricsCollection(metricsDef, metricsDef.getName());
}
configmapFromPreCollectData = CollectUtil.getConfigmapFromPreCollectData(metricsData); configmapFromPreCollectData = CollectUtil.getConfigmapFromPreCollectData(metricsData);
} }
} }
@@ -38,6 +38,8 @@ import org.testcontainers.lifecycle.Startables;
import org.testcontainers.utility.DockerImageName; import org.testcontainers.utility.DockerImageName;
import java.time.Duration; import java.time.Duration;
import java.util.Arrays;
import java.util.List;
import java.util.Random; import java.util.Random;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeoutException; import java.util.concurrent.TimeoutException;
@@ -54,6 +56,7 @@ public class SshCollectE2eTest extends AbstractCollectE2eTest {
private static final String ROOT_USER = "root"; private static final String ROOT_USER = "root";
private static final int SSH_PORT = 22; private static final int SSH_PORT = 22;
private static final int PASSWORD_LENGTH = 12; private static final int PASSWORD_LENGTH = 12;
private static final List<String> ALLOW_EMPTY_WHITE_LIST = Arrays.asList("top_mem_process", "top_cpu_process");
private static GenericContainer<?> linuxContainer; private static GenericContainer<?> linuxContainer;
@@ -88,8 +91,13 @@ public class SshCollectE2eTest extends AbstractCollectE2eTest {
Assertions.assertTrue(linuxContainer.isRunning(), "Ubuntu container should be running"); Assertions.assertTrue(linuxContainer.isRunning(), "Ubuntu container should be running");
Job ubuntuJob = appService.getAppDefine("ubuntu"); Job ubuntuJob = appService.getAppDefine("ubuntu");
ubuntuJob.getMetrics().forEach(metricsDef -> ubuntuJob.getMetrics().forEach(metricsDef -> {
validateMetricsCollection(metricsDef, metricsDef.getName())); if (ALLOW_EMPTY_WHITE_LIST.contains(metricsDef.getName())) {
validateMetricsCollection(metricsDef, metricsDef.getName(), true);
} else {
validateMetricsCollection(metricsDef, metricsDef.getName());
}
});
} }
@Override @Override
@@ -0,0 +1,128 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.apache.hertzbeat.collector.collect.basic.telnet;
import lombok.extern.slf4j.Slf4j;
import org.apache.hertzbeat.collector.collect.AbstractCollectE2eTest;
import org.apache.hertzbeat.collector.collect.telnet.TelnetCollectImpl;
import org.apache.hertzbeat.collector.util.CollectUtil;
import org.apache.hertzbeat.common.entity.job.Configmap;
import org.apache.hertzbeat.common.entity.job.Job;
import org.apache.hertzbeat.common.entity.job.Metrics;
import org.apache.hertzbeat.common.entity.job.protocol.Protocol;
import org.apache.hertzbeat.common.entity.job.protocol.TelnetProtocol;
import org.apache.hertzbeat.common.entity.message.CollectRep;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.junit.jupiter.MockitoExtension;
import org.testcontainers.containers.GenericContainer;
import org.testcontainers.containers.wait.strategy.Wait;
import org.testcontainers.utility.DockerImageName;
import java.time.Duration;
import java.util.HashMap;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
/**
* Integration test for Zookeeper monitoring functionality
*/
@Slf4j
@ExtendWith(MockitoExtension.class)
public class ZookeeperMonitorE2eTest extends AbstractCollectE2eTest {
private static final String ZOOKEEPER_IMAGE_NAME = "zookeeper:3.8.4";
private static final String ZOOKEEPER_NAME = "zookeeper";
private static final Integer ZOOKEEPER_PORT = 2181;
private static GenericContainer<?> zookeeperContainer;
@AfterAll
public static void tearDown() {
if (zookeeperContainer != null) {
zookeeperContainer.stop();
}
}
@BeforeEach
public void setUp() throws Exception {
super.setUp();
collect = new TelnetCollectImpl();
metrics = new Metrics();
try {
// Start Zookeeper container with custom configuration
zookeeperContainer = new GenericContainer<>(DockerImageName.parse(ZOOKEEPER_IMAGE_NAME))
.withExposedPorts(ZOOKEEPER_PORT)
.withEnv("ZOO_4LW_COMMANDS_WHITELIST", "*")
.withNetworkAliases(ZOOKEEPER_NAME)
.waitingFor(
Wait.forLogMessage(".*Started AdminServer on address.*\\n", 1)
.withStartupTimeout(Duration.ofSeconds(60))
)
.withLogConsumer(outputFrame -> {
log.info(outputFrame.getUtf8String());
});
zookeeperContainer.start();
log.info("Zookeeper container started at {}:{}",
zookeeperContainer.getHost(),
zookeeperContainer.getMappedPort(ZOOKEEPER_PORT));
} catch (Exception e) {
e.printStackTrace();
log.error("Failed to start Zookeeper container", e);
throw e;
}
Thread.sleep(30000);
}
@Override
protected CollectRep.MetricsData.Builder collectMetrics(Metrics metricsDef) {
TelnetProtocol telnetProtocol = (TelnetProtocol) buildProtocol(metricsDef);
metrics.setTelnet(telnetProtocol);
CollectRep.MetricsData.Builder metricsData = CollectRep.MetricsData.newBuilder();
metricsData.setApp(ZOOKEEPER_NAME);
metrics.setAliasFields(metricsDef.getAliasFields());
return collectMetricsData(metrics, metricsDef, metricsData);
}
@Override
protected Protocol buildProtocol(Metrics metricsDef) {
TelnetProtocol protocol = new TelnetProtocol();
protocol.setHost(zookeeperContainer.getHost());
protocol.setPort(String.valueOf(zookeeperContainer.getMappedPort(ZOOKEEPER_PORT)));
protocol.setCmd(metricsDef.getTelnet().getCmd());
return protocol;
}
@Test
public void testZookeeperMonitor() {
Assertions.assertTrue(zookeeperContainer.isRunning(), "Zookeeper container should be running");
Job dockerJob = appService.getAppDefine("zookeeper");
List<Map<String, Configmap>> configmapFromPreCollectData = new LinkedList<>();
for (Metrics metricsDef : dockerJob.getMetrics()) {
metricsDef = CollectUtil.replaceCryPlaceholderToMetrics(metricsDef, configmapFromPreCollectData.size() > 0 ? configmapFromPreCollectData.get(0) : new HashMap<>());
CollectRep.MetricsData metricsData = validateMetricsCollection(metricsDef, metricsDef.getName());
configmapFromPreCollectData = CollectUtil.getConfigmapFromPreCollectData(metricsData);
}
}
}
@@ -22,6 +22,7 @@ import org.apache.hertzbeat.collector.dispatch.CollectDataDispatch;
import org.apache.hertzbeat.collector.dispatch.MetricsCollect; import org.apache.hertzbeat.collector.dispatch.MetricsCollect;
import org.apache.hertzbeat.collector.dispatch.timer.Timeout; import org.apache.hertzbeat.collector.dispatch.timer.Timeout;
import org.apache.hertzbeat.collector.dispatch.timer.WheelTimerTask; import org.apache.hertzbeat.collector.dispatch.timer.WheelTimerTask;
import org.apache.hertzbeat.common.constants.CommonConstants;
import org.apache.hertzbeat.common.entity.job.Job; import org.apache.hertzbeat.common.entity.job.Job;
import org.apache.hertzbeat.common.entity.job.Metrics; import org.apache.hertzbeat.common.entity.job.Metrics;
import org.apache.hertzbeat.common.entity.job.protocol.Protocol; import org.apache.hertzbeat.common.entity.job.protocol.Protocol;
@@ -77,20 +78,39 @@ public abstract class AbstractCollectE2eTest {
/** /**
* Validate metrics collection, check if the metrics values are not empty <br/> * Validate metrics collection, check if the metrics values are not empty <br/>
* We believe that all monitoring metrics should have data * @param metricsDef metrics definition
* @param metricName metric name
* @return metrics data
*/ */
protected CollectRep.MetricsData validateMetricsCollection(Metrics metricsDef, String metricName) { protected CollectRep.MetricsData validateMetricsCollection(Metrics metricsDef, String metricName) {
// By default, we do not allow empty values
return validateMetricsCollection(metricsDef, metricName, false);
}
/**
* Validate metrics collection, check if the metrics values are not empty <br/>
* We believe that all monitoring metrics should have data
*
* @param metricsDef metrics definition
* @param metricName metric name
* @param allowEmpty In some special scenarios, it is not necessary to check if the value is `&nbsp;`
*/
protected CollectRep.MetricsData validateMetricsCollection(Metrics metricsDef, String metricName, boolean allowEmpty) {
CollectRep.MetricsData.Builder metricsData = collectMetrics(metricsDef); CollectRep.MetricsData.Builder metricsData = collectMetrics(metricsDef);
metricsCollect.calculateFields(metricsDef, metricsData); metricsCollect.calculateFields(metricsDef, metricsData);
Assertions.assertTrue(metricsData.getValuesList().size() > 0, Assertions.assertTrue(metricsData.getValuesList().size() > 0,
String.format("%s metrics values should not be empty", metricName)); String.format("%s metrics values should not be empty, detail: %s", metricName, metricsData.getMsg()));
for (CollectRep.ValueRow valueRow : metricsData.getValuesList()) { for (CollectRep.ValueRow valueRow : metricsData.getValuesList()) {
for (int i = 0; i < valueRow.getColumnsCount(); i++) { for (int i = 0; i < valueRow.getColumnsCount(); i++) {
Assertions.assertFalse(valueRow.getColumns(i).isEmpty(), Assertions.assertFalse(valueRow.getColumns(i).isEmpty(),
String.format("%s metric column %d should not be empty", metricName, i)); String.format("%s metric column %d should not be empty", metricName, i));
if (!allowEmpty) {
// Check if the value is not null
Assertions.assertNotEquals(CommonConstants.NULL_VALUE, valueRow.getColumns(i), String.format("%s metric column %d should not be null", metricName, i));
}
} }
} }
@@ -98,21 +118,31 @@ public abstract class AbstractCollectE2eTest {
return metricsData.build(); return metricsData.build();
} }
/**
* Set alias fields for metrics
*
* @param metrics metrics
* @param metricsDef metrics definition
*/
protected void setMetricsAliasFields(Metrics metrics, Metrics metricsDef) { protected void setMetricsAliasFields(Metrics metrics, Metrics metricsDef) {
metrics.setAliasFields(metricsDef.getAliasFields() == null List<String> aliasFields = metricsDef.getAliasFields() == null
? metricsDef.getFields().stream() ? metricsDef.getFields().stream().map(Metrics.Field::getField).collect(Collectors.toList())
.map(Metrics.Field::getField) : metricsDef.getAliasFields();
.collect(Collectors.toList()) : metrics.setAliasFields(aliasFields);
metricsDef.getAliasFields()); metricsDef.setAliasFields(aliasFields);
} }
protected abstract CollectRep.MetricsData.Builder collectMetrics(Metrics metricsDef); protected abstract CollectRep.MetricsData.Builder collectMetrics(Metrics metricsDef);
protected CollectRep.MetricsData.Builder collectMetricsData(Metrics metrics, Metrics metricsDef) { protected CollectRep.MetricsData.Builder collectMetricsData(Metrics metrics, Metrics metricsDef) {
CollectRep.MetricsData.Builder metricsData = CollectRep.MetricsData.newBuilder();
return this.collectMetricsData(metrics, metricsDef, metricsData);
}
protected CollectRep.MetricsData.Builder collectMetricsData(Metrics metrics, Metrics metricsDef, CollectRep.MetricsData.Builder metricsData) {
setMetricsAliasFields(metrics, metricsDef); setMetricsAliasFields(metrics, metricsDef);
// Collect metrics // Collect metrics
CollectRep.MetricsData.Builder metricsData = CollectRep.MetricsData.newBuilder();
collect.collect(metricsData, metrics); collect.collect(metricsData, metrics);
return metricsData; return metricsData;
} }