Browse Source

Save time series strategies: add tests for changes in save time series node

pull/12413/head
Dmytro Skarzhynets 2 years ago
parent
commit
0b277661bd
  1. 6
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java
  2. 11
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNodeConfiguration.java
  3. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/DeduplicatePersistenceStrategy.java
  4. 365
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNodeTest.java

6
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java

@ -90,11 +90,11 @@ public class TbMsgTimeseriesNode implements TbNode {
onTenantProfileUpdate(ctx.getTenantProfile());
persistenceSettings = config.getPersistenceSettings();
if (persistenceSettings == null) {
throw new TbNodeException("Persistence settings cannot be null!", true);
throw new TbNodeException("Persistence settings cannot be null", true);
}
}
void onTenantProfileUpdate(TenantProfile tenantProfile) {
private void onTenantProfileUpdate(TenantProfile tenantProfile) {
DefaultTenantProfileConfiguration configuration = (DefaultTenantProfileConfiguration) tenantProfile.getProfileData().getConfiguration();
tenantProfileDefaultStorageTtl = TimeUnit.DAYS.toSeconds(configuration.getDefaultStorageTtlDays());
}
@ -193,7 +193,7 @@ public class TbMsgTimeseriesNode implements TbNode {
boolean hasChanges = false;
switch (fromVersion) {
case 0:
if (oldConfiguration.has("persistenceSettings")) {
if (oldConfiguration.has("persistenceSettings") && !oldConfiguration.has("skipLatestPersistence")) {
break;
}
hasChanges = true;

11
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNodeConfiguration.java

@ -16,6 +16,7 @@
package org.thingsboard.rule.engine.telemetry;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonProperty;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
@ -67,11 +68,15 @@ public class TbMsgTimeseriesNodeConfiguration implements NodeConfiguration<TbMsg
@Getter
final class Deduplicate implements PersistenceSettings {
public final PersistenceStrategy deduplicateStrategy;
private final int deduplicationIntervalSecs;
@JsonIgnore
private final PersistenceStrategy deduplicateStrategy;
@JsonCreator
private Deduplicate(@JsonProperty("deduplicationIntervalSecs") int deduplicationIntervalSecs) {
this.deduplicateStrategy = PersistenceStrategy.deduplicate(deduplicationIntervalSecs);
Deduplicate(@JsonProperty("deduplicationIntervalSecs") int deduplicationIntervalSecs) {
this.deduplicationIntervalSecs = deduplicationIntervalSecs;
deduplicateStrategy = PersistenceStrategy.deduplicate(deduplicationIntervalSecs);
}
}

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/DeduplicatePersistenceStrategy.java

@ -45,6 +45,11 @@ final class DeduplicatePersistenceStrategy implements PersistenceStrategy {
.build(__ -> Sets.newConcurrentHashSet());
}
@JsonProperty("deduplicationIntervalSecs")
public long getDeduplicationIntervalSecs() {
return Duration.ofMillis(deduplicationIntervalMillis).toSeconds();
}
@Override
public boolean shouldPersist(long ts, UUID originatorUuid) {
long intervalNumber = ts / deduplicationIntervalMillis;

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

@ -41,6 +41,7 @@ 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.DoubleDataEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.msg.TbMsgType;
@ -57,25 +58,30 @@ import java.util.concurrent.TimeUnit;
import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.assertArg;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.then;
import static org.mockito.Mockito.clearInvocations;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoMoreInteractions;
import static org.mockito.Mockito.when;
@ExtendWith(MockitoExtension.class)
public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("c8f34868-603a-4433-876a-7d356e5cf377"));
private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("e5095e9a-04f4-44c9-b443-1cf1b97d3384"));
private final TenantProfileId TENANT_PROFILE_ID = new TenantProfileId(UUID.fromString("ab78dd78-83d0-43fa-869f-d42ec9ed1744"));
private TenantProfile tenantProfile;
private TbMsgTimeseriesNode node;
private TbMsgTimeseriesNodeConfiguration config;
private long tenantProfileDefaultStorageTtl;
@Mock
private TbContext ctxMock;
@ -84,6 +90,17 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
@BeforeEach
public void setUp() throws TbNodeException {
tenantProfile = new TenantProfile(new TenantProfileId(UUID.fromString("ab78dd78-83d0-43fa-869f-d42ec9ed1744")));
var tenantProfileConfiguration = new DefaultTenantProfileConfiguration();
tenantProfileConfiguration.setDefaultStorageTtlDays(5);
var tenantProfileData = new TenantProfileData();
tenantProfileData.setConfiguration(tenantProfileConfiguration);
tenantProfile.setProfileData(tenantProfileData);
lenient().when(ctxMock.getTenantProfile()).thenReturn(tenantProfile);
lenient().when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
lenient().when(ctxMock.getTelemetryService()).thenReturn(telemetryServiceMock);
node = spy(new TbMsgTimeseriesNode());
config = new TbMsgTimeseriesNodeConfiguration().defaultConfiguration();
}
@ -95,18 +112,47 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
assertThat(config.isUseServerTs()).isFalse();
}
@Test
public void whenInit_thenShouldAddTenantProfileListener() throws Exception {
// GIVEN-WHEN
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
// THEN
then(ctxMock).should().addTenantProfileListener(any());
}
@Test
public void givenPersistenceSettingsAreNull_whenInit_thenThrowsException() {
// GIVEN
config.setPersistenceSettings(null);
// WHEN-THEN
assertThatThrownBy(() -> node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config))))
.isInstanceOf(TbNodeException.class)
.matches(e -> ((TbNodeException) e).isUnrecoverable())
.hasMessage("Persistence settings cannot be null");
}
@ParameterizedTest
@EnumSource(TbMsgType.class)
public void givenMsgTypeAndEmptyMsgData_whenOnMsg_thenVerifyFailureMsg(TbMsgType msgType) throws TbNodeException {
init();
// GIVEN
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
TbMsg msg = TbMsg.newMsg()
.type(msgType)
.originator(DEVICE_ID)
.copyMetaData(TbMsgMetaData.EMPTY)
.data(TbMsg.EMPTY_JSON_ARRAY)
.build();
// WHEN
node.onMsg(ctxMock, msg);
// THEN
then(ctxMock).should().addTenantProfileListener(any());
then(ctxMock).should().getTenantProfile();
ArgumentCaptor<Throwable> throwableCaptor = ArgumentCaptor.forClass(Throwable.class);
verify(ctxMock).tellFailure(eq(msg), throwableCaptor.capture());
@ -120,9 +166,11 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
}
@Test
public void givenTtlFromConfigIsZeroAndUseServiceTsIsTrue_whenOnMsg_thenSaveTimeseriesUsingTenantProfileDefaultTtl() throws TbNodeException {
public void givenTtlFromConfigIsZeroAndUseServerTsIsTrue_whenOnMsg_thenSaveTimeseriesUsingTenantProfileDefaultTtl() throws TbNodeException {
// GIVEN
config.setUseServerTs(true);
init();
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
String data = """
{
@ -137,23 +185,28 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
.data(data)
.build();
when(ctxMock.getTelemetryService()).thenReturn(telemetryServiceMock);
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
doAnswer(invocation -> {
TimeseriesSaveRequest request = invocation.getArgument(0);
request.getCallback().onSuccess(null);
return null;
}).when(telemetryServiceMock).saveTimeseries(any(TimeseriesSaveRequest.class));
// WHEN
node.onMsg(ctxMock, msg);
// THEN
then(ctxMock).should().getTenantId();
then(ctxMock).should().getTelemetryService();
then(ctxMock).should().addTenantProfileListener(any());
then(ctxMock).should().getTenantProfile();
List<TsKvEntry> expectedList = getTsKvEntriesListWithTs(data, System.currentTimeMillis());
verify(telemetryServiceMock).saveTimeseries(assertArg(request -> {
assertThat(request.getTenantId()).isEqualTo(TENANT_ID);
assertThat(request.getCustomerId()).isNull();
assertThat(request.getEntityId()).isEqualTo(DEVICE_ID);
assertThat(request.getEntries()).usingRecursiveFieldByFieldElementComparatorIgnoringFields("ts").containsExactlyElementsOf(expectedList);
assertThat(request.getTtl()).isEqualTo(tenantProfileDefaultStorageTtl);
assertThat(request.getTtl()).isEqualTo(extractTtlAsSeconds(tenantProfile));
assertThat(request.isSaveLatest()).isTrue();
assertThat(request.getCallback()).isInstanceOf(TelemetryNodeCallback.class);
}));
@ -162,9 +215,9 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
}
@Test
public void givenSkipLatestPersistenceIsTrueAndTtlFromConfig_whenOnMsg_thenSaveTimeseriesUsingTtlFromConfig() throws TbNodeException {
long ttlFromConfig = 5L;
config.setDefaultTTL(ttlFromConfig);
public void givenSkipLatestPersistenceSettingsAndTtlFromConfig_whenOnMsg_thenSaveTimeseriesUsingTtlFromConfig() throws TbNodeException {
// GIVEN
config.setDefaultTTL(10L);
var timeseriesStrategy = PersistenceStrategy.onEveryMessage();
var latestStrategy = PersistenceStrategy.skip();
@ -172,7 +225,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
var persistenceSettings = new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Advanced(timeseriesStrategy, latestStrategy, webSockets);
config.setPersistenceSettings(persistenceSettings);
init();
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
String data = """
{
@ -189,23 +242,28 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
.data(data)
.build();
when(ctxMock.getTelemetryService()).thenReturn(telemetryServiceMock);
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
doAnswer(invocation -> {
TimeseriesSaveRequest request = invocation.getArgument(0);
request.getCallback().onSuccess(null);
return null;
}).when(telemetryServiceMock).saveTimeseries(any(TimeseriesSaveRequest.class));
// WHEN
node.onMsg(ctxMock, msg);
// THEN
then(ctxMock).should().getTenantId();
then(ctxMock).should().getTelemetryService();
then(ctxMock).should().addTenantProfileListener(any());
then(ctxMock).should().getTenantProfile();
List<TsKvEntry> expectedList = getTsKvEntriesListWithTs(data, ts);
verify(telemetryServiceMock).saveTimeseries(assertArg(request -> {
assertThat(request.getTenantId()).isEqualTo(TENANT_ID);
assertThat(request.getCustomerId()).isNull();
assertThat(request.getEntityId()).isEqualTo(DEVICE_ID);
assertThat(request.getEntries()).containsExactlyElementsOf(expectedList);
assertThat(request.getTtl()).isEqualTo(ttlFromConfig);
assertThat(request.getTtl()).isEqualTo(config.getDefaultTTL());
assertThat(request.isSaveTimeseries()).isTrue();
assertThat(request.isSaveLatest()).isFalse();
assertThat(request.isSendWsUpdate()).isTrue();
@ -218,11 +276,10 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
@ParameterizedTest
@MethodSource
public void givenTtlFromConfigAndTtlFromMd_whenOnMsg_thenVerifyTtl(String ttlFromMd, long ttlFromConfig, long expectedTtl) throws TbNodeException {
// GIVEN
config.setDefaultTTL(ttlFromConfig);
init();
when(ctxMock.getTelemetryService()).thenReturn(telemetryServiceMock);
when(ctxMock.getTenantId()).thenReturn(TENANT_ID);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
String data = """
{
@ -238,8 +295,11 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
.copyMetaData(metadata)
.data(data)
.build();
// WHEN
node.onMsg(ctxMock, msg);
// THEN
verify(telemetryServiceMock).saveTimeseries(assertArg(request -> {
assertThat(request.getTenantId()).isEqualTo(TENANT_ID);
assertThat(request.getCustomerId()).isNull();
@ -262,26 +322,6 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
);
}
private void init() throws TbNodeException {
var configuration = new TbNodeConfiguration(JacksonUtil.valueToTree(config));
var tenantProfile = getTenantProfile();
when(ctxMock.getTenantProfile()).thenReturn(tenantProfile);
tenantProfile.getProfileConfiguration().ifPresent(profileConfiguration ->
tenantProfileDefaultStorageTtl = TimeUnit.DAYS.toSeconds(profileConfiguration.getDefaultStorageTtlDays()));
node.init(ctxMock, configuration);
verify(ctxMock).addTenantProfileListener(any());
}
private TenantProfile getTenantProfile() {
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<>();
@ -293,6 +333,253 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
return expectedList;
}
@Test
public void givenOnEveryMessagePersistenceSettingsAndSameMessageTwoTimes_whenOnMsg_thenPersistSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.OnEveryMessage());
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(DEVICE_ID)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of("ts", "123")))
.build();
// WHEN-THEN
var expectedSaveRequest = TimeseriesSaveRequest.builder()
.tenantId(TENANT_ID)
.customerId(msg.getCustomerId())
.entityId(msg.getOriginator())
.entry(new BasicTsKvEntry(123L, new DoubleDataEntry("temperature", 22.3)))
.ttl(extractTtlAsSeconds(tenantProfile))
.saveTimeseries(true)
.saveLatest(true)
.sendWsUpdate(true)
.build();
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(1)).saveTimeseries(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest).usingRecursiveComparison().ignoringFields("callback").isEqualTo(expectedSaveRequest)
));
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(2)).saveTimeseries(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest).usingRecursiveComparison().ignoringFields("callback").isEqualTo(expectedSaveRequest)
));
}
@Test
public void givenDeduplicatePersistenceSettingsAndSameMessageTwoTimes_whenOnMsg_thenPersistThisMessageOnlyFirstTime() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Deduplicate(10));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(DEVICE_ID)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of("ts", "123")))
.build();
// WHEN-THEN
var expectedSaveRequest = TimeseriesSaveRequest.builder()
.tenantId(TENANT_ID)
.customerId(msg.getCustomerId())
.entityId(msg.getOriginator())
.entry(new BasicTsKvEntry(123L, new DoubleDataEntry("temperature", 22.3)))
.ttl(extractTtlAsSeconds(tenantProfile))
.saveTimeseries(true)
.saveLatest(true)
.sendWsUpdate(true)
.build();
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should().saveTimeseries(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest).usingRecursiveComparison().ignoringFields("callback").isEqualTo(expectedSaveRequest)
));
clearInvocations(telemetryServiceMock, ctxMock);
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(never()).saveTimeseries(any());
}
@Test
public void givenWebsocketsOnlyPersistenceSettingsAndSameMessageTwoTimes_whenOnMsg_thenSendsOnlyWsUpdateTwoTimes() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.WebSocketsOnly());
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(DEVICE_ID)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of("ts", "123")))
.build();
// WHEN-THEN
var expectedSaveRequest = TimeseriesSaveRequest.builder()
.tenantId(TENANT_ID)
.customerId(msg.getCustomerId())
.entityId(msg.getOriginator())
.entry(new BasicTsKvEntry(123L, new DoubleDataEntry("temperature", 22.3)))
.ttl(extractTtlAsSeconds(tenantProfile))
.saveTimeseries(false)
.saveLatest(false)
.sendWsUpdate(true)
.build();
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(1)).saveTimeseries(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest).usingRecursiveComparison().ignoringFields("callback").isEqualTo(expectedSaveRequest)
));
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(2)).saveTimeseries(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest).usingRecursiveComparison().ignoringFields("callback").isEqualTo(expectedSaveRequest)
));
}
@Test
public void givenAdvancedPersistenceSettingsWithOnEveryMessageStrategiesForAllActionsAndSameMessageTwoTimes_whenOnMsg_thenPersistSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Advanced(
PersistenceStrategy.onEveryMessage(),
PersistenceStrategy.onEveryMessage(),
PersistenceStrategy.onEveryMessage()
));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(DEVICE_ID)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of("ts", "123")))
.build();
// WHEN-THEN
var expectedSaveRequest = TimeseriesSaveRequest.builder()
.tenantId(TENANT_ID)
.customerId(msg.getCustomerId())
.entityId(msg.getOriginator())
.entry(new BasicTsKvEntry(123L, new DoubleDataEntry("temperature", 22.3)))
.ttl(extractTtlAsSeconds(tenantProfile))
.saveTimeseries(true)
.saveLatest(true)
.sendWsUpdate(true)
.build();
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(1)).saveTimeseries(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest).usingRecursiveComparison().ignoringFields("callback").isEqualTo(expectedSaveRequest)
));
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(2)).saveTimeseries(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest).usingRecursiveComparison().ignoringFields("callback").isEqualTo(expectedSaveRequest)
));
}
@Test
public void givenAdvancedPersistenceSettingsWithDifferentDeduplicateStrategyForEachAction_whenOnMsg_thenEvaluatesStrategiesForEachActionsIndependently() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Advanced(
PersistenceStrategy.deduplicate(1),
PersistenceStrategy.deduplicate(2),
PersistenceStrategy.deduplicate(3)
));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
long ts1 = 500L;
long ts2 = 1500L;
long ts3 = 2500L;
// WHEN-THEN
node.onMsg(ctxMock, TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(DEVICE_ID)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of("ts", Long.toString(ts1))))
.build());
then(telemetryServiceMock).should().saveTimeseries(assertArg(
actualSaveRequest -> {
assertThat(actualSaveRequest.isSaveTimeseries()).isTrue();
assertThat(actualSaveRequest.isSaveLatest()).isTrue();
assertThat(actualSaveRequest.isSendWsUpdate()).isTrue();
}
));
clearInvocations(telemetryServiceMock);
node.onMsg(ctxMock, TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(DEVICE_ID)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of("ts", Long.toString(ts2))))
.build());
then(telemetryServiceMock).should().saveTimeseries(assertArg(
actualSaveRequest -> {
assertThat(actualSaveRequest.isSaveTimeseries()).isTrue();
assertThat(actualSaveRequest.isSaveLatest()).isFalse();
assertThat(actualSaveRequest.isSendWsUpdate()).isFalse();
}
));
clearInvocations(telemetryServiceMock);
node.onMsg(ctxMock, TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(DEVICE_ID)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of("ts", Long.toString(ts3))))
.build());
then(telemetryServiceMock).should().saveTimeseries(assertArg(
actualSaveRequest -> {
assertThat(actualSaveRequest.isSaveTimeseries()).isTrue();
assertThat(actualSaveRequest.isSaveLatest()).isTrue();
assertThat(actualSaveRequest.isSendWsUpdate()).isFalse();
}
));
}
@Test
public void givenAdvancedPersistenceSettingsWithSkipStrategiesForAllActionsAndSameMessageTwoTimes_whenOnMsg_thenSkipsSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Advanced(
PersistenceStrategy.skip(),
PersistenceStrategy.skip(),
PersistenceStrategy.skip()
));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
.originator(DEVICE_ID)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of("ts", "123")))
.build();
// WHEN-THEN
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(never()).saveTimeseries(any());
then(ctxMock).should(times(1)).tellSuccess(msg);
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(never()).saveTimeseries(any());
then(ctxMock).should(times(2)).tellSuccess(msg);
}
private static long extractTtlAsSeconds(TenantProfile tenantProfile) {
return TimeUnit.DAYS.toSeconds(tenantProfile.getDefaultProfileConfiguration().getDefaultStorageTtlDays());
}
@Override
protected TbNode getTestNode() {
return node;

Loading…
Cancel
Save