Browse Source

Rule Engine stats: exception message truncation

pull/9353/head
ViacheslavKlimov 3 years ago
committed by Andrii Shvaika
parent
commit
12d2c26279
  1. 6
      application/src/main/java/org/thingsboard/server/service/stats/DefaultRuleEngineStatisticsService.java
  2. 1
      application/src/main/resources/thingsboard.yml
  3. 47
      application/src/test/java/org/thingsboard/server/controller/BaseQueueControllerTest.java
  4. 13
      common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java
  5. 11
      common/data/src/test/java/org/thingsboard/server/common/data/StringUtilsTest.java
  6. 16
      common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleEngineException.java
  7. 4
      common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleNodeException.java
  8. 27
      common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java
  9. 8
      dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java

6
application/src/main/java/org/thingsboard/server/service/stats/DefaultRuleEngineStatisticsService.java

@ -19,6 +19,7 @@ import com.google.common.util.concurrent.FutureCallback;
import lombok.Data; import lombok.Data;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetId;
@ -72,6 +73,9 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS
private final Lock lock = new ReentrantLock(); private final Lock lock = new ReentrantLock();
private final ConcurrentMap<TenantQueueKey, AssetId> tenantQueueAssets = new ConcurrentHashMap<>(); private final ConcurrentMap<TenantQueueKey, AssetId> tenantQueueAssets = new ConcurrentHashMap<>();
@Value("${queue.rule-engine.stats.max-error-message-length:4096}")
private int maxErrorMessageLength;
@Override @Override
public void reportQueueStats(long ts, TbRuleEngineConsumerStats ruleEngineStats) { public void reportQueueStats(long ts, TbRuleEngineConsumerStats ruleEngineStats) {
String queueName = ruleEngineStats.getQueueName(); String queueName = ruleEngineStats.getQueueName();
@ -97,7 +101,7 @@ public class DefaultRuleEngineStatisticsService implements RuleEngineStatisticsS
}); });
ruleEngineStats.getTenantExceptions().forEach((tenantId, e) -> { ruleEngineStats.getTenantExceptions().forEach((tenantId, e) -> {
try { try {
TsKvEntry tsKv = new BasicTsKvEntry(e.getTs(), new JsonDataEntry(RULE_ENGINE_EXCEPTION, e.toJsonString())); TsKvEntry tsKv = new BasicTsKvEntry(e.getTs(), new JsonDataEntry(RULE_ENGINE_EXCEPTION, e.toJsonString(maxErrorMessageLength)));
long ttl = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getRuleEngineExceptionsTtlDays); long ttl = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getRuleEngineExceptionsTtlDays);
ttl = TimeUnit.DAYS.toSeconds(ttl); ttl = TimeUnit.DAYS.toSeconds(ttl);
tsService.saveAndNotifyInternal(tenantId, getServiceAssetId(tenantId, queueName), Collections.singletonList(tsKv), ttl, CALLBACK); tsService.saveAndNotifyInternal(tenantId, getServiceAssetId(tenantId, queueName), Collections.singletonList(tsKv), ttl, CALLBACK);

1
application/src/main/resources/thingsboard.yml

@ -1248,6 +1248,7 @@ queue:
stats: stats:
enabled: "${TB_QUEUE_RULE_ENGINE_STATS_ENABLED:true}" enabled: "${TB_QUEUE_RULE_ENGINE_STATS_ENABLED:true}"
print-interval-ms: "${TB_QUEUE_RULE_ENGINE_STATS_PRINT_INTERVAL_MS:60000}" print-interval-ms: "${TB_QUEUE_RULE_ENGINE_STATS_PRINT_INTERVAL_MS:60000}"
max-error-message-length: "${TB_QUEUE_RULE_ENGINE_MAX_ERROR_MESSAGE_LENGTH:4096}"
queues: queues:
- name: "${TB_QUEUE_RE_MAIN_QUEUE_NAME:Main}" - name: "${TB_QUEUE_RE_MAIN_QUEUE_NAME:Main}"
topic: "${TB_QUEUE_RE_MAIN_TOPIC:tb_rule_engine.main}" topic: "${TB_QUEUE_RE_MAIN_TOPIC:tb_rule_engine.main}"

47
application/src/test/java/org/thingsboard/server/controller/BaseQueueControllerTest.java

@ -16,15 +16,19 @@
package org.thingsboard.server.controller; package org.thingsboard.server.controller;
import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.core.type.TypeReference;
import org.apache.commons.lang3.RandomStringUtils;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Test; import org.junit.Test;
import org.mockito.ArgumentCaptor; import org.mockito.ArgumentCaptor;
import org.mockito.Mockito; import org.mockito.Mockito;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.mock.mockito.SpyBean; import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.TestPropertySource;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.queue.ProcessingStrategy; import org.thingsboard.server.common.data.queue.ProcessingStrategy;
@ -48,10 +52,13 @@ import java.util.Map;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import java.util.stream.Stream; import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.argThat; import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verify;
@ -60,6 +67,9 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.
import static org.thingsboard.server.dao.asset.BaseAssetService.TB_SERVICE_QUEUE; import static org.thingsboard.server.dao.asset.BaseAssetService.TB_SERVICE_QUEUE;
@DaoSqlTest @DaoSqlTest
@TestPropertySource(properties = {
"queue.rule-engine.stats.max-error-message-length=100"
})
public class BaseQueueControllerTest extends AbstractControllerTest { public class BaseQueueControllerTest extends AbstractControllerTest {
@Autowired @Autowired
@ -183,9 +193,44 @@ public class BaseQueueControllerTest extends AbstractControllerTest {
verify(timeseriesDao).save(eq(tenantId), eq(serviceAsset.getId()), argThat(tsKvEntry -> { verify(timeseriesDao).save(eq(tenantId), eq(serviceAsset.getId()), argThat(tsKvEntry -> {
return tsKvEntry.getKey().equals(DefaultRuleEngineStatisticsService.RULE_ENGINE_EXCEPTION) && return tsKvEntry.getKey().equals(DefaultRuleEngineStatisticsService.RULE_ENGINE_EXCEPTION) &&
tsKvEntry.getJsonValue().get().equals(ruleEngineException.toJsonString()); tsKvEntry.getJsonValue().get().equals(ruleEngineException.toJsonString(0));
}), ttlCaptor.capture()); }), ttlCaptor.capture());
assertThat(ttlCaptor.getValue()).isEqualTo(TimeUnit.DAYS.toSeconds(ruleEngineExceptionsTtlDays)); assertThat(ttlCaptor.getValue()).isEqualTo(TimeUnit.DAYS.toSeconds(ruleEngineExceptionsTtlDays));
} }
@Test
public void testRuleEngineExceptionTruncation() {
Queue queue = new Queue();
queue.setName("Test-2");
queue.setTenantId(TenantId.SYS_TENANT_ID);
TbRuleEngineProcessingResult testProcessingResult = Mockito.mock(TbRuleEngineProcessingResult.class);
when(testProcessingResult.getSuccessMap()).thenReturn(new ConcurrentHashMap<>());
when(testProcessingResult.getFailedMap()).thenReturn(new ConcurrentHashMap<>());
when(testProcessingResult.getPendingMap()).thenReturn(new ConcurrentHashMap<>());
String largeExceptionMessage = RandomStringUtils.randomAlphabetic(150);
RuleEngineException ruleEngineException = new RuleEngineException(largeExceptionMessage);
when(testProcessingResult.getExceptionsMap()).thenReturn(new ConcurrentHashMap<>(Map.of(
tenantId, ruleEngineException
)));
TbRuleEngineConsumerStats testStats = new TbRuleEngineConsumerStats(queue, statsFactory);
testStats.log(testProcessingResult, true);
ruleEngineStatisticsService.reportQueueStats(System.currentTimeMillis(), testStats);
AtomicReference<TsKvEntry> reExceptionTsKvEntryCaptor = new AtomicReference<>();
verify(timeseriesDao).save(eq(tenantId), any(), argThat(tsKvEntry -> {
if (tsKvEntry.getKey().equals(DefaultRuleEngineStatisticsService.RULE_ENGINE_EXCEPTION)) {
reExceptionTsKvEntryCaptor.set(tsKvEntry);
return true;
}
return false;
}), anyLong());
TsKvEntry reExceptionTsKvEntry = reExceptionTsKvEntryCaptor.get();
String finalErrorMessage = JacksonUtil.toJsonNode(reExceptionTsKvEntry.getJsonValue().get()).get("message").asText();
assertThat(finalErrorMessage).isEqualTo(largeExceptionMessage.substring(0, 100) + "...[truncated 50 symbols]");
}
} }

13
common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java

@ -20,6 +20,7 @@ import org.apache.commons.lang3.RandomStringUtils;
import java.security.SecureRandom; import java.security.SecureRandom;
import java.util.Base64; import java.util.Base64;
import java.util.function.Function;
import static org.apache.commons.lang3.StringUtils.repeat; import static org.apache.commons.lang3.StringUtils.repeat;
@ -228,4 +229,16 @@ public class StringUtils {
return generateSafeToken(DEFAULT_TOKEN_LENGTH); return generateSafeToken(DEFAULT_TOKEN_LENGTH);
} }
public static String truncate(String string, int maxLength) {
return truncate(string, maxLength, n -> "...[truncated " + n + " symbols]");
}
public static String truncate(String string, int maxLength, Function<Integer, String> truncationMarkerFunc) {
if (string == null || maxLength <= 0 || string.length() <= maxLength) {
return string;
}
int truncatedSymbols = string.length() - maxLength;
return string.substring(0, maxLength) + truncationMarkerFunc.apply(truncatedSymbols);
}
} }

11
common/data/src/test/java/org/thingsboard/server/common/data/StringUtilsTest.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.common.data; package org.thingsboard.server.common.data;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.ValueSource; import org.junit.jupiter.params.provider.ValueSource;
@ -37,4 +38,14 @@ class StringUtilsTest {
assertThat(StringUtils.contains0x00(sample)).isFalse(); assertThat(StringUtils.contains0x00(sample)).isFalse();
} }
@Test
void testTruncate() {
int maxLength = 5;
assertThat(StringUtils.truncate(null, maxLength)).isNull();
assertThat(StringUtils.truncate("", maxLength)).isEmpty();
assertThat(StringUtils.truncate("123", maxLength)).isEqualTo("123");
assertThat(StringUtils.truncate("1234567", maxLength)).isEqualTo("12345...[truncated 2 symbols]");
assertThat(StringUtils.truncate("1234567", 0)).isEqualTo("1234567");
}
} }

16
common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleEngineException.java

@ -19,17 +19,17 @@ import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.Getter; import lombok.Getter;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.StringUtils;
@Slf4j @Slf4j
public class RuleEngineException extends Exception { public class RuleEngineException extends Exception {
protected static final ObjectMapper mapper = new ObjectMapper(); protected static final ObjectMapper mapper = new ObjectMapper();
@Getter @Getter
private long ts; private final long ts;
public RuleEngineException(String message) { public RuleEngineException(String message) {
super(message != null ? message : "Unknown"); this(message, null);
ts = System.currentTimeMillis();
} }
public RuleEngineException(String message, Throwable t) { public RuleEngineException(String message, Throwable t) {
@ -37,12 +37,18 @@ public class RuleEngineException extends Exception {
ts = System.currentTimeMillis(); ts = System.currentTimeMillis();
} }
public String toJsonString() { public String toJsonString(int maxMessageLength) {
try { try {
return mapper.writeValueAsString(mapper.createObjectNode().put("message", getMessage())); return mapper.writeValueAsString(mapper.createObjectNode()
.put("message", truncateIfNecessary(getMessage(), maxMessageLength)));
} catch (JsonProcessingException e) { } catch (JsonProcessingException e) {
log.warn("Failed to serialize exception ", e); log.warn("Failed to serialize exception ", e);
throw new RuntimeException(e); throw new RuntimeException(e);
} }
} }
protected String truncateIfNecessary(String message, int maxMessageLength) {
return StringUtils.truncate(message, maxMessageLength);
}
} }

4
common/message/src/main/java/org/thingsboard/server/common/msg/queue/RuleNodeException.java

@ -52,14 +52,14 @@ public class RuleNodeException extends RuleEngineException {
} }
} }
public String toJsonString() { public String toJsonString(int maxMessageLength) {
try { try {
return mapper.writeValueAsString(mapper.createObjectNode() return mapper.writeValueAsString(mapper.createObjectNode()
.put("ruleNodeId", ruleNodeId.toString()) .put("ruleNodeId", ruleNodeId.toString())
.put("ruleChainId", ruleChainId.toString()) .put("ruleChainId", ruleChainId.toString())
.put("ruleNodeName", ruleNodeName) .put("ruleNodeName", ruleNodeName)
.put("ruleChainName", ruleChainName) .put("ruleChainName", ruleChainName)
.put("message", getMessage())); .put("message", truncateIfNecessary(getMessage(), maxMessageLength)));
} catch (JsonProcessingException e) { } catch (JsonProcessingException e) {
log.warn("Failed to serialize exception ", e); log.warn("Failed to serialize exception ", e);
throw new RuntimeException(e); throw new RuntimeException(e);

27
common/util/src/main/java/org/thingsboard/common/util/JacksonUtil.java

@ -70,8 +70,7 @@ public class JacksonUtil {
try { try {
return fromValue != null ? OBJECT_MAPPER.convertValue(fromValue, toValueType) : null; return fromValue != null ? OBJECT_MAPPER.convertValue(fromValue, toValueType) : null;
} catch (IllegalArgumentException e) { } catch (IllegalArgumentException e) {
throw new IllegalArgumentException("The given object value: " throw new IllegalArgumentException("The given object value cannot be converted to " + toValueType + ": " + fromValue, e);
+ fromValue + " cannot be converted to " + toValueType, e);
} }
} }
@ -79,8 +78,7 @@ public class JacksonUtil {
try { try {
return fromValue != null ? OBJECT_MAPPER.convertValue(fromValue, toValueTypeRef) : null; return fromValue != null ? OBJECT_MAPPER.convertValue(fromValue, toValueTypeRef) : null;
} catch (IllegalArgumentException e) { } catch (IllegalArgumentException e) {
throw new IllegalArgumentException("The given object value: " throw new IllegalArgumentException("The given object value cannot be converted to " + toValueTypeRef + ": " + fromValue, e);
+ fromValue + " cannot be converted to " + toValueTypeRef, e);
} }
} }
@ -88,8 +86,7 @@ public class JacksonUtil {
try { try {
return string != null ? OBJECT_MAPPER.readValue(string, clazz) : null; return string != null ? OBJECT_MAPPER.readValue(string, clazz) : null;
} catch (IOException e) { } catch (IOException e) {
throw new IllegalArgumentException("The given string value: " throw new IllegalArgumentException("The given string value cannot be transformed to Json object: " + string, e);
+ string + " cannot be transformed to Json object", e);
} }
} }
@ -97,8 +94,7 @@ public class JacksonUtil {
try { try {
return string != null ? OBJECT_MAPPER.readValue(string, valueTypeRef) : null; return string != null ? OBJECT_MAPPER.readValue(string, valueTypeRef) : null;
} catch (IOException e) { } catch (IOException e) {
throw new IllegalArgumentException("The given string value: " throw new IllegalArgumentException("The given string value cannot be transformed to Json object: " + string, e);
+ string + " cannot be transformed to Json object", e);
} }
} }
@ -106,8 +102,7 @@ public class JacksonUtil {
try { try {
return string != null ? OBJECT_MAPPER.readValue(string, javaType) : null; return string != null ? OBJECT_MAPPER.readValue(string, javaType) : null;
} catch (IOException e) { } catch (IOException e) {
throw new IllegalArgumentException("The given String value: " throw new IllegalArgumentException("The given String value cannot be transformed to Json object: " + string, e);
+ string + " cannot be transformed to Json object", e);
} }
} }
@ -115,8 +110,7 @@ public class JacksonUtil {
try { try {
return bytes != null ? OBJECT_MAPPER.readValue(bytes, clazz) : null; return bytes != null ? OBJECT_MAPPER.readValue(bytes, clazz) : null;
} catch (IOException e) { } catch (IOException e) {
throw new IllegalArgumentException("The given string value: " throw new IllegalArgumentException("The given string value cannot be transformed to Json object: " + Arrays.toString(bytes), e);
+ Arrays.toString(bytes) + " cannot be transformed to Json object", e);
} }
} }
@ -124,8 +118,7 @@ public class JacksonUtil {
try { try {
return OBJECT_MAPPER.readTree(bytes); return OBJECT_MAPPER.readTree(bytes);
} catch (IOException e) { } catch (IOException e) {
throw new IllegalArgumentException("The given byte[] value: " throw new IllegalArgumentException("The given byte[] value cannot be transformed to Json object: " + Arrays.toString(bytes), e);
+ Arrays.toString(bytes) + " cannot be transformed to Json object", e);
} }
} }
@ -133,8 +126,7 @@ public class JacksonUtil {
try { try {
return value != null ? OBJECT_MAPPER.writeValueAsString(value) : null; return value != null ? OBJECT_MAPPER.writeValueAsString(value) : null;
} catch (JsonProcessingException e) { } catch (JsonProcessingException e) {
throw new IllegalArgumentException("The given Json object value: " throw new IllegalArgumentException("The given Json object value cannot be transformed to a String: " + value, e);
+ value + " cannot be transformed to a String", e);
} }
} }
@ -208,8 +200,7 @@ public class JacksonUtil {
try { try {
return OBJECT_MAPPER.writeValueAsBytes(value); return OBJECT_MAPPER.writeValueAsBytes(value);
} catch (JsonProcessingException e) { } catch (JsonProcessingException e) {
throw new IllegalArgumentException("The given Json object value: " throw new IllegalArgumentException("The given Json object value cannot be transformed to a String: " + value, e);
+ value + " cannot be transformed to a String", e);
} }
} }

8
dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java

@ -92,12 +92,8 @@ public class BaseEventService implements EventService {
private <T extends Event> void truncateField(T event, Function<T, String> getter, BiConsumer<T, String> setter) { private <T extends Event> void truncateField(T event, Function<T, String> getter, BiConsumer<T, String> setter) {
var str = getter.apply(event); var str = getter.apply(event);
if (StringUtils.isNotEmpty(str)) { str = StringUtils.truncate(str, maxDebugEventSymbols);
var length = str.length(); setter.accept(event, str);
if (length > maxDebugEventSymbols) {
setter.accept(event, str.substring(0, maxDebugEventSymbols) + "...[truncated " + (length - maxDebugEventSymbols) + " symbols]");
}
}
} }
@Override @Override

Loading…
Cancel
Save