Browse Source

Save time series strategies: rename SaveActions to Strategy

pull/12413/head
Dmytro Skarzhynets 2 years ago
parent
commit
e009967fa7
  1. 2
      application/src/main/java/org/thingsboard/server/service/entitiy/entityview/DefaultTbEntityViewService.java
  2. 18
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  3. 2
      application/src/test/java/org/thingsboard/server/service/entitiy/entityview/DefaultTbEntityViewServiceTest.java
  4. 14
      application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java
  5. 20
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TimeseriesSaveRequest.java
  6. 20
      rule-engine/rule-engine-api/src/test/java/org/thingsboard/rule/engine/api/TimeseriesSaveRequestTest.java
  7. 16
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java
  8. 4
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java
  9. 20
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNodeTest.java

2
application/src/main/java/org/thingsboard/server/service/entitiy/entityview/DefaultTbEntityViewService.java

@ -348,7 +348,7 @@ public class DefaultTbEntityViewService extends AbstractTbEntityService implemen
.tenantId(entityView.getTenantId())
.entityId(entityId)
.entries(latestValues)
.saveActions(TimeseriesSaveRequest.SaveActions.LATEST_AND_WS)
.strategy(TimeseriesSaveRequest.Strategy.LATEST_AND_WS)
.callback(new FutureCallback<Void>() {
@Override
public void onSuccess(@Nullable Void tmp) {

18
application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java

@ -118,10 +118,10 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
EntityId entityId = request.getEntityId();
checkInternalEntity(entityId);
boolean sysTenant = TenantId.SYS_TENANT_ID.equals(tenantId) || tenantId == null;
if (sysTenant || !request.getSaveActions().saveTimeseries() || apiUsageStateService.getApiUsageState(tenantId).isDbStorageEnabled()) {
if (sysTenant || !request.getStrategy().saveTimeseries() || apiUsageStateService.getApiUsageState(tenantId).isDbStorageEnabled()) {
KvUtils.validate(request.getEntries(), valueNoXssValidation);
ListenableFuture<Integer> future = saveTimeseriesInternal(request);
if (request.getSaveActions().saveTimeseries()) {
if (request.getStrategy().saveTimeseries()) {
FutureCallback<Integer> callback = getApiUsageCallback(tenantId, request.getCustomerId(), sysTenant, request.getCallback());
Futures.addCallback(future, callback, tsCallBackExecutor);
}
@ -134,23 +134,23 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
public ListenableFuture<Integer> saveTimeseriesInternal(TimeseriesSaveRequest request) {
TenantId tenantId = request.getTenantId();
EntityId entityId = request.getEntityId();
TimeseriesSaveRequest.SaveActions saveActions = request.getSaveActions();
TimeseriesSaveRequest.Strategy strategy = request.getStrategy();
ListenableFuture<Integer> saveFuture;
if (saveActions.saveTimeseries() && saveActions.saveLatest()) {
if (strategy.saveTimeseries() && strategy.saveLatest()) {
saveFuture = tsService.save(tenantId, entityId, request.getEntries(), request.getTtl());
} else if (saveActions.saveLatest()) {
} else if (strategy.saveLatest()) {
saveFuture = Futures.transform(tsService.saveLatest(tenantId, entityId, request.getEntries()), result -> 0, MoreExecutors.directExecutor());
} else if (saveActions.saveTimeseries()) {
} else if (strategy.saveTimeseries()) {
saveFuture = tsService.saveWithoutLatest(tenantId, entityId, request.getEntries(), request.getTtl());
} else {
saveFuture = Futures.immediateFuture(0);
}
addMainCallback(saveFuture, request.getCallback());
if (saveActions.sendWsUpdate()) {
if (strategy.sendWsUpdate()) {
addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, request.getEntries()));
}
if (saveActions.saveLatest()) {
if (strategy.saveLatest()) {
copyLatestToEntityViews(tenantId, entityId, request.getEntries());
}
return saveFuture;
@ -237,7 +237,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
.tenantId(tenantId)
.entityId(entityView.getId())
.entries(entityViewLatest)
.saveActions(TimeseriesSaveRequest.SaveActions.LATEST_AND_WS)
.strategy(TimeseriesSaveRequest.Strategy.LATEST_AND_WS)
.callback(new FutureCallback<>() {
@Override
public void onSuccess(@Nullable Void tmp) {}

2
application/src/test/java/org/thingsboard/server/service/entitiy/entityview/DefaultTbEntityViewServiceTest.java

@ -93,7 +93,7 @@ class DefaultTbEntityViewServiceTest {
.entityId(entityView.getId())
.entries(latest)
.ttl(0L)
.saveActions(TimeseriesSaveRequest.SaveActions.LATEST_AND_WS)
.strategy(TimeseriesSaveRequest.Strategy.LATEST_AND_WS)
.build();
var actualCopyLatestRequest = captor.getValue();

14
application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java

@ -173,7 +173,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.entityId(entityId)
.entries(sampleTelemetry)
.ttl(sampleTtl)
.saveActions(new TimeseriesSaveRequest.SaveActions(true, false, false))
.strategy(new TimeseriesSaveRequest.Strategy(true, false, false))
.callback(emptyCallback)
.build();
@ -193,7 +193,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.entityId(entityId)
.entries(sampleTelemetry)
.ttl(sampleTtl)
.saveActions(TimeseriesSaveRequest.SaveActions.LATEST_AND_WS)
.strategy(TimeseriesSaveRequest.Strategy.LATEST_AND_WS)
.callback(emptyCallback)
.build();
@ -216,7 +216,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.entityId(entityId)
.entries(sampleTelemetry)
.ttl(sampleTtl)
.saveActions(TimeseriesSaveRequest.SaveActions.SAVE_ALL)
.strategy(TimeseriesSaveRequest.Strategy.SAVE_ALL)
.future(future)
.build();
@ -242,7 +242,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.entityId(entityId)
.entries(sampleTelemetry)
.ttl(sampleTtl)
.saveActions(TimeseriesSaveRequest.SaveActions.LATEST_AND_WS)
.strategy(TimeseriesSaveRequest.Strategy.LATEST_AND_WS)
.future(future)
.build();
@ -275,7 +275,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.entityId(entityId)
.entries(sampleTelemetry)
.ttl(sampleTtl)
.saveActions(new TimeseriesSaveRequest.SaveActions(false, true, false))
.strategy(new TimeseriesSaveRequest.Strategy(false, true, false))
.callback(emptyCallback)
.build();
@ -302,7 +302,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.entityId(entityId)
.entries(sampleTelemetry)
.ttl(sampleTtl)
.saveActions(new TimeseriesSaveRequest.SaveActions(true, false, false))
.strategy(new TimeseriesSaveRequest.Strategy(true, false, false))
.callback(emptyCallback)
.build();
@ -328,7 +328,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.entityId(entityId)
.entries(sampleTelemetry)
.ttl(sampleTtl)
.saveActions(new TimeseriesSaveRequest.SaveActions(saveTimeseries, saveLatest, sendWsUpdate))
.strategy(new TimeseriesSaveRequest.Strategy(saveTimeseries, saveLatest, sendWsUpdate))
.callback(emptyCallback)
.build();

20
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TimeseriesSaveRequest.java

@ -38,15 +38,15 @@ public class TimeseriesSaveRequest {
private final EntityId entityId;
private final List<TsKvEntry> entries;
private final long ttl;
private final SaveActions saveActions;
private final Strategy strategy;
private final FutureCallback<Void> callback;
public record SaveActions(boolean saveTimeseries, boolean saveLatest, boolean sendWsUpdate) {
public record Strategy(boolean saveTimeseries, boolean saveLatest, boolean sendWsUpdate) {
public static final SaveActions SAVE_ALL = new SaveActions(true, true, true);
public static final SaveActions WS_ONLY = new SaveActions(false, false, true);
public static final SaveActions LATEST_AND_WS = new SaveActions(false, true, true);
public static final SaveActions SKIP_ALL = new SaveActions(false, false, false);
public static final Strategy SAVE_ALL = new Strategy(true, true, true);
public static final Strategy WS_ONLY = new Strategy(false, false, true);
public static final Strategy LATEST_AND_WS = new Strategy(false, true, true);
public static final Strategy SKIP_ALL = new Strategy(false, false, false);
}
@ -61,7 +61,7 @@ public class TimeseriesSaveRequest {
private EntityId entityId;
private List<TsKvEntry> entries;
private long ttl;
private SaveActions saveActions = SaveActions.SAVE_ALL;
private Strategy strategy = Strategy.SAVE_ALL;
private FutureCallback<Void> callback;
Builder() {}
@ -99,8 +99,8 @@ public class TimeseriesSaveRequest {
return this;
}
public Builder saveActions(SaveActions settings) {
this.saveActions = settings;
public Builder strategy(Strategy strategy) {
this.strategy = strategy;
return this;
}
@ -124,7 +124,7 @@ public class TimeseriesSaveRequest {
}
public TimeseriesSaveRequest build() {
return new TimeseriesSaveRequest(tenantId, customerId, entityId, entries, ttl, saveActions, callback);
return new TimeseriesSaveRequest(tenantId, customerId, entityId, entries, ttl, strategy, callback);
}
}

20
rule-engine/rule-engine-api/src/test/java/org/thingsboard/rule/engine/api/TimeseriesSaveRequestTest.java

@ -22,30 +22,30 @@ import static org.assertj.core.api.Assertions.assertThat;
class TimeseriesSaveRequestTest {
@Test
void testDefaultSaveActionsAreSaveAll() {
void testDefaultSaveStrategyIsSaveAll() {
var request = TimeseriesSaveRequest.builder().build();
assertThat(request.getSaveActions()).isEqualTo(TimeseriesSaveRequest.SaveActions.SAVE_ALL);
assertThat(request.getStrategy()).isEqualTo(TimeseriesSaveRequest.Strategy.SAVE_ALL);
}
@Test
void testSaveActionsSaveAll() {
assertThat(TimeseriesSaveRequest.SaveActions.SAVE_ALL).isEqualTo(new TimeseriesSaveRequest.SaveActions(true, true, true));
void testSaveAllStrategy() {
assertThat(TimeseriesSaveRequest.Strategy.SAVE_ALL).isEqualTo(new TimeseriesSaveRequest.Strategy(true, true, true));
}
@Test
void testSaveActionsWsOnly() {
assertThat(TimeseriesSaveRequest.SaveActions.WS_ONLY).isEqualTo(new TimeseriesSaveRequest.SaveActions(false, false, true));
void testWsOnlyStrategy() {
assertThat(TimeseriesSaveRequest.Strategy.WS_ONLY).isEqualTo(new TimeseriesSaveRequest.Strategy(false, false, true));
}
@Test
void testSaveActionsLatestAndWs() {
assertThat(TimeseriesSaveRequest.SaveActions.LATEST_AND_WS).isEqualTo(new TimeseriesSaveRequest.SaveActions(false, true, true));
void testLatestAndWsStrategy() {
assertThat(TimeseriesSaveRequest.Strategy.LATEST_AND_WS).isEqualTo(new TimeseriesSaveRequest.Strategy(false, true, true));
}
@Test
void testSaveActionsSkipAll() {
assertThat(TimeseriesSaveRequest.SaveActions.SKIP_ALL).isEqualTo(new TimeseriesSaveRequest.SaveActions(false, false, false));
void testSkipAllStrategy() {
assertThat(TimeseriesSaveRequest.Strategy.SKIP_ALL).isEqualTo(new TimeseriesSaveRequest.Strategy(false, false, false));
}
}

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

@ -104,10 +104,10 @@ public class TbMsgTimeseriesNode implements TbNode {
}
long ts = computeTs(msg, config.isUseServerTs());
TimeseriesSaveRequest.SaveActions saveActions = determineSaveActions(ts, msg.getOriginator().getId());
TimeseriesSaveRequest.Strategy strategy = determineSaveActions(ts, msg.getOriginator().getId());
// short-circuit
if (!saveActions.saveTimeseries() && !saveActions.saveLatest() && !saveActions.sendWsUpdate()) {
if (!strategy.saveTimeseries() && !strategy.saveLatest() && !strategy.sendWsUpdate()) {
ctx.tellSuccess(msg);
return;
}
@ -135,7 +135,7 @@ public class TbMsgTimeseriesNode implements TbNode {
.entityId(msg.getOriginator())
.entries(tsKvEntryList)
.ttl(ttl)
.saveActions(saveActions)
.strategy(strategy)
.callback(new TelemetryNodeCallback(ctx, msg))
.build());
}
@ -144,19 +144,19 @@ public class TbMsgTimeseriesNode implements TbNode {
return ignoreMetadataTs ? System.currentTimeMillis() : msg.getMetaDataTs();
}
private TimeseriesSaveRequest.SaveActions determineSaveActions(long ts, UUID originatorUuid) {
private TimeseriesSaveRequest.Strategy determineSaveActions(long ts, UUID originatorUuid) {
if (persistenceSettings instanceof OnEveryMessage) {
return TimeseriesSaveRequest.SaveActions.SAVE_ALL;
return TimeseriesSaveRequest.Strategy.SAVE_ALL;
}
if (persistenceSettings instanceof WebSocketsOnly) {
return TimeseriesSaveRequest.SaveActions.WS_ONLY;
return TimeseriesSaveRequest.Strategy.WS_ONLY;
}
if (persistenceSettings instanceof Deduplicate deduplicate) {
boolean isFirstMsgInInterval = deduplicate.getDeduplicateStrategy().shouldPersist(ts, originatorUuid);
return isFirstMsgInInterval ? TimeseriesSaveRequest.SaveActions.SAVE_ALL : TimeseriesSaveRequest.SaveActions.SKIP_ALL;
return isFirstMsgInInterval ? TimeseriesSaveRequest.Strategy.SAVE_ALL : TimeseriesSaveRequest.Strategy.SKIP_ALL;
}
if (persistenceSettings instanceof Advanced advanced) {
return new TimeseriesSaveRequest.SaveActions(
return new TimeseriesSaveRequest.Strategy(
advanced.timeseries().shouldPersist(ts, originatorUuid),
advanced.latest().shouldPersist(ts, originatorUuid),
advanced.webSockets().shouldPersist(ts, originatorUuid)

4
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java

@ -533,7 +533,7 @@ public class TbMathNodeTest {
verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture());
verify(telemetryService, times(1)).saveTimeseries(assertArg(request -> {
assertThat(request.getEntries()).size().isOne();
assertThat(request.getSaveActions()).isEqualTo(TimeseriesSaveRequest.SaveActions.SAVE_ALL);
assertThat(request.getStrategy()).isEqualTo(TimeseriesSaveRequest.Strategy.SAVE_ALL);
}));
TbMsg resultMsg = msgCaptor.getValue();
@ -569,7 +569,7 @@ public class TbMathNodeTest {
verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture());
verify(telemetryService, times(1)).saveTimeseries(assertArg(request -> {
assertThat(request.getEntries()).size().isOne();
assertThat(request.getSaveActions()).isEqualTo(TimeseriesSaveRequest.SaveActions.SAVE_ALL);
assertThat(request.getStrategy()).isEqualTo(TimeseriesSaveRequest.Strategy.SAVE_ALL);
}));
TbMsg resultMsg = msgCaptor.getValue();

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

@ -208,7 +208,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
assertThat(request.getEntityId()).isEqualTo(DEVICE_ID);
assertThat(request.getEntries()).usingRecursiveFieldByFieldElementComparatorIgnoringFields("ts").containsExactlyElementsOf(expectedList);
assertThat(request.getTtl()).isEqualTo(extractTtlAsSeconds(tenantProfile));
assertThat(request.getSaveActions()).isEqualTo(TimeseriesSaveRequest.SaveActions.SAVE_ALL);
assertThat(request.getStrategy()).isEqualTo(TimeseriesSaveRequest.Strategy.SAVE_ALL);
assertThat(request.getCallback()).isInstanceOf(TelemetryNodeCallback.class);
}));
verify(ctxMock).tellSuccess(msg);
@ -265,7 +265,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
assertThat(request.getEntityId()).isEqualTo(DEVICE_ID);
assertThat(request.getEntries()).containsExactlyElementsOf(expectedList);
assertThat(request.getTtl()).isEqualTo(config.getDefaultTTL());
assertThat(request.getSaveActions()).isEqualTo(new TimeseriesSaveRequest.SaveActions(true, false, true));
assertThat(request.getStrategy()).isEqualTo(new TimeseriesSaveRequest.Strategy(true, false, true));
assertThat(request.getCallback()).isInstanceOf(TelemetryNodeCallback.class);
}));
verify(ctxMock).tellSuccess(msg);
@ -304,7 +304,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
assertThat(request.getCustomerId()).isNull();
assertThat(request.getEntityId()).isEqualTo(DEVICE_ID);
assertThat(request.getTtl()).isEqualTo(expectedTtl);
assertThat(request.getSaveActions()).isEqualTo(TimeseriesSaveRequest.SaveActions.SAVE_ALL);
assertThat(request.getStrategy()).isEqualTo(TimeseriesSaveRequest.Strategy.SAVE_ALL);
assertThat(request.getCallback()).isInstanceOf(TelemetryNodeCallback.class);
}));
}
@ -353,7 +353,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
.entityId(msg.getOriginator())
.entry(new BasicTsKvEntry(123L, new DoubleDataEntry("temperature", 22.3)))
.ttl(extractTtlAsSeconds(tenantProfile))
.saveActions(TimeseriesSaveRequest.SaveActions.SAVE_ALL)
.strategy(TimeseriesSaveRequest.Strategy.SAVE_ALL)
.build();
node.onMsg(ctxMock, msg);
@ -388,7 +388,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
.entityId(msg.getOriginator())
.entry(new BasicTsKvEntry(123L, new DoubleDataEntry("temperature", 22.3)))
.ttl(extractTtlAsSeconds(tenantProfile))
.saveActions(TimeseriesSaveRequest.SaveActions.SAVE_ALL)
.strategy(TimeseriesSaveRequest.Strategy.SAVE_ALL)
.build();
node.onMsg(ctxMock, msg);
@ -423,7 +423,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
.entityId(msg.getOriginator())
.entry(new BasicTsKvEntry(123L, new DoubleDataEntry("temperature", 22.3)))
.ttl(extractTtlAsSeconds(tenantProfile))
.saveActions(TimeseriesSaveRequest.SaveActions.WS_ONLY)
.strategy(TimeseriesSaveRequest.Strategy.WS_ONLY)
.build();
node.onMsg(ctxMock, msg);
@ -462,7 +462,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
.entityId(msg.getOriginator())
.entry(new BasicTsKvEntry(123L, new DoubleDataEntry("temperature", 22.3)))
.ttl(extractTtlAsSeconds(tenantProfile))
.saveActions(TimeseriesSaveRequest.SaveActions.SAVE_ALL)
.strategy(TimeseriesSaveRequest.Strategy.SAVE_ALL)
.build();
node.onMsg(ctxMock, msg);
@ -499,7 +499,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
.metaData(new TbMsgMetaData(Map.of("ts", Long.toString(ts1))))
.build());
then(telemetryServiceMock).should().saveTimeseries(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest.getSaveActions()).isEqualTo(TimeseriesSaveRequest.SaveActions.SAVE_ALL)
actualSaveRequest -> assertThat(actualSaveRequest.getStrategy()).isEqualTo(TimeseriesSaveRequest.Strategy.SAVE_ALL)
));
clearInvocations(telemetryServiceMock);
@ -511,7 +511,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
.metaData(new TbMsgMetaData(Map.of("ts", Long.toString(ts2))))
.build());
then(telemetryServiceMock).should().saveTimeseries(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest.getSaveActions()).isEqualTo(new TimeseriesSaveRequest.SaveActions(true, false, false))
actualSaveRequest -> assertThat(actualSaveRequest.getStrategy()).isEqualTo(new TimeseriesSaveRequest.Strategy(true, false, false))
));
clearInvocations(telemetryServiceMock);
@ -523,7 +523,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
.metaData(new TbMsgMetaData(Map.of("ts", Long.toString(ts3))))
.build());
then(telemetryServiceMock).should().saveTimeseries(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest.getSaveActions()).isEqualTo(new TimeseriesSaveRequest.SaveActions(true, true, false))
actualSaveRequest -> assertThat(actualSaveRequest.getStrategy()).isEqualTo(new TimeseriesSaveRequest.Strategy(true, true, false))
));
}

Loading…
Cancel
Save