Browse Source

added test for default config and set different values for check logic of setting ttl

pull/10556/head
IrynaMatveieva 2 years ago
parent
commit
4fc6929922
  1. 89
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNodeTest.java

89
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNodeTest.java

@ -15,7 +15,7 @@
*/
package org.thingsboard.rule.engine.telemetry;
import com.fasterxml.jackson.databind.JsonNode;
import com.google.gson.JsonParser;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
@ -31,10 +31,13 @@ import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.adaptor.JsonConverter;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
@ -42,14 +45,14 @@ import org.thingsboard.server.common.data.tenant.profile.TenantProfileData;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.HashMap;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import java.util.stream.Stream;
import static org.assertj.core.api.AssertionsForClassTypes.assertThat;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.ArgumentMatchers.anyLong;
@ -82,17 +85,26 @@ public class TbMsgTimeseriesNodeTest {
config = new TbMsgTimeseriesNodeConfiguration().defaultConfiguration();
}
@Test
public void verifyDefaultConfig() {
assertThat(config.getDefaultTTL()).isEqualTo(0L);
assertThat(config.isSkipLatestPersistence()).isFalse();
assertThat(config.isUseServerTs()).isFalse();
}
@ParameterizedTest
@EnumSource(TbMsgType.class)
void givenMsgTypeAndEmptyMsgData_whenOnMsg_thenVerifyFailureMsg(TbMsgType msgType) throws TbNodeException {
public void givenMsgTypeAndEmptyMsgData_whenOnMsg_thenVerifyFailureMsg(TbMsgType msgType) throws TbNodeException {
init();
TbMsg msg = TbMsg.newMsg(msgType, DEVICE_ID, TbMsgMetaData.EMPTY, TbMsg.EMPTY_JSON_ARRAY);
node.onMsg(ctxMock, msg);
ArgumentCaptor<Throwable> throwableCaptor = ArgumentCaptor.forClass(Throwable.class);
verify(ctxMock).tellFailure(eq(msg), throwableCaptor.capture());
if (TbMsgType.POST_TELEMETRY_REQUEST.equals(msgType)) {
assertThat(throwableCaptor.getValue()).isInstanceOf(IllegalArgumentException.class).hasMessage("Msg body is empty: " + msg.getData());
verifyNoMoreInteractions(ctxMock);
return;
}
assertThat(throwableCaptor.getValue()).isInstanceOf(IllegalArgumentException.class).hasMessage("Unsupported msg type: " + msgType);
@ -100,9 +112,10 @@ public class TbMsgTimeseriesNodeTest {
}
@Test
void givenTtlFromConfigIsZeroAndUseServiceTsIsTrue_whenOnMsg_thenSaveTimeseriesUsingTenantProfileDefaultTtl() throws TbNodeException {
public void givenTtlFromConfigIsZeroAndUseServiceTsIsTrue_whenOnMsg_thenSaveTimeseriesUsingTenantProfileDefaultTtl() throws TbNodeException {
config.setUseServerTs(true);
init();
String data = """
{
"temp": 45,
@ -121,25 +134,22 @@ public class TbMsgTimeseriesNodeTest {
node.onMsg(ctxMock, msg);
List<TsKvEntry> expectedList = getTsKvEntriesListWithTs(data, System.currentTimeMillis());
ArgumentCaptor<List<TsKvEntry>> entryListCaptor = ArgumentCaptor.forClass(List.class);
verify(telemetryServiceMock).saveAndNotify(eq(TENANT_ID), isNull(), eq(DEVICE_ID), entryListCaptor.capture(),
eq(tenantProfileDefaultStorageTtl), any(TelemetryNodeCallback.class));
List<TsKvEntry> entryListCaptorValue = entryListCaptor.getValue();
assertThat(entryListCaptorValue.size()).isEqualTo(2);
verifyTimeseriesToSave(entryListCaptorValue, msg);
assertThat(entryListCaptor.getValue()).usingRecursiveFieldByFieldElementComparatorIgnoringFields("ts")
.containsExactlyElementsOf(expectedList);
verify(ctxMock).tellSuccess(msg);
verifyNoMoreInteractions(ctxMock, telemetryServiceMock);
}
@Test
void givenSkipLatestPersistenceIsTrueAndTtlFromConfig_whenOnMsg_thenSaveTimeseriesUsingTtlFromConfig() throws TbNodeException {
public void givenSkipLatestPersistenceIsTrueAndTtlFromConfig_whenOnMsg_thenSaveTimeseriesUsingTtlFromConfig() throws TbNodeException {
long ttlFromConfig = 5L;
config.setDefaultTTL(ttlFromConfig);
config.setSkipLatestPersistence(true);
var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
var tenantProfile = getTenantProfile();
when(ctxMock.getTenantProfile()).thenReturn(tenantProfile);
node.init(ctxMock, configuration);
init();
String data = """
{
@ -161,20 +171,18 @@ public class TbMsgTimeseriesNodeTest {
node.onMsg(ctxMock, msg);
verify(ctxMock).addTenantProfileListener(any());
List<TsKvEntry> expectedList = getTsKvEntriesListWithTs(data, ts);
ArgumentCaptor<List<TsKvEntry>> entryListCaptor = ArgumentCaptor.forClass(List.class);
verify(telemetryServiceMock).saveWithoutLatestAndNotify(
eq(TENANT_ID), isNull(), eq(DEVICE_ID), entryListCaptor.capture(), eq(ttlFromConfig), any(TelemetryNodeCallback.class));
List<TsKvEntry> entryListCaptorValue = entryListCaptor.getValue();
assertThat(entryListCaptorValue.size()).isEqualTo(2);
verifyTimeseriesToSave(entryListCaptorValue, msg, ts);
assertThat(entryListCaptor.getValue()).containsExactlyElementsOf(expectedList);
verify(ctxMock).tellSuccess(msg);
verifyNoMoreInteractions(ctxMock, telemetryServiceMock);
}
@ParameterizedTest
@MethodSource
void givenTtlFromConfigAndTtlFromMd_whenOnMsg_thenVerifyTtl(String ttlFromMd, long ttlFromConfig, long expectedTtl) throws TbNodeException {
public void givenTtlFromConfigAndTtlFromMd_whenOnMsg_thenVerifyTtl(String ttlFromMd, long ttlFromConfig, long expectedTtl) throws TbNodeException {
config.setDefaultTTL(ttlFromConfig);
init();
@ -187,9 +195,9 @@ public class TbMsgTimeseriesNodeTest {
"humidity": 77
}
""";
var metadata = new HashMap<String, String>();
metadata.put("TTL", ttlFromMd);
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, new TbMsgMetaData(metadata), data);
var metadata = new TbMsgMetaData();
metadata.putValue("TTL", ttlFromMd);
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, DEVICE_ID, metadata, data);
node.onMsg(ctxMock, msg);
verify(telemetryServiceMock).saveAndNotify(eq(TENANT_ID), isNull(), eq(DEVICE_ID), anyList(), eq(expectedTtl), any(TelemetryNodeCallback.class));
@ -197,36 +205,23 @@ public class TbMsgTimeseriesNodeTest {
private static Stream<Arguments> givenTtlFromConfigAndTtlFromMd_whenOnMsg_thenVerifyTtl() {
return Stream.of(
Arguments.of("5", 1L, 5L),
// when ttl is present in metadata and it is not zero then ttl = ttl from metadata
Arguments.of("1", 2L, 1L),
// when ttl is absent in metadata and present in config and it is not zero then ttl = ttl from config
Arguments.of("", 3L, 3L),
Arguments.of(null, 8L, 8L)
Arguments.of(null, 4L, 4L),
// when ttl is present in metadata or config but it is zero then ttl = default ttl from tenant profile
Arguments.of("0", 0L, TimeUnit.DAYS.toSeconds(5L))
);
}
private void verifyTimeseriesToSave(List<TsKvEntry> tsKvEntryList, TbMsg incomingMsg) {
verifyTimeseriesToSave(tsKvEntryList, incomingMsg, null);
}
private void verifyTimeseriesToSave(List<TsKvEntry> tsKvEntryList, TbMsg incomingMsg, Long ts) {
JsonNode msgData = JacksonUtil.toJsonNode(incomingMsg.getData());
tsKvEntryList.forEach(tsKvEntry -> {
if (ts != null) {
assertThat(tsKvEntry.getTs()).isEqualTo(ts);
}
String key = tsKvEntry.getKey();
assertThat(msgData.has(key)).isTrue();
String value = tsKvEntry.getValueAsString();
assertThat(value).isEqualTo(msgData.findValue(key).asText());
});
}
private void init() throws TbNodeException {
var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
var tenantProfile = getTenantProfile();
when(ctxMock.getTenantProfile()).thenReturn(tenantProfile);
node.init(ctxMock, configuration);
tenantProfile.getProfileConfiguration().ifPresent(profileConfiguration ->
tenantProfileDefaultStorageTtl = TimeUnit.DAYS.toSeconds(profileConfiguration.getDefaultStorageTtlDays()));
node.init(ctxMock, configuration);
verify(ctxMock).addTenantProfileListener(any());
}
@ -234,9 +229,21 @@ public class TbMsgTimeseriesNodeTest {
var tenantProfile = new TenantProfile(TENANT_PROFILE_ID);
var tenantProfileData = new TenantProfileData();
var tenantProfileConfiguration = new DefaultTenantProfileConfiguration();
tenantProfileConfiguration.setDefaultStorageTtlDays(5);
tenantProfileData.setConfiguration(tenantProfileConfiguration);
tenantProfile.setProfileData(tenantProfileData);
return tenantProfile;
}
private static List<TsKvEntry> getTsKvEntriesListWithTs(String data, long ts) {
Map<Long, List<KvEntry>> tsKvMap = JsonConverter.convertToTelemetry(JsonParser.parseString(data), ts);
List<TsKvEntry> expectedList = new ArrayList<>();
for (Map.Entry<Long, List<KvEntry>> tsKvEntry : tsKvMap.entrySet()) {
for (KvEntry kvEntry : tsKvEntry.getValue()) {
expectedList.add(new BasicTsKvEntry(tsKvEntry.getKey(), kvEntry));
}
}
return expectedList;
}
}

Loading…
Cancel
Save