Skip to content
Merged
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 @@ -67,14 +67,20 @@ public void collect(CollectRep.MetricsData.Builder builder, Metrics metrics) {
long responseTime = System.currentTimeMillis() - startTime;
List<String> aliasFields = metrics.getAliasFields();
String app = builder.getApp();
Map<String, String> resultMap = execCmdAndParseResult(telnetClient, telnet.getCmd(), app);
resultMap.put(CollectorConstants.RESPONSE_TIME, Long.toString(responseTime));
if (resultMap.size() < aliasFields.size()) {
log.error("telnet response data not enough: {}", resultMap);
CmdResult cmdResult = execCmdAndParseResult(telnetClient, telnet.getCmd(), app);
Map<String, String> resultMap = cmdResult.values();
boolean expectsCmdMetrics = StringUtils.isNotBlank(telnet.getCmd())
&& aliasFields.stream().anyMatch(field -> !CollectorConstants.RESPONSE_TIME.equalsIgnoreCase(field));
boolean hasExpectedMetric = aliasFields.stream().anyMatch(resultMap::containsKey);
if (expectsCmdMetrics && !hasExpectedMetric) {
// e.g. zookeeper refusing a 4lw command not in its 4lw.commands.whitelist
String reply = sanitizeReply(cmdResult.rawResponse());
log.warn("telnet cmd [{}] returned no expected metrics: {}", telnet.getCmd(), reply);
builder.setCode(CollectRep.Code.FAIL);
builder.setMsg("The cmd execution results do not match the expected number of metrics.");
builder.setMsg("Cmd [" + telnet.getCmd() + "] returned no expected metrics. Response: " + reply);
return;
}
resultMap.put(CollectorConstants.RESPONSE_TIME, Long.toString(responseTime));
CollectRep.ValueRow.Builder valueRowBuilder = CollectRep.ValueRow.newBuilder();
for (String field : aliasFields) {
String fieldValue = resultMap.get(field);
Expand Down Expand Up @@ -118,30 +124,36 @@ public String supportProtocol() {
return DispatchConstants.PROTOCOL_TELNET;
}

private static Map<String, String> execCmdAndParseResult(TelnetClient telnetClient, String cmd, String app) throws IOException {
record CmdResult(Map<String, String> values, String rawResponse) {
}

private static String sanitizeReply(String raw) {
return StringUtils.abbreviate(raw.trim().replaceAll("[\\p{Cntrl}]+", " "), 300);
}

private static CmdResult execCmdAndParseResult(TelnetClient telnetClient, String cmd, String app) throws IOException {
if (cmd == null || StringUtils.isEmpty(cmd.trim())) {
return new HashMap<>(16);
return new CmdResult(new HashMap<>(16), "");
}
OutputStream outputStream = telnetClient.getOutputStream();
outputStream.write(cmd.getBytes(StandardCharsets.UTF_8));
outputStream.flush();
String result = new String(telnetClient.getInputStream().readAllBytes());
String[] lines = result.split("\n");
if (CollectorConstants.ZOOKEEPER_APP.equals(app) && CollectorConstants.ZOOKEEPER_ENVI_HEAD.equals(lines[0])) {
if (lines.length > 0 && CollectorConstants.ZOOKEEPER_APP.equals(app)
&& CollectorConstants.ZOOKEEPER_ENVI_HEAD.equals(lines[0])) {
lines = Arrays.stream(lines)
.skip(1)
.toArray(String[]::new);
}
boolean contains = lines[0].contains("=");
return Arrays.stream(lines)
.map(item -> {
if (contains) {
return item.split("=");
} else {
return item.split("\t");
}
})
if (lines.length == 0) {
return new CmdResult(new HashMap<>(16), result);
}
String separator = lines[0].contains("=") ? "=" : "\t";
Map<String, String> values = Arrays.stream(lines)
.map(item -> item.split(separator, 2))
.filter(item -> item.length == 2)
.collect(Collectors.toMap(x -> x[0], x -> x[1]));
.collect(Collectors.toMap(x -> x[0], x -> x[1], (first, second) -> first, HashMap::new));
return new CmdResult(values, result);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;

import java.io.ByteArrayInputStream;
import java.io.InputStream;
Expand All @@ -30,6 +31,7 @@
import java.util.List;
import org.apache.commons.net.telnet.TelnetClient;
import org.apache.hertzbeat.collector.dispatch.DispatchConstants;
import org.apache.hertzbeat.common.constants.CommonConstants;
import org.apache.hertzbeat.common.entity.job.Metrics;
import org.apache.hertzbeat.common.entity.job.protocol.TelnetProtocol;
import org.apache.hertzbeat.common.entity.message.CollectRep;
Expand Down Expand Up @@ -155,6 +157,114 @@ void testCollectWithTab() {
mocked.close();
}

@Test
void testCollectPadsMissingMetrics() {
CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder();
Metrics metrics = telnetMetrics("mntr", List.of("responseTime", "a", "b", "c"));
try (MockedConstruction<TelnetClient> mocked = mockTelnetReply("a=1")) {
telnetCollect.collect(builder, metrics);
}
assertEquals(1, builder.getValuesCount());
assertEquals("1", builder.getValues(0).getColumns(1));
assertEquals(CommonConstants.NULL_VALUE, builder.getValues(0).getColumns(2));
assertEquals(CommonConstants.NULL_VALUE, builder.getValues(0).getColumns(3));
}

@Test
void testCollectFailsWithRawReplyWhenNothingParsed() {
CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder();
Metrics metrics = telnetMetrics("conf", List.of("responseTime", "a"));
try (MockedConstruction<TelnetClient> mocked =
mockTelnetReply("conf is not executed because it is not in the whitelist.")) {
telnetCollect.collect(builder, metrics);
}
assertEquals(CollectRep.Code.FAIL, builder.getCode());
assertTrue(builder.getMsg().contains("not in the whitelist"));
}

@Test
void testCollectKeepsFirstOnDuplicateKeys() {
CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder();
Metrics metrics = telnetMetrics("mntr", List.of("a", "b"));
try (MockedConstruction<TelnetClient> mocked = mockTelnetReply("a=1\na=2\nb=3")) {
telnetCollect.collect(builder, metrics);
}
assertEquals(1, builder.getValuesCount());
assertEquals("1", builder.getValues(0).getColumns(0));
assertEquals("3", builder.getValues(0).getColumns(1));
}

@Test
void testCollectFailsWhenReplyHasOnlyUnrelatedPairs() {
CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder();
Metrics metrics = telnetMetrics("mntr", List.of("responseTime", "a"));
try (MockedConstruction<TelnetClient> mocked = mockTelnetReply("error=conf disabled")) {
telnetCollect.collect(builder, metrics);
}
assertEquals(CollectRep.Code.FAIL, builder.getCode());
}

@Test
void testCollectSurvivesHeaderOnlyEnviReply() {
CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder().setApp("zookeeper");
Metrics metrics = telnetMetrics("envi", List.of("responseTime", "a"));
try (MockedConstruction<TelnetClient> mocked = mockTelnetReply("Environment:")) {
telnetCollect.collect(builder, metrics);
}
assertEquals(CollectRep.Code.FAIL, builder.getCode());
}

@Test
void testCollectFailsCleanlyOnNewlineOnlyReply() {
CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder().setApp("zookeeper");
Metrics metrics = telnetMetrics("conf", List.of("responseTime", "a"));
try (MockedConstruction<TelnetClient> mocked = mockTelnetReply("\n")) {
telnetCollect.collect(builder, metrics);
}
assertEquals(CollectRep.Code.FAIL, builder.getCode());
assertTrue(builder.getMsg().contains("returned no expected metrics"));
}

@Test
void testCollectKeepsValueContainingSeparator() {
CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder();
Metrics metrics = telnetMetrics("conf", List.of("dataDir", "secureClientPort"));
try (MockedConstruction<TelnetClient> mocked = mockTelnetReply("dataDir=/data/zk=a\nsecureClientPort=")) {
telnetCollect.collect(builder, metrics);
}
assertEquals(1, builder.getValuesCount());
assertEquals("/data/zk=a", builder.getValues(0).getColumns(0));
assertEquals("", builder.getValues(0).getColumns(1));
}

@Test
void testCollectBlankCmdKeepsResponseTimeRow() {
CollectRep.MetricsData.Builder builder = CollectRep.MetricsData.newBuilder();
Metrics metrics = telnetMetrics("", List.of("responseTime"));
try (MockedConstruction<TelnetClient> mocked =
Mockito.mockConstruction(TelnetClient.class, (telnetClient, context) ->
Mockito.when(telnetClient.isConnected()).thenReturn(true))) {
telnetCollect.collect(builder, metrics);
}
assertEquals(1, builder.getValuesCount());
}

private static Metrics telnetMetrics(String cmd, List<String> aliasFields) {
Metrics metrics = new Metrics();
metrics.setTelnet(TelnetProtocol.builder().timeout("10").port("2181").cmd(cmd).build());
metrics.setAliasFields(aliasFields);
return metrics;
}

private static MockedConstruction<TelnetClient> mockTelnetReply(String reply) {
InputStream inputStream = new ByteArrayInputStream(reply.getBytes(StandardCharsets.UTF_8));
return Mockito.mockConstruction(TelnetClient.class, (telnetClient, context) -> {
Mockito.when(telnetClient.isConnected()).thenReturn(true);
Mockito.when(telnetClient.getOutputStream()).thenReturn(Mockito.mock(OutputStream.class));
Mockito.when(telnetClient.getInputStream()).thenReturn(inputStream);
});
}

@Test
void preCheck() throws IllegalArgumentException {
// metrics is null
Expand Down
Loading