Browse Source

Merge pull request #12621 from dskarzh/feature/save-timeseries-strategies

Save time series strategies: corrections after review
pull/12690/head
Viacheslav Klimov 1 year ago
committed by GitHub
parent
commit
b4edfeb725
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 2
      application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json
  2. 2
      application/src/main/data/json/tenant/device_profile/rule_chain_template.json
  3. 2
      application/src/main/data/json/tenant/rule_chains/root_rule_chain.json
  4. 4
      application/src/main/data/upgrade/basic/schema_update.sql
  5. 6
      monitoring/src/main/resources/root_rule_chain.json
  6. 2
      msa/black-box-tests/src/test/resources/MqttRuleNodeTestMetadata.json
  7. 52
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java
  8. 28
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNodeConfiguration.java
  9. 6
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/DeduplicateProcessingStrategy.java
  10. 10
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/OnEveryMessageProcessingStrategy.java
  11. 22
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/ProcessingStrategy.java
  12. 10
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/SkipProcessingStrategy.java
  13. 72
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNodeTest.java
  14. 76
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/DeduplicateProcessingStrategyTest.java
  15. 6
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/OnEveryMessageProcessingStrategyTest.java
  16. 14
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/ProcessingStrategyTest.java
  17. 6
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/SkipProcessingStrategyTest.java
  18. 22
      ui-ngx/src/app/modules/home/components/rule-node/action/advanced-persistence-setting-row.component.ts
  19. 5
      ui-ngx/src/app/modules/home/components/rule-node/action/advanced-persistence-setting.component.html
  20. 5
      ui-ngx/src/app/modules/home/components/rule-node/action/advanced-persistence-setting.component.ts
  21. 13
      ui-ngx/src/app/modules/home/components/rule-node/action/timeseries-config.component.html
  22. 78
      ui-ngx/src/app/modules/home/components/rule-node/action/timeseries-config.component.ts
  23. 54
      ui-ngx/src/app/modules/home/components/rule-node/action/timeseries-config.models.ts
  24. 6
      ui-ngx/src/app/modules/home/components/rule-node/common/example-hint.component.html
  25. 31
      ui-ngx/src/app/modules/home/components/rule-node/common/example-hint.component.scss
  26. 2
      ui-ngx/src/app/modules/home/components/rule-node/common/example-hint.component.ts
  27. 29
      ui-ngx/src/assets/help/en_US/rulenode/save_timeseries_node_advanced.md
  28. 7
      ui-ngx/src/assets/locale/locale.constant-en_US.json

2
application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json

@ -37,7 +37,7 @@
"configuration": {
"defaultTTL": 0,
"useServerTs": false,
"persistenceSettings": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
}
},

2
application/src/main/data/json/tenant/device_profile/rule_chain_template.json

@ -23,7 +23,7 @@
"configuration": {
"defaultTTL": 0,
"useServerTs": false,
"persistenceSettings": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
}
}

2
application/src/main/data/json/tenant/rule_chains/root_rule_chain.json

@ -22,7 +22,7 @@
"configuration": {
"defaultTTL": 0,
"useServerTs": false,
"persistenceSettings": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
}
}

4
application/src/main/data/upgrade/basic/schema_update.sql

@ -29,7 +29,7 @@ DO $$
SET configuration = (
(configuration::jsonb - 'skipLatestPersistence')
|| jsonb_build_object(
'persistenceSettings', jsonb_build_object(
'processingSettings', jsonb_build_object(
'type', 'ADVANCED',
'timeseries', jsonb_build_object('type', 'ON_EVERY_MESSAGE'),
'latest', jsonb_build_object('type', 'SKIP'),
@ -46,7 +46,7 @@ DO $$
SET configuration = (
(configuration::jsonb - 'skipLatestPersistence')
|| jsonb_build_object(
'persistenceSettings', jsonb_build_object(
'processingSettings', jsonb_build_object(
'type', 'ON_EVERY_MESSAGE'
)
)

6
monitoring/src/main/resources/root_rule_chain.json

@ -25,7 +25,7 @@
"configuration": {
"defaultTTL": 0,
"useServerTs": false,
"persistenceSettings": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
}
},
@ -281,7 +281,7 @@
"configuration": {
"defaultTTL": 0,
"useServerTs": false,
"persistenceSettings": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
}
},
@ -319,7 +319,7 @@
"configuration": {
"defaultTTL": 180,
"useServerTs": false,
"persistenceSettings": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
}
},

2
msa/black-box-tests/src/test/resources/MqttRuleNodeTestMetadata.json

@ -40,7 +40,7 @@
"configuration": {
"defaultTTL": 0,
"useServerTs": false,
"persistenceSettings": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
}
},

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

@ -27,7 +27,7 @@ import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.rule.engine.telemetry.strategy.PersistenceStrategy;
import org.thingsboard.rule.engine.telemetry.strategy.ProcessingStrategy;
import org.thingsboard.server.common.adaptor.JsonConverter;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.TenantProfile;
@ -45,11 +45,11 @@ import java.util.Map;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.PersistenceSettings;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Advanced;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Deduplicate;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.PersistenceSettings.OnEveryMessage;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.PersistenceSettings.WebSocketsOnly;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.ProcessingSettings;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Advanced;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Deduplicate;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.ProcessingSettings.OnEveryMessage;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.ProcessingSettings.WebSocketsOnly;
import static org.thingsboard.server.common.data.msg.TbMsgType.POST_TELEMETRY_REQUEST;
@Slf4j
@ -58,7 +58,7 @@ import static org.thingsboard.server.common.data.msg.TbMsgType.POST_TELEMETRY_RE
name = "save time series",
configClazz = TbMsgTimeseriesNodeConfiguration.class,
nodeDescription = """
Saves time series data with a configurable TTL and according to configured persistence strategies.
Saves time series data with a configurable TTL and according to configured processing strategies.
""",
nodeDetails = """
Node performs three <strong>actions:</strong>
@ -68,14 +68,14 @@ import static org.thingsboard.server.common.data.msg.TbMsgType.POST_TELEMETRY_RE
<li><strong>WebSockets:</strong> notify WebSockets subscriptions about time series data updates.</li>
</ul>
For each <em>action</em>, three <strong>persistence strategies</strong> are available:
For each <em>action</em>, three <strong>processing strategies</strong> are available:
<ul>
<li><strong>On every message:</strong> perform the action for every message.</li>
<li><strong>Deduplicate:</strong> perform the action only for the first message from a particular originator within a configurable interval.</li>
<li><strong>Skip:</strong> never perform the action.</li>
</ul>
<strong>Persistence strategies</strong> are configured using <em>persistence settings</em>, which support two modes:
<strong>Processing strategies</strong> are configured using <em>processing settings</em>, which support two modes:
<ul>
<li><strong>Basic</strong>
<ul>
@ -110,7 +110,7 @@ public class TbMsgTimeseriesNode implements TbNode {
private TbContext ctx;
private long tenantProfileDefaultStorageTtl;
private PersistenceSettings persistenceSettings;
private ProcessingSettings processingSettings;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
@ -118,7 +118,7 @@ public class TbMsgTimeseriesNode implements TbNode {
this.ctx = ctx;
ctx.addTenantProfileListener(this::onTenantProfileUpdate);
onTenantProfileUpdate(ctx.getTenantProfile());
persistenceSettings = config.getPersistenceSettings();
processingSettings = config.getProcessingSettings();
}
private void onTenantProfileUpdate(TenantProfile tenantProfile) {
@ -175,25 +175,25 @@ public class TbMsgTimeseriesNode implements TbNode {
}
private TimeseriesSaveRequest.Strategy determineSaveStrategy(long ts, UUID originatorUuid) {
if (persistenceSettings instanceof OnEveryMessage) {
if (processingSettings instanceof OnEveryMessage) {
return TimeseriesSaveRequest.Strategy.SAVE_ALL;
}
if (persistenceSettings instanceof WebSocketsOnly) {
if (processingSettings instanceof WebSocketsOnly) {
return TimeseriesSaveRequest.Strategy.WS_ONLY;
}
if (persistenceSettings instanceof Deduplicate deduplicate) {
boolean isFirstMsgInInterval = deduplicate.getPersistenceStrategy().shouldPersist(ts, originatorUuid);
if (processingSettings instanceof Deduplicate deduplicate) {
boolean isFirstMsgInInterval = deduplicate.getProcessingStrategy().shouldProcess(ts, originatorUuid);
return isFirstMsgInInterval ? TimeseriesSaveRequest.Strategy.SAVE_ALL : TimeseriesSaveRequest.Strategy.SKIP_ALL;
}
if (persistenceSettings instanceof Advanced advanced) {
if (processingSettings instanceof Advanced advanced) {
return new TimeseriesSaveRequest.Strategy(
advanced.timeseries().shouldPersist(ts, originatorUuid),
advanced.latest().shouldPersist(ts, originatorUuid),
advanced.webSockets().shouldPersist(ts, originatorUuid)
advanced.timeseries().shouldProcess(ts, originatorUuid),
advanced.latest().shouldProcess(ts, originatorUuid),
advanced.webSockets().shouldProcess(ts, originatorUuid)
);
}
// should not happen
throw new IllegalArgumentException("Unknown persistence settings type: " + persistenceSettings.getClass().getSimpleName());
throw new IllegalArgumentException("Unknown processing settings type: " + processingSettings.getClass().getSimpleName());
}
@Override
@ -209,14 +209,14 @@ public class TbMsgTimeseriesNode implements TbNode {
hasChanges = true;
JsonNode skipLatestPersistence = oldConfiguration.get("skipLatestPersistence");
if (skipLatestPersistence != null && "true".equals(skipLatestPersistence.asText())) {
var skipLatestPersistenceSettings = new Advanced(
PersistenceStrategy.onEveryMessage(),
PersistenceStrategy.skip(),
PersistenceStrategy.onEveryMessage()
var skipLatestProcessingSettings = new Advanced(
ProcessingStrategy.onEveryMessage(),
ProcessingStrategy.skip(),
ProcessingStrategy.onEveryMessage()
);
((ObjectNode) oldConfiguration).set("persistenceSettings", JacksonUtil.valueToTree(skipLatestPersistenceSettings));
((ObjectNode) oldConfiguration).set("processingSettings", JacksonUtil.valueToTree(skipLatestProcessingSettings));
} else {
((ObjectNode) oldConfiguration).set("persistenceSettings", JacksonUtil.valueToTree(new OnEveryMessage()));
((ObjectNode) oldConfiguration).set("processingSettings", JacksonUtil.valueToTree(new OnEveryMessage()));
}
((ObjectNode) oldConfiguration).remove("skipLatestPersistence");
break;

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

@ -24,14 +24,14 @@ import jakarta.validation.constraints.NotNull;
import lombok.Data;
import lombok.Getter;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import org.thingsboard.rule.engine.telemetry.strategy.PersistenceStrategy;
import org.thingsboard.rule.engine.telemetry.strategy.ProcessingStrategy;
import java.util.Objects;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Advanced;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Deduplicate;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.PersistenceSettings.OnEveryMessage;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.PersistenceSettings.WebSocketsOnly;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Advanced;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Deduplicate;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.ProcessingSettings.OnEveryMessage;
import static org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNodeConfiguration.ProcessingSettings.WebSocketsOnly;
@Data
public class TbMsgTimeseriesNodeConfiguration implements NodeConfiguration<TbMsgTimeseriesNodeConfiguration> {
@ -39,14 +39,14 @@ public class TbMsgTimeseriesNodeConfiguration implements NodeConfiguration<TbMsg
private long defaultTTL;
private boolean useServerTs;
@NotNull
private PersistenceSettings persistenceSettings;
private TbMsgTimeseriesNodeConfiguration.ProcessingSettings processingSettings;
@Override
public TbMsgTimeseriesNodeConfiguration defaultConfiguration() {
TbMsgTimeseriesNodeConfiguration configuration = new TbMsgTimeseriesNodeConfiguration();
configuration.setDefaultTTL(0L);
configuration.setUseServerTs(false);
configuration.setPersistenceSettings(new OnEveryMessage());
configuration.setProcessingSettings(new OnEveryMessage());
return configuration;
}
@ -61,29 +61,29 @@ public class TbMsgTimeseriesNodeConfiguration implements NodeConfiguration<TbMsg
@JsonSubTypes.Type(value = Deduplicate.class, name = "DEDUPLICATE"),
@JsonSubTypes.Type(value = Advanced.class, name = "ADVANCED")
})
sealed interface PersistenceSettings permits OnEveryMessage, Deduplicate, WebSocketsOnly, Advanced {
sealed interface ProcessingSettings permits OnEveryMessage, Deduplicate, WebSocketsOnly, Advanced {
record OnEveryMessage() implements PersistenceSettings {}
record OnEveryMessage() implements ProcessingSettings {}
record WebSocketsOnly() implements PersistenceSettings {}
record WebSocketsOnly() implements ProcessingSettings {}
@Getter
final class Deduplicate implements PersistenceSettings {
final class Deduplicate implements ProcessingSettings {
private final int deduplicationIntervalSecs;
@JsonIgnore
private final PersistenceStrategy persistenceStrategy;
private final ProcessingStrategy processingStrategy;
@JsonCreator
Deduplicate(@JsonProperty("deduplicationIntervalSecs") int deduplicationIntervalSecs) {
this.deduplicationIntervalSecs = deduplicationIntervalSecs;
persistenceStrategy = PersistenceStrategy.deduplicate(deduplicationIntervalSecs);
processingStrategy = ProcessingStrategy.deduplicate(deduplicationIntervalSecs);
}
}
record Advanced(PersistenceStrategy timeseries, PersistenceStrategy latest, PersistenceStrategy webSockets) implements PersistenceSettings {
record Advanced(ProcessingStrategy timeseries, ProcessingStrategy latest, ProcessingStrategy webSockets) implements ProcessingSettings {
public Advanced {
Objects.requireNonNull(timeseries);

6
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/DeduplicatePersistenceStrategy.java → rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/DeduplicateProcessingStrategy.java

@ -26,7 +26,7 @@ import java.time.Duration;
import java.util.Set;
import java.util.UUID;
final class DeduplicatePersistenceStrategy implements PersistenceStrategy {
final class DeduplicateProcessingStrategy implements ProcessingStrategy {
private static final int MIN_DEDUPLICATION_INTERVAL_SECS = 1;
private static final int MAX_DEDUPLICATION_INTERVAL_SECS = (int) Duration.ofDays(1L).toSeconds();
@ -42,7 +42,7 @@ final class DeduplicatePersistenceStrategy implements PersistenceStrategy {
private final LoadingCache<Long, Set<UUID>> deduplicationCache;
@JsonCreator
public DeduplicatePersistenceStrategy(@JsonProperty("deduplicationIntervalSecs") int deduplicationIntervalSecs) {
public DeduplicateProcessingStrategy(@JsonProperty("deduplicationIntervalSecs") int deduplicationIntervalSecs) {
if (deduplicationIntervalSecs < MIN_DEDUPLICATION_INTERVAL_SECS || deduplicationIntervalSecs > MAX_DEDUPLICATION_INTERVAL_SECS) {
throw new IllegalArgumentException("Deduplication interval must be at least " + MIN_DEDUPLICATION_INTERVAL_SECS + " second(s) " +
"and at most " + MAX_DEDUPLICATION_INTERVAL_SECS + " second(s), was " + deduplicationIntervalSecs + " second(s)");
@ -81,7 +81,7 @@ final class DeduplicatePersistenceStrategy implements PersistenceStrategy {
}
@Override
public boolean shouldPersist(long ts, UUID originatorUuid) {
public boolean shouldProcess(long ts, UUID originatorUuid) {
long intervalNumber = ts / deduplicationIntervalMillis;
return deduplicationCache.get(intervalNumber).add(originatorUuid);
}

10
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/OnEveryMessagePersistenceStrategy.java → rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/OnEveryMessageProcessingStrategy.java

@ -19,19 +19,19 @@ import com.fasterxml.jackson.annotation.JsonCreator;
import java.util.UUID;
final class OnEveryMessagePersistenceStrategy implements PersistenceStrategy {
final class OnEveryMessageProcessingStrategy implements ProcessingStrategy {
private static final OnEveryMessagePersistenceStrategy INSTANCE = new OnEveryMessagePersistenceStrategy();
private static final OnEveryMessageProcessingStrategy INSTANCE = new OnEveryMessageProcessingStrategy();
private OnEveryMessagePersistenceStrategy() {}
private OnEveryMessageProcessingStrategy() {}
@JsonCreator
public static OnEveryMessagePersistenceStrategy getInstance() {
public static OnEveryMessageProcessingStrategy getInstance() {
return INSTANCE;
}
@Override
public boolean shouldPersist(long ts, UUID originatorUuid) {
public boolean shouldProcess(long ts, UUID originatorUuid) {
return true;
}

22
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/PersistenceStrategy.java → rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/ProcessingStrategy.java

@ -26,24 +26,24 @@ import java.util.UUID;
property = "type"
)
@JsonSubTypes({
@JsonSubTypes.Type(value = OnEveryMessagePersistenceStrategy.class, name = "ON_EVERY_MESSAGE"),
@JsonSubTypes.Type(value = DeduplicatePersistenceStrategy.class, name = "DEDUPLICATE"),
@JsonSubTypes.Type(value = SkipPersistenceStrategy.class, name = "SKIP")
@JsonSubTypes.Type(value = OnEveryMessageProcessingStrategy.class, name = "ON_EVERY_MESSAGE"),
@JsonSubTypes.Type(value = DeduplicateProcessingStrategy.class, name = "DEDUPLICATE"),
@JsonSubTypes.Type(value = SkipProcessingStrategy.class, name = "SKIP")
})
public sealed interface PersistenceStrategy permits OnEveryMessagePersistenceStrategy, DeduplicatePersistenceStrategy, SkipPersistenceStrategy {
public sealed interface ProcessingStrategy permits OnEveryMessageProcessingStrategy, DeduplicateProcessingStrategy, SkipProcessingStrategy {
static PersistenceStrategy onEveryMessage() {
return OnEveryMessagePersistenceStrategy.getInstance();
static ProcessingStrategy onEveryMessage() {
return OnEveryMessageProcessingStrategy.getInstance();
}
static PersistenceStrategy deduplicate(int deduplicationIntervalSecs) {
return new DeduplicatePersistenceStrategy(deduplicationIntervalSecs);
static ProcessingStrategy deduplicate(int deduplicationIntervalSecs) {
return new DeduplicateProcessingStrategy(deduplicationIntervalSecs);
}
static PersistenceStrategy skip() {
return SkipPersistenceStrategy.getInstance();
static ProcessingStrategy skip() {
return SkipProcessingStrategy.getInstance();
}
boolean shouldPersist(long ts, UUID originatorUuid);
boolean shouldProcess(long ts, UUID originatorUuid);
}

10
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/SkipPersistenceStrategy.java → rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/strategy/SkipProcessingStrategy.java

@ -19,19 +19,19 @@ import com.fasterxml.jackson.annotation.JsonCreator;
import java.util.UUID;
final class SkipPersistenceStrategy implements PersistenceStrategy {
final class SkipProcessingStrategy implements ProcessingStrategy {
private static final SkipPersistenceStrategy INSTANCE = new SkipPersistenceStrategy();
private static final SkipProcessingStrategy INSTANCE = new SkipProcessingStrategy();
private SkipPersistenceStrategy() {}
private SkipProcessingStrategy() {}
@JsonCreator
public static SkipPersistenceStrategy getInstance() {
public static SkipProcessingStrategy getInstance() {
return INSTANCE;
}
@Override
public boolean shouldPersist(long ts, UUID originatorUuid) {
public boolean shouldProcess(long ts, UUID originatorUuid) {
return false;
}

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

@ -34,7 +34,7 @@ import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest;
import org.thingsboard.rule.engine.telemetry.strategy.PersistenceStrategy;
import org.thingsboard.rule.engine.telemetry.strategy.ProcessingStrategy;
import org.thingsboard.server.common.adaptor.JsonConverter;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.DeviceId;
@ -110,7 +110,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
@Test
public void verifyDefaultConfig() {
assertThat(config.getDefaultTTL()).isEqualTo(0L);
assertThat(config.getPersistenceSettings()).isInstanceOf(TbMsgTimeseriesNodeConfiguration.PersistenceSettings.OnEveryMessage.class);
assertThat(config.getProcessingSettings()).isInstanceOf(TbMsgTimeseriesNodeConfiguration.ProcessingSettings.OnEveryMessage.class);
assertThat(config.isUseServerTs()).isFalse();
}
@ -124,14 +124,14 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
}
@Test
public void givenPersistenceSettingsAreNull_whenValidatingConstraints_thenThrowsException() {
public void givenProcessingSettingsAreNull_whenValidatingConstraints_thenThrowsException() {
// GIVEN
config.setPersistenceSettings(null);
config.setProcessingSettings(null);
// WHEN-THEN
assertThatThrownBy(() -> ConstraintValidator.validateFields(config))
.isInstanceOf(DataValidationException.class)
.hasMessage("Validation error: persistenceSettings must not be null");
.hasMessage("Validation error: processingSettings must not be null");
}
@ParameterizedTest
@ -216,15 +216,15 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
}
@Test
public void givenSkipLatestPersistenceSettingsAndTtlFromConfig_whenOnMsg_thenSaveTimeseriesUsingTtlFromConfig() throws TbNodeException {
public void givenSkipLatestProcessingSettingsAndTtlFromConfig_whenOnMsg_thenSaveTimeseriesUsingTtlFromConfig() throws TbNodeException {
// GIVEN
config.setDefaultTTL(10L);
var timeseriesStrategy = PersistenceStrategy.onEveryMessage();
var latestStrategy = PersistenceStrategy.skip();
var webSockets = PersistenceStrategy.onEveryMessage();
var persistenceSettings = new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Advanced(timeseriesStrategy, latestStrategy, webSockets);
config.setPersistenceSettings(persistenceSettings);
var timeseriesStrategy = ProcessingStrategy.onEveryMessage();
var latestStrategy = ProcessingStrategy.skip();
var webSockets = ProcessingStrategy.onEveryMessage();
var processingSettings = new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Advanced(timeseriesStrategy, latestStrategy, webSockets);
config.setProcessingSettings(processingSettings);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -333,9 +333,9 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
}
@Test
public void givenOnEveryMessagePersistenceSettingsAndSameMessageTwoTimes_whenOnMsg_thenPersistSameMessageTwoTimes() throws TbNodeException {
public void givenOnEveryMessageProcessingSettingsAndSameMessageTwoTimes_whenOnMsg_thenPersistSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.OnEveryMessage());
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.OnEveryMessage());
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -368,9 +368,9 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
}
@Test
public void givenDeduplicatePersistenceSettingsAndSameMessageTwoTimes_whenOnMsg_thenPersistThisMessageOnlyFirstTime() throws TbNodeException {
public void givenDeduplicateProcessingSettingsAndSameMessageTwoTimes_whenOnMsg_thenPersistThisMessageOnlyFirstTime() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Deduplicate(10));
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Deduplicate(10));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -403,9 +403,9 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
}
@Test
public void givenWebsocketsOnlyPersistenceSettingsAndSameMessageTwoTimes_whenOnMsg_thenSendsOnlyWsUpdateTwoTimes() throws TbNodeException {
public void givenWebSocketsOnlyProcessingSettingsAndSameMessageTwoTimes_whenOnMsg_thenSendsOnlyWsUpdateTwoTimes() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.WebSocketsOnly());
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.WebSocketsOnly());
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -438,12 +438,12 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
}
@Test
public void givenAdvancedPersistenceSettingsWithOnEveryMessageStrategiesForAllActionsAndSameMessageTwoTimes_whenOnMsg_thenPersistSameMessageTwoTimes() throws TbNodeException {
public void givenAdvancedProcessingSettingsWithOnEveryMessageStrategiesForAllActionsAndSameMessageTwoTimes_whenOnMsg_thenPersistSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Advanced(
PersistenceStrategy.onEveryMessage(),
PersistenceStrategy.onEveryMessage(),
PersistenceStrategy.onEveryMessage()
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Advanced(
ProcessingStrategy.onEveryMessage(),
ProcessingStrategy.onEveryMessage(),
ProcessingStrategy.onEveryMessage()
));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -477,12 +477,12 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
}
@Test
public void givenAdvancedPersistenceSettingsWithDifferentDeduplicateStrategyForEachAction_whenOnMsg_thenEvaluatesStrategiesForEachActionsIndependently() throws TbNodeException {
public void givenAdvancedProcessingSettingsWithDifferentDeduplicateStrategyForEachAction_whenOnMsg_thenEvaluatesStrategiesForEachActionsIndependently() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Advanced(
PersistenceStrategy.deduplicate(1),
PersistenceStrategy.deduplicate(2),
PersistenceStrategy.deduplicate(3)
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Advanced(
ProcessingStrategy.deduplicate(1),
ProcessingStrategy.deduplicate(2),
ProcessingStrategy.deduplicate(3)
));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -528,12 +528,12 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
}
@Test
public void givenAdvancedPersistenceSettingsWithSkipStrategiesForAllActionsAndSameMessageTwoTimes_whenOnMsg_thenSkipsSameMessageTwoTimes() throws TbNodeException {
public void givenAdvancedProcessingSettingsWithSkipStrategiesForAllActionsAndSameMessageTwoTimes_whenOnMsg_thenSkipsSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setPersistenceSettings(new TbMsgTimeseriesNodeConfiguration.PersistenceSettings.Advanced(
PersistenceStrategy.skip(),
PersistenceStrategy.skip(),
PersistenceStrategy.skip()
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Advanced(
ProcessingStrategy.skip(),
ProcessingStrategy.skip(),
ProcessingStrategy.skip()
));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -577,7 +577,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
{
"defaultTTL": 0,
"useServerTs": false,
"persistenceSettings": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
}
}"""),
@ -591,7 +591,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
{
"defaultTTL": 0,
"useServerTs": false,
"persistenceSettings": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
}
}"""),
@ -606,7 +606,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
{
"defaultTTL": 0,
"useServerTs": false,
"persistenceSettings": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
}
}"""),
@ -621,7 +621,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
{
"defaultTTL": 0,
"useServerTs": false,
"persistenceSettings": {
"processingSettings": {
"type": "ADVANCED",
"timeseries": {
"type": "ON_EVERY_MESSAGE"

76
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/DeduplicatePersistenceStrategyTest.java → rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/DeduplicateProcessingStrategyTest.java

@ -28,27 +28,27 @@ import java.util.UUID;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
class DeduplicatePersistenceStrategyTest {
class DeduplicateProcessingStrategyTest {
final int deduplicationIntervalSecs = 10;
DeduplicatePersistenceStrategy strategy;
DeduplicateProcessingStrategy strategy;
@BeforeEach
void setup() {
strategy = new DeduplicatePersistenceStrategy(deduplicationIntervalSecs);
strategy = new DeduplicateProcessingStrategy(deduplicationIntervalSecs);
}
@Test
void shouldThrowWhenDeduplicationIntervalIsLessThanOneSecond() {
assertThatThrownBy(() -> new DeduplicatePersistenceStrategy(0))
assertThatThrownBy(() -> new DeduplicateProcessingStrategy(0))
.isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("Deduplication interval must be at least 1 second(s) and at most 86400 second(s), was 0 second(s)");
}
@Test
void shouldThrowWhenDeduplicationIntervalIsMoreThan24Hours() {
assertThatThrownBy(() -> new DeduplicatePersistenceStrategy(86401))
assertThatThrownBy(() -> new DeduplicateProcessingStrategy(86401))
.isInstanceOf(IllegalArgumentException.class)
.hasMessageContaining("Deduplication interval must be at least 1 second(s) and at most 86400 second(s), was 86401 second(s)");
}
@ -59,7 +59,7 @@ class DeduplicatePersistenceStrategyTest {
int deduplicationIntervalSecs = 1; // min deduplication interval duration
// WHEN
strategy = new DeduplicatePersistenceStrategy(deduplicationIntervalSecs);
strategy = new DeduplicateProcessingStrategy(deduplicationIntervalSecs);
// THEN
var deduplicationCache = (LoadingCache<Long, Set<UUID>>) ReflectionTestUtils.getField(strategy, "deduplicationCache");
@ -76,7 +76,7 @@ class DeduplicatePersistenceStrategyTest {
int deduplicationIntervalSecs = (int) Duration.ofHours(1L).toSeconds(); // max deduplication interval duration
// WHEN
strategy = new DeduplicatePersistenceStrategy(deduplicationIntervalSecs);
strategy = new DeduplicateProcessingStrategy(deduplicationIntervalSecs);
// THEN
var deduplicationCache = (LoadingCache<Long, Set<UUID>>) ReflectionTestUtils.getField(strategy, "deduplicationCache");
@ -93,7 +93,7 @@ class DeduplicatePersistenceStrategyTest {
int deduplicationIntervalSecs = (int) Duration.ofDays(1L).toSeconds(); // max deduplication interval duration
// WHEN
strategy = new DeduplicatePersistenceStrategy(deduplicationIntervalSecs);
strategy = new DeduplicateProcessingStrategy(deduplicationIntervalSecs);
// THEN
var deduplicationCache = (LoadingCache<Long, Set<UUID>>) ReflectionTestUtils.getField(strategy, "deduplicationCache");
@ -110,7 +110,7 @@ class DeduplicatePersistenceStrategyTest {
int deduplicationIntervalSecs = 1; // min deduplication interval duration
// WHEN
strategy = new DeduplicatePersistenceStrategy(deduplicationIntervalSecs);
strategy = new DeduplicateProcessingStrategy(deduplicationIntervalSecs);
// THEN
var deduplicationCache = (LoadingCache<Long, Set<UUID>>) ReflectionTestUtils.getField(strategy, "deduplicationCache");
@ -127,7 +127,7 @@ class DeduplicatePersistenceStrategyTest {
int deduplicationIntervalSecs = (int) Duration.ofHours(1L).toSeconds();
// WHEN
strategy = new DeduplicatePersistenceStrategy(deduplicationIntervalSecs);
strategy = new DeduplicateProcessingStrategy(deduplicationIntervalSecs);
// THEN
var deduplicationCache = (LoadingCache<Long, Set<UUID>>) ReflectionTestUtils.getField(strategy, "deduplicationCache");
@ -144,7 +144,7 @@ class DeduplicatePersistenceStrategyTest {
int deduplicationIntervalSecs = (int) Duration.ofDays(1L).toSeconds(); // max deduplication interval duration
// WHEN
strategy = new DeduplicatePersistenceStrategy(deduplicationIntervalSecs);
strategy = new DeduplicateProcessingStrategy(deduplicationIntervalSecs);
// THEN
var deduplicationCache = (LoadingCache<Long, Set<UUID>>) ReflectionTestUtils.getField(strategy, "deduplicationCache");
@ -160,7 +160,7 @@ class DeduplicatePersistenceStrategyTest {
long ts = 1_000_000L;
UUID originator = UUID.randomUUID();
assertThat(strategy.shouldPersist(ts, originator)).isTrue();
assertThat(strategy.shouldProcess(ts, originator)).isTrue();
}
@Test
@ -169,11 +169,11 @@ class DeduplicatePersistenceStrategyTest {
UUID originator = UUID.randomUUID();
// Initial call should return true
assertThat(strategy.shouldPersist(baseTs, originator)).isTrue();
assertThat(strategy.shouldProcess(baseTs, originator)).isTrue();
// Subsequent call within the same interval should return false for the same originator
long withinSameIntervalTs = baseTs + 1000L;
assertThat(strategy.shouldPersist(withinSameIntervalTs, originator)).isFalse();
assertThat(strategy.shouldProcess(withinSameIntervalTs, originator)).isFalse();
}
@Test
@ -183,12 +183,12 @@ class DeduplicatePersistenceStrategyTest {
UUID originator2 = UUID.randomUUID();
// First call for different originators in the same interval should return true independently
assertThat(strategy.shouldPersist(baseTs, originator1)).isTrue();
assertThat(strategy.shouldPersist(baseTs, originator2)).isTrue();
assertThat(strategy.shouldProcess(baseTs, originator1)).isTrue();
assertThat(strategy.shouldProcess(baseTs, originator2)).isTrue();
// Subsequent calls for the same originators within the same interval should return false
assertThat(strategy.shouldPersist(baseTs + 500L, originator1)).isFalse();
assertThat(strategy.shouldPersist(baseTs + 500L, originator2)).isFalse();
assertThat(strategy.shouldProcess(baseTs + 500L, originator1)).isFalse();
assertThat(strategy.shouldProcess(baseTs + 500L, originator2)).isFalse();
}
@Test
@ -197,11 +197,11 @@ class DeduplicatePersistenceStrategyTest {
long maxTs = Long.MAX_VALUE;
UUID originator = UUID.randomUUID();
assertThat(strategy.shouldPersist(minTs, originator)).isTrue();
assertThat(strategy.shouldPersist(minTs + 1L, originator)).isFalse();
assertThat(strategy.shouldProcess(minTs, originator)).isTrue();
assertThat(strategy.shouldProcess(minTs + 1L, originator)).isFalse();
assertThat(strategy.shouldPersist(maxTs, originator)).isTrue();
assertThat(strategy.shouldPersist(maxTs - 1L, originator)).isFalse();
assertThat(strategy.shouldProcess(maxTs, originator)).isTrue();
assertThat(strategy.shouldProcess(maxTs - 1L, originator)).isFalse();
}
@Test
@ -213,22 +213,22 @@ class DeduplicatePersistenceStrategyTest {
long firstIntervalEnd = firstIntervalStart + Duration.ofSeconds(deduplicationIntervalSecs).toMillis() - 1L;
long firstIntervalMiddle = calculateMiddle(firstIntervalStart, firstIntervalEnd);
assertThat(strategy.shouldPersist(firstIntervalStart, originator)).isTrue();
assertThat(strategy.shouldPersist(firstIntervalStart + 1, originator)).isFalse();
assertThat(strategy.shouldPersist(firstIntervalMiddle, originator)).isFalse();
assertThat(strategy.shouldPersist(firstIntervalEnd - 1, originator)).isFalse();
assertThat(strategy.shouldPersist(firstIntervalEnd, originator)).isFalse();
assertThat(strategy.shouldProcess(firstIntervalStart, originator)).isTrue();
assertThat(strategy.shouldProcess(firstIntervalStart + 1, originator)).isFalse();
assertThat(strategy.shouldProcess(firstIntervalMiddle, originator)).isFalse();
assertThat(strategy.shouldProcess(firstIntervalEnd - 1, originator)).isFalse();
assertThat(strategy.shouldProcess(firstIntervalEnd, originator)).isFalse();
// check 2nd interval
long secondIntervalStart = firstIntervalEnd + 1L;
long secondIntervalEnd = secondIntervalStart + Duration.ofSeconds(deduplicationIntervalSecs).toMillis() - 1L;
long secondIntervalMiddle = calculateMiddle(secondIntervalStart, secondIntervalEnd);
assertThat(strategy.shouldPersist(secondIntervalStart, originator)).isTrue();
assertThat(strategy.shouldPersist(secondIntervalStart + 1, originator)).isFalse();
assertThat(strategy.shouldPersist(secondIntervalMiddle, originator)).isFalse();
assertThat(strategy.shouldPersist(secondIntervalEnd - 1, originator)).isFalse();
assertThat(strategy.shouldPersist(secondIntervalEnd, originator)).isFalse();
assertThat(strategy.shouldProcess(secondIntervalStart, originator)).isTrue();
assertThat(strategy.shouldProcess(secondIntervalStart + 1, originator)).isFalse();
assertThat(strategy.shouldProcess(secondIntervalMiddle, originator)).isFalse();
assertThat(strategy.shouldProcess(secondIntervalEnd - 1, originator)).isFalse();
assertThat(strategy.shouldProcess(secondIntervalEnd, originator)).isFalse();
}
@Test
@ -238,19 +238,19 @@ class DeduplicatePersistenceStrategyTest {
long baseTs = 0L;
// First interval for both originators
assertThat(strategy.shouldPersist(baseTs, originator1)).isTrue();
assertThat(strategy.shouldPersist(baseTs, originator2)).isTrue();
assertThat(strategy.shouldProcess(baseTs, originator1)).isTrue();
assertThat(strategy.shouldProcess(baseTs, originator2)).isTrue();
// Move to the next interval
long nextIntervalTs = baseTs + Duration.ofSeconds(10).toMillis();
// Each originator should be allowed again in the new interval
assertThat(strategy.shouldPersist(nextIntervalTs, originator1)).isTrue();
assertThat(strategy.shouldPersist(nextIntervalTs, originator2)).isTrue();
assertThat(strategy.shouldProcess(nextIntervalTs, originator1)).isTrue();
assertThat(strategy.shouldProcess(nextIntervalTs, originator2)).isTrue();
// Subsequent calls in the same new interval should return false
assertThat(strategy.shouldPersist(nextIntervalTs + 500L, originator1)).isFalse();
assertThat(strategy.shouldPersist(nextIntervalTs + 500L, originator2)).isFalse();
assertThat(strategy.shouldProcess(nextIntervalTs + 500L, originator1)).isFalse();
assertThat(strategy.shouldProcess(nextIntervalTs + 500L, originator2)).isFalse();
}
private static long calculateMiddle(long start, long end) {

6
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/OnEveryMessagePersistenceStrategyTest.java → rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/OnEveryMessageProcessingStrategyTest.java

@ -24,13 +24,13 @@ import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
class OnEveryMessagePersistenceStrategyTest {
class OnEveryMessageProcessingStrategyTest {
@ParameterizedTest
@MethodSource("edgeCaseProvider")
void shouldAlwaysReturnTrueForAnyInput(long timestamp, UUID originator) {
var onEveryMessage = OnEveryMessagePersistenceStrategy.getInstance();
assertThat(onEveryMessage.shouldPersist(timestamp, originator)).isTrue();
var onEveryMessage = OnEveryMessageProcessingStrategy.getInstance();
assertThat(onEveryMessage.shouldProcess(timestamp, originator)).isTrue();
}
private static Stream<Arguments> edgeCaseProvider() {

14
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/PersistenceStrategyTest.java → rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/ProcessingStrategyTest.java

@ -22,23 +22,23 @@ import java.time.Duration;
import static org.assertj.core.api.Assertions.assertThat;
class PersistenceStrategyTest {
class ProcessingStrategyTest {
@Test
void testOnEveryMessageReturnsCorrectInstance() {
PersistenceStrategy strategy = PersistenceStrategy.onEveryMessage();
ProcessingStrategy strategy = ProcessingStrategy.onEveryMessage();
assertThat(strategy)
.isNotNull()
.isInstanceOf(OnEveryMessagePersistenceStrategy.class);
.isInstanceOf(OnEveryMessageProcessingStrategy.class);
}
@Test
void testDeduplicateReturnsCorrectInstance() {
int validDeduplicationIntervalSecs = 5;
PersistenceStrategy strategy = PersistenceStrategy.deduplicate(validDeduplicationIntervalSecs);
ProcessingStrategy strategy = ProcessingStrategy.deduplicate(validDeduplicationIntervalSecs);
assertThat(strategy)
.isNotNull()
.isInstanceOf(DeduplicatePersistenceStrategy.class);
.isInstanceOf(DeduplicateProcessingStrategy.class);
long actualDeduplicationIntervalMillis = (long) ReflectionTestUtils.getField(strategy, "deduplicationIntervalMillis");
assertThat(actualDeduplicationIntervalMillis).isEqualTo(Duration.ofSeconds(validDeduplicationIntervalSecs).toMillis());
@ -46,10 +46,10 @@ class PersistenceStrategyTest {
@Test
void testSkipReturnsCorrectInstance() {
PersistenceStrategy strategy = PersistenceStrategy.skip();
ProcessingStrategy strategy = ProcessingStrategy.skip();
assertThat(strategy)
.isNotNull()
.isInstanceOf(SkipPersistenceStrategy.class);
.isInstanceOf(SkipProcessingStrategy.class);
}
}

6
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/SkipPersistenceStrategyTest.java → rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/strategy/SkipProcessingStrategyTest.java

@ -24,13 +24,13 @@ import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
class SkipPersistenceStrategyTest {
class SkipProcessingStrategyTest {
@ParameterizedTest
@MethodSource("edgeCaseProvider")
void shouldAlwaysReturnFalseForAnyInput(long timestamp, UUID originator) {
var skipStrategy = SkipPersistenceStrategy.getInstance();
assertThat(skipStrategy.shouldPersist(timestamp, originator)).isFalse();
var skipStrategy = SkipProcessingStrategy.getInstance();
assertThat(skipStrategy.shouldProcess(timestamp, originator)).isFalse();
}
private static Stream<Arguments> edgeCaseProvider() {

22
ui-ngx/src/app/modules/home/components/rule-node/action/advanced-persistence-setting-row.component.ts

@ -24,11 +24,11 @@ import {
Validator
} from '@angular/forms';
import {
AdvancedPersistenceConfig,
defaultAdvancedPersistenceConfig,
AdvancedProcessingConfig,
defaultAdvancedProcessingConfig,
maxDeduplicateTimeSecs,
PersistenceType,
PersistenceTypeTranslationMap
ProcessingType,
ProcessingTypeTranslationMap
} from '@home/components/rule-node/action/timeseries-config.models';
import { isDefinedAndNotNull } from '@core/utils';
import { takeUntilDestroyed } from '@angular/core/rxjs-interop';
@ -52,13 +52,13 @@ export class AdvancedPersistenceSettingRowComponent implements ControlValueAcces
title: string;
persistenceSettingRowForm = this.fb.group({
type: [defaultAdvancedPersistenceConfig.type],
type: [defaultAdvancedProcessingConfig.type],
deduplicationIntervalSecs: [{value: 60, disabled: true}]
});
PersistenceType = PersistenceType;
persistenceStrategies = [PersistenceType.ON_EVERY_MESSAGE, PersistenceType.DEDUPLICATE, PersistenceType.SKIP];
PersistenceTypeTranslationMap = PersistenceTypeTranslationMap;
PersistenceType = ProcessingType;
persistenceStrategies = [ProcessingType.ON_EVERY_MESSAGE, ProcessingType.DEDUPLICATE, ProcessingType.SKIP];
PersistenceTypeTranslationMap = ProcessingTypeTranslationMap;
maxDeduplicateTime = maxDeduplicateTimeSecs;
@ -96,16 +96,16 @@ export class AdvancedPersistenceSettingRowComponent implements ControlValueAcces
};
}
writeValue(value: AdvancedPersistenceConfig) {
writeValue(value: AdvancedProcessingConfig) {
if (isDefinedAndNotNull(value)) {
this.persistenceSettingRowForm.patchValue(value, {emitEvent: false});
} else {
this.persistenceSettingRowForm.patchValue(defaultAdvancedPersistenceConfig);
this.persistenceSettingRowForm.patchValue(defaultAdvancedProcessingConfig);
}
}
private updatedValidation() {
if (this.persistenceSettingRowForm.get('type').value === PersistenceType.DEDUPLICATE) {
if (this.persistenceSettingRowForm.get('type').value === ProcessingType.DEDUPLICATE) {
this.persistenceSettingRowForm.get('deduplicationIntervalSecs').enable({emitEvent: false});
} else {
this.persistenceSettingRowForm.get('deduplicationIntervalSecs').disable({emitEvent: false})

5
ui-ngx/src/app/modules/home/components/rule-node/action/advanced-persistence-setting.component.html

@ -16,6 +16,11 @@
-->
<section [formGroup]="persistenceForm" class="tb-form-panel no-border no-padding">
<tb-example-hint
[hintText]="'rule-node-config.save-time-series.advanced-settings-hint'"
[popupHelpLink]="'rulenode/save_timeseries_node_advanced'"
>
</tb-example-hint>
<tb-advanced-persistence-setting-row
formControlName="timeseries"
title="{{ 'rule-node-config.save-time-series.time-series' | translate }}"

5
ui-ngx/src/app/modules/home/components/rule-node/action/advanced-persistence-setting.component.ts

@ -24,7 +24,7 @@ import {
} from '@angular/forms';
import { Component, forwardRef } from '@angular/core';
import { takeUntilDestroyed } from '@angular/core/rxjs-interop';
import { AdvancedPersistenceStrategy } from '@home/components/rule-node/action/timeseries-config.models';
import { AdvancedProcessingStrategy } from '@home/components/rule-node/action/timeseries-config.models';
@Component({
selector: 'tb-advanced-persistence-settings',
@ -76,8 +76,7 @@ export class AdvancedPersistenceSettingComponent implements ControlValueAccessor
};
}
writeValue(value: AdvancedPersistenceStrategy) {
writeValue(value: AdvancedProcessingStrategy) {
this.persistenceForm.patchValue(value, {emitEvent: false});
}
}

13
ui-ngx/src/app/modules/home/components/rule-node/action/timeseries-config.component.html

@ -16,10 +16,10 @@
-->
<section [formGroup]="timeseriesConfigForm" class="tb-form-panel no-border no-padding">
<div class="tb-form-panel stroked no-padding-bottom no-gap" formGroupName="persistenceSettings">
<div class="tb-form-panel stroked no-padding-bottom no-gap" formGroupName="processingSettings">
<div class="mb-4 flex flex-row items-center justify-between">
<div class="tb-form-panel-title" tb-hint-tooltip-icon="{{ 'rule-node-config.save-time-series.persistence-settings-hint' | translate}}" translate>
rule-node-config.save-time-series.persistence-settings
<div class="tb-form-panel-title" tb-hint-tooltip-icon="{{ 'rule-node-config.save-time-series.processing-settings-hint' | translate}}" translate>
rule-node-config.save-time-series.processing-settings
</div>
<tb-toggle-select appearance="fill" selectMediaBreakpoint="xs"
formControlName="isAdvanced">
@ -27,7 +27,7 @@
<tb-toggle-option [value]=true>{{ 'rule-node-config.advanced-mode' | translate }}</tb-toggle-option>
</tb-toggle-select>
</div>
@if(!timeseriesConfigForm.get('persistenceSettings.isAdvanced').value) {
@if(!timeseriesConfigForm.get('processingSettings.isAdvanced').value) {
<mat-form-field>
<mat-label translate>rule-node-config.save-time-series.strategy</mat-label>
<mat-select formControlName="type">
@ -37,7 +37,7 @@
</mat-select>
</mat-form-field>
@if(timeseriesConfigForm.get('persistenceSettings.type').value === PersistenceType.DEDUPLICATE) {
@if(timeseriesConfigForm.get('processingSettings.type').value === PersistenceType.DEDUPLICATE) {
<tb-time-unit-input
required
labelText="{{ 'rule-node-config.save-time-series.deduplication-interval' | translate }}"
@ -49,8 +49,7 @@
formControlName="deduplicationIntervalSecs">
</tb-time-unit-input>
}
}
@else {
} @else {
<tb-advanced-persistence-settings
class="mb-4"
formControlName="advanced"

78
ui-ngx/src/app/modules/home/components/rule-node/action/timeseries-config.component.ts

@ -20,10 +20,10 @@ import { RuleNodeConfigurationComponent } from '@shared/models/rule-node.models'
import {
defaultAdvancedPersistenceStrategy,
maxDeduplicateTimeSecs,
PersistenceSettings,
PersistenceSettingsForm,
PersistenceType,
PersistenceTypeTranslationMap,
ProcessingSettings,
ProcessingSettingsForm,
ProcessingType,
ProcessingTypeTranslationMap,
TimeseriesNodeConfiguration,
TimeseriesNodeConfigurationForm
} from '@home/components/rule-node/action/timeseries-config.models';
@ -37,9 +37,9 @@ export class TimeseriesConfigComponent extends RuleNodeConfigurationComponent {
timeseriesConfigForm: FormGroup;
PersistenceType = PersistenceType;
persistenceStrategies = [PersistenceType.ON_EVERY_MESSAGE, PersistenceType.DEDUPLICATE, PersistenceType.WEBSOCKETS_ONLY];
PersistenceTypeTranslationMap = PersistenceTypeTranslationMap;
PersistenceType = ProcessingType;
persistenceStrategies = [ProcessingType.ON_EVERY_MESSAGE, ProcessingType.DEDUPLICATE, ProcessingType.WEBSOCKETS_ONLY];
PersistenceTypeTranslationMap = ProcessingTypeTranslationMap;
maxDeduplicateTime = maxDeduplicateTimeSecs
@ -52,22 +52,22 @@ export class TimeseriesConfigComponent extends RuleNodeConfigurationComponent {
}
protected validatorTriggers(): string[] {
return ['persistenceSettings.isAdvanced', 'persistenceSettings.type'];
return ['processingSettings.isAdvanced', 'processingSettings.type'];
}
protected prepareInputConfig(config: TimeseriesNodeConfiguration): TimeseriesNodeConfigurationForm {
let persistenceSettings: PersistenceSettingsForm;
if (config?.persistenceSettings) {
const isAdvanced = config?.persistenceSettings?.type === PersistenceType.ADVANCED;
persistenceSettings = {
type: isAdvanced ? PersistenceType.ON_EVERY_MESSAGE : config.persistenceSettings.type,
let processingSettings: ProcessingSettingsForm;
if (config?.processingSettings) {
const isAdvanced = config?.processingSettings?.type === ProcessingType.ADVANCED;
processingSettings = {
type: isAdvanced ? ProcessingType.ON_EVERY_MESSAGE : config.processingSettings.type,
isAdvanced: isAdvanced,
deduplicationIntervalSecs: config.persistenceSettings?.deduplicationIntervalSecs ?? 60,
advanced: isAdvanced ? config.persistenceSettings : defaultAdvancedPersistenceStrategy
deduplicationIntervalSecs: config.processingSettings?.deduplicationIntervalSecs ?? 60,
advanced: isAdvanced ? config.processingSettings : defaultAdvancedPersistenceStrategy
}
} else {
persistenceSettings = {
type: PersistenceType.ON_EVERY_MESSAGE,
processingSettings = {
type: ProcessingType.ON_EVERY_MESSAGE,
isAdvanced: false,
deduplicationIntervalSecs: 60,
advanced: defaultAdvancedPersistenceStrategy
@ -75,36 +75,36 @@ export class TimeseriesConfigComponent extends RuleNodeConfigurationComponent {
}
return {
...config,
persistenceSettings: persistenceSettings
processingSettings: processingSettings
}
}
protected prepareOutputConfig(config: TimeseriesNodeConfigurationForm): TimeseriesNodeConfiguration {
let persistenceSettings: PersistenceSettings;
if (config.persistenceSettings.isAdvanced) {
persistenceSettings = {
...config.persistenceSettings.advanced,
type: PersistenceType.ADVANCED
let processingSettings: ProcessingSettings;
if (config.processingSettings.isAdvanced) {
processingSettings = {
...config.processingSettings.advanced,
type: ProcessingType.ADVANCED
};
} else {
persistenceSettings = {
type: config.persistenceSettings.type,
deduplicationIntervalSecs: config.persistenceSettings?.deduplicationIntervalSecs
processingSettings = {
type: config.processingSettings.type,
deduplicationIntervalSecs: config.processingSettings?.deduplicationIntervalSecs
};
}
return {
...config,
persistenceSettings
processingSettings
};
}
protected onConfigurationSet(config: TimeseriesNodeConfigurationForm) {
this.timeseriesConfigForm = this.fb.group({
persistenceSettings: this.fb.group({
isAdvanced: [config?.persistenceSettings?.isAdvanced ?? false],
type: [config?.persistenceSettings?.type ?? PersistenceType.ON_EVERY_MESSAGE],
processingSettings: this.fb.group({
isAdvanced: [config?.processingSettings?.isAdvanced ?? false],
type: [config?.processingSettings?.type ?? ProcessingType.ON_EVERY_MESSAGE],
deduplicationIntervalSecs: [
{value: config?.persistenceSettings?.deduplicationIntervalSecs ?? 60, disabled: true},
{value: config?.processingSettings?.deduplicationIntervalSecs ?? 60, disabled: true},
[Validators.required, Validators.max(maxDeduplicateTimeSecs)]
],
advanced: [{value: null, disabled: true}]
@ -115,18 +115,18 @@ export class TimeseriesConfigComponent extends RuleNodeConfigurationComponent {
}
protected updateValidators(emitEvent: boolean, _trigger?: string) {
const persistenceForm = this.timeseriesConfigForm.get('persistenceSettings') as FormGroup;
const isAdvanced: boolean = persistenceForm.get('isAdvanced').value;
const type: PersistenceType = persistenceForm.get('type').value;
if (!isAdvanced && type === PersistenceType.DEDUPLICATE) {
persistenceForm.get('deduplicationIntervalSecs').enable({emitEvent});
const processingForm = this.timeseriesConfigForm.get('processingSettings') as FormGroup;
const isAdvanced: boolean = processingForm.get('isAdvanced').value;
const type: ProcessingType = processingForm.get('type').value;
if (!isAdvanced && type === ProcessingType.DEDUPLICATE) {
processingForm.get('deduplicationIntervalSecs').enable({emitEvent});
} else {
persistenceForm.get('deduplicationIntervalSecs').disable({emitEvent});
processingForm.get('deduplicationIntervalSecs').disable({emitEvent});
}
if (isAdvanced) {
persistenceForm.get('advanced').enable({emitEvent});
processingForm.get('advanced').enable({emitEvent});
} else {
persistenceForm.get('advanced').disable({emitEvent});
processingForm.get('advanced').disable({emitEvent});
}
}
}

54
ui-ngx/src/app/modules/home/components/rule-node/action/timeseries-config.models.ts

@ -19,24 +19,24 @@ import { DAY, SECOND } from '@shared/models/time/time.models';
export const maxDeduplicateTimeSecs = DAY / SECOND;
export interface TimeseriesNodeConfiguration {
persistenceSettings: PersistenceSettings;
processingSettings: ProcessingSettings;
defaultTTL: number;
useServerTs: boolean;
}
export interface TimeseriesNodeConfigurationForm extends Omit<TimeseriesNodeConfiguration, 'persistenceSettings'> {
persistenceSettings: PersistenceSettingsForm
export interface TimeseriesNodeConfigurationForm extends Omit<TimeseriesNodeConfiguration, 'processingSettings'> {
processingSettings: ProcessingSettingsForm
}
export type PersistenceSettings = BasicPersistenceSettings & Partial<DeduplicatePersistenceStrategy> & Partial<AdvancedPersistenceStrategy>;
export type ProcessingSettings = BasicProcessingSettings & Partial<DeduplicateProcessingStrategy> & Partial<AdvancedProcessingStrategy>;
export type PersistenceSettingsForm = Omit<PersistenceSettings, keyof AdvancedPersistenceStrategy> & {
export type ProcessingSettingsForm = Omit<ProcessingSettings, keyof AdvancedProcessingStrategy> & {
isAdvanced: boolean;
advanced?: Partial<AdvancedPersistenceStrategy>;
type: PersistenceType;
advanced?: Partial<AdvancedProcessingStrategy>;
type: ProcessingType;
};
export enum PersistenceType {
export enum ProcessingType {
ON_EVERY_MESSAGE = 'ON_EVERY_MESSAGE',
DEDUPLICATE = 'DEDUPLICATE',
WEBSOCKETS_ONLY = 'WEBSOCKETS_ONLY',
@ -44,35 +44,35 @@ export enum PersistenceType {
SKIP = 'SKIP'
}
export const PersistenceTypeTranslationMap = new Map<PersistenceType, string>([
[PersistenceType.ON_EVERY_MESSAGE, 'rule-node-config.save-time-series.strategy-type.every-message'],
[PersistenceType.DEDUPLICATE, 'rule-node-config.save-time-series.strategy-type.deduplicate'],
[PersistenceType.WEBSOCKETS_ONLY, 'rule-node-config.save-time-series.strategy-type.web-sockets-only'],
[PersistenceType.SKIP, 'rule-node-config.save-time-series.strategy-type.skip'],
export const ProcessingTypeTranslationMap = new Map<ProcessingType, string>([
[ProcessingType.ON_EVERY_MESSAGE, 'rule-node-config.save-time-series.strategy-type.every-message'],
[ProcessingType.DEDUPLICATE, 'rule-node-config.save-time-series.strategy-type.deduplicate'],
[ProcessingType.WEBSOCKETS_ONLY, 'rule-node-config.save-time-series.strategy-type.web-sockets-only'],
[ProcessingType.SKIP, 'rule-node-config.save-time-series.strategy-type.skip'],
])
export interface BasicPersistenceSettings {
type: PersistenceType;
export interface BasicProcessingSettings {
type: ProcessingType;
}
export interface DeduplicatePersistenceStrategy extends BasicPersistenceSettings{
export interface DeduplicateProcessingStrategy extends BasicProcessingSettings{
deduplicationIntervalSecs: number;
}
export interface AdvancedPersistenceStrategy extends BasicPersistenceSettings{
timeseries: AdvancedPersistenceConfig;
latest: AdvancedPersistenceConfig;
webSockets: AdvancedPersistenceConfig;
export interface AdvancedProcessingStrategy extends BasicProcessingSettings{
timeseries: AdvancedProcessingConfig;
latest: AdvancedProcessingConfig;
webSockets: AdvancedProcessingConfig;
}
export type AdvancedPersistenceConfig = WithOptional<DeduplicatePersistenceStrategy, 'deduplicationIntervalSecs'>;
export type AdvancedProcessingConfig = WithOptional<DeduplicateProcessingStrategy, 'deduplicationIntervalSecs'>;
export const defaultAdvancedPersistenceConfig: AdvancedPersistenceConfig = {
type: PersistenceType.ON_EVERY_MESSAGE
export const defaultAdvancedProcessingConfig: AdvancedProcessingConfig = {
type: ProcessingType.ON_EVERY_MESSAGE
}
export const defaultAdvancedPersistenceStrategy: Omit<AdvancedPersistenceStrategy, 'type'> = {
timeseries: defaultAdvancedPersistenceConfig,
latest: defaultAdvancedPersistenceConfig,
webSockets: defaultAdvancedPersistenceConfig,
export const defaultAdvancedPersistenceStrategy: Omit<AdvancedProcessingStrategy, 'type'> = {
timeseries: defaultAdvancedProcessingConfig,
latest: defaultAdvancedProcessingConfig,
webSockets: defaultAdvancedProcessingConfig,
}

6
ui-ngx/src/app/modules/home/components/rule-node/common/example-hint.component.html

@ -15,11 +15,11 @@
limitations under the License.
-->
<div [hidden]="!hintText" class="tb-form-hint tb-primary-fill space-between">
<div [hidden]="!hintText" class="tb-form-hint tb-primary-fill flex justify-between gap-5">
<div [innerHTML]=" hintText | translate | safe: 'html'"
[style.text-align]="textAlign"
class="hint-text"></div>
<div *ngIf="popupHelpLink" class="see-example" tb-help-popup="{{ popupHelpLink }}"
class="w-full"></div>
<div *ngIf="popupHelpLink" class="flex shrink-0" tb-help-popup="{{ popupHelpLink }}"
hintMode
tb-help-popup-placement="right"
trigger-style="letter-spacing:0.25px; font-size:12px"

31
ui-ngx/src/app/modules/home/components/rule-node/common/example-hint.component.scss

@ -1,31 +0,0 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
:host {
.space-between {
display: flex;
justify-content: space-between;
gap: 20px;
.see-example {
display: flex;
flex-shrink: 0;
}
}
.hint-text {
width: 100%;
}
}

2
ui-ngx/src/app/modules/home/components/rule-node/common/example-hint.component.ts

@ -19,7 +19,7 @@ import { Component, Input } from '@angular/core';
@Component({
selector: 'tb-example-hint',
templateUrl: './example-hint.component.html',
styleUrls: ['./example-hint.component.scss']
styleUrls: []
})
export class ExampleHintComponent {
@Input() hintText: string;

29
ui-ngx/src/assets/help/en_US/rulenode/save_timeseries_node_advanced.md

@ -0,0 +1,29 @@
#### Potential unexpected behavior with mixed processing strategies
When configuring the processing strategies, certain combinations can lead to unexpected behavior. Consider the following scenarios:
- **Disabling WebSocket (WS) updates**
If WS updates are disabled, any changes to the time series data won’t be pushed to dashboards (or other WS subscriptions).
This means that even if a database is updated, dashboards may not display the updated data until browser page is reloaded.
- **Different deduplication intervals across actions**
When you configure different deduplication intervals for actions, the same incoming message might be processed differently for each action.
For example, a message might be stored immediately in the Time series table (if set to *On every message*) while not being stored in the Latest values table because its deduplication interval hasn’t elapsed.
Also, if the WebSocket updates are configured with a different interval, dashboards might show updates that do not match what is stored in the database.
- **Skipping database storage**
Choosing to disable one or more persistence actions (for instance, skipping database storage for Time series or Latest values while keeping WS updates enabled) introduces the risk of having only partial data available:
- If a message is processed only for real-time notifications (WebSockets) and not stored in the database, historical queries may not match data on the dashboard.
- When processing strategies for Time series and Latest values are out-of-sync, telemetry data may be stored in one table (e.g., Time series) while the same data is absent in the other (e.g., Latest values).
- **Deduplication cache clearing**
The deduplication mechanism uses a cache to track processed messages within each interval.
For performance and system stability reasons, this cache is periodically cleared.
As a result, if a cache entry is removed during the deduplication period, messages from the same originator may be processed more than once within that interval.
This means deduplication should be used as a performance optimization rather than an absolute guarantee of single processing per interval.
We recommend using deduplication only when the occasional repeated processing is acceptable and won't cause system correctness issue or data inconsistencies.

7
ui-ngx/src/assets/locale/locale.constant-en_US.json

@ -5134,8 +5134,9 @@
"basic-mode": "Basic",
"advanced-mode": "Advanced",
"save-time-series": {
"persistence-settings": "Persistence settings",
"persistence-settings-hint": "Define how and when time series data is saved. In Basic mode, apply a single persistence strategy to all actions or enable only WebSockets updates. Advanced mode allows you to configure individual persistence strategies for each action.",
"processing-settings": "Processing settings",
"processing-settings-hint": "Define how incoming messages are processed. In Basic mode, select a preconfigured processing strategy or enable only WebSocket updates. Advanced mode allows you to select individual processing strategies for each action.",
"advanced-settings-hint": "Be cautious when configuring processing strategies. Certain combinations can lead to unexpected behavior.",
"strategy": "Strategy",
"deduplication-interval": "Deduplication interval",
"deduplication-interval-required": "Deduplication interval is required",
@ -5147,7 +5148,7 @@
"web-sockets-only": "WebSockets only"
},
"time-series": "Time series",
"latest": "Latest",
"latest": "Latest values",
"web-sockets": "WebSockets"
},
"key-val": {

Loading…
Cancel
Save