Browse Source

Save attributes strategies: BE initial implementation

pull/12764/head
Dmytro Skarzhynets 2 years ago
parent
commit
400e74b00d
No known key found for this signature in database GPG Key ID: 2B51652F224037DF
  1. 5
      application/src/main/data/json/edge/rule_chains/edge_root_rule_chain.json
  2. 5
      application/src/main/data/json/tenant/device_profile/rule_chain_template.json
  3. 5
      application/src/main/data/json/tenant/rule_chains/root_rule_chain.json
  4. 30
      application/src/main/data/upgrade/basic/schema_update.sql
  5. 12
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  6. 96
      application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java
  7. 5
      monitoring/src/main/resources/root_rule_chain.json
  8. 17
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/AttributesSaveRequest.java
  9. 46
      rule-engine/rule-engine-api/src/test/java/org/thingsboard/rule/engine/api/AttributesSaveRequestTest.java
  10. 109
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java
  11. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNodeConfiguration.java
  12. 12
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java
  13. 62
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNodeConfiguration.java
  14. 51
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/settings/AttributesProcessingSettings.java
  15. 47
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/settings/BaseProcessingSettings.java
  16. 52
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/settings/TimeseriesProcessingSettings.java
  17. 29
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNodeConfigurationTest.java
  18. 680
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNodeTest.java
  19. 20
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNodeTest.java

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

@ -50,8 +50,11 @@
},
"type": "org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode",
"name": "Save Client Attributes",
"configurationVersion": 2,
"configurationVersion": 3,
"configuration": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "CLIENT_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,

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

@ -35,8 +35,11 @@
},
"type": "org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode",
"name": "Save Client Attributes",
"configurationVersion": 2,
"configurationVersion": 3,
"configuration": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "CLIENT_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,

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

@ -34,8 +34,11 @@
},
"type": "org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode",
"name": "Save Client Attributes",
"configurationVersion": 2,
"configurationVersion": 3,
"configuration": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "CLIENT_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,

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

@ -62,4 +62,32 @@ $$;
-- UPDATE SAVE TIME SERIES NODES END
ALTER TABLE api_usage_state ADD COLUMN IF NOT EXISTS version BIGINT DEFAULT 1;
-- UPDATE SAVE ATTRIBUTES NODES START
DO $$
BEGIN
-- Check if the rule_node table exists
IF EXISTS (
SELECT 1
FROM information_schema.tables
WHERE table_name = 'rule_node'
) THEN
UPDATE rule_node
SET configuration = (
configuration::jsonb
|| jsonb_build_object(
'processingSettings', jsonb_build_object('type', 'ON_EVERY_MESSAGE')
)
)::text,
configuration_version = 3
WHERE type = 'org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode'
AND configuration_version = 2;
END IF;
END;
$$;
-- UPDATE SAVE ATTRIBUTES NODES END
ALTER TABLE api_usage_state ADD COLUMN IF NOT EXISTS version BIGINT DEFAULT 1;

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

@ -53,6 +53,7 @@ import org.thingsboard.server.service.entitiy.entityview.TbEntityViewService;
import org.thingsboard.server.service.subscription.TbSubscriptionUtils;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Comparator;
import java.util.HashMap;
import java.util.List;
@ -165,9 +166,16 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
@Override
public void saveAttributesInternal(AttributesSaveRequest request) {
log.trace("Executing saveInternal [{}]", request);
ListenableFuture<List<Long>> saveFuture = attrService.save(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries());
ListenableFuture<List<Long>> saveFuture;
if (request.getStrategy().saveAttributes()) {
saveFuture = attrService.save(request.getTenantId(), request.getEntityId(), request.getScope(), request.getEntries());
} else {
saveFuture = Futures.immediateFuture(Collections.emptyList());
}
addMainCallback(saveFuture, request.getCallback());
addWsCallback(saveFuture, success -> onAttributesUpdate(request.getTenantId(), request.getEntityId(), request.getScope().name(), request.getEntries(), request.isNotifyDevice()));
if (request.getStrategy().sendWsUpdate()) {
addWsCallback(saveFuture, success -> onAttributesUpdate(request.getTenantId(), request.getEntityId(), request.getScope().name(), request.getEntries(), request.isNotifyDevice()));
}
}
@Override

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

@ -29,11 +29,13 @@ import org.junit.jupiter.params.provider.MethodSource;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.rule.engine.api.AttributesSaveRequest;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.ApiUsageState;
import org.thingsboard.server.common.data.ApiUsageStateValue;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
@ -86,7 +88,7 @@ class DefaultTelemetrySubscriptionServiceTest {
final long sampleTtl = 10_000L;
final List<TsKvEntry> sampleTelemetry = List.of(
final List<TsKvEntry> sampleTimeseries = List.of(
new BasicTsKvEntry(100L, new DoubleDataEntry("temperature", 65.2)),
new BasicTsKvEntry(100L, new DoubleDataEntry("humidity", 33.1))
);
@ -147,9 +149,9 @@ class DefaultTelemetrySubscriptionServiceTest {
lenient().when(partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId)).thenReturn(tpi);
lenient().when(tsService.save(tenantId, entityId, sampleTelemetry, sampleTtl)).thenReturn(immediateFuture(sampleTelemetry.size()));
lenient().when(tsService.saveWithoutLatest(tenantId, entityId, sampleTelemetry, sampleTtl)).thenReturn(immediateFuture(sampleTelemetry.size()));
lenient().when(tsService.saveLatest(tenantId, entityId, sampleTelemetry)).thenReturn(immediateFuture(listOfNNumbers(sampleTelemetry.size())));
lenient().when(tsService.save(tenantId, entityId, sampleTimeseries, sampleTtl)).thenReturn(immediateFuture(sampleTimeseries.size()));
lenient().when(tsService.saveWithoutLatest(tenantId, entityId, sampleTimeseries, sampleTtl)).thenReturn(immediateFuture(sampleTimeseries.size()));
lenient().when(tsService.saveLatest(tenantId, entityId, sampleTimeseries)).thenReturn(immediateFuture(listOfNNumbers(sampleTimeseries.size())));
// mock no entity views
lenient().when(tbEntityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId)).thenReturn(immediateFuture(Collections.emptyList()));
@ -171,7 +173,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.tenantId(tenantId)
.customerId(customerId)
.entityId(entityId)
.entries(sampleTelemetry)
.entries(sampleTimeseries)
.ttl(sampleTtl)
.strategy(new TimeseriesSaveRequest.Strategy(true, false, false))
.callback(emptyCallback)
@ -181,7 +183,7 @@ class DefaultTelemetrySubscriptionServiceTest {
telemetryService.saveTimeseries(request);
// THEN
then(apiUsageClient).should().report(tenantId, customerId, ApiUsageRecordKey.STORAGE_DP_COUNT, sampleTelemetry.size());
then(apiUsageClient).should().report(tenantId, customerId, ApiUsageRecordKey.STORAGE_DP_COUNT, sampleTimeseries.size());
}
@Test
@ -191,7 +193,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.tenantId(tenantId)
.customerId(customerId)
.entityId(entityId)
.entries(sampleTelemetry)
.entries(sampleTimeseries)
.ttl(sampleTtl)
.strategy(TimeseriesSaveRequest.Strategy.LATEST_AND_WS)
.callback(emptyCallback)
@ -214,7 +216,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.tenantId(tenantId)
.customerId(customerId)
.entityId(entityId)
.entries(sampleTelemetry)
.entries(sampleTimeseries)
.ttl(sampleTtl)
.strategy(TimeseriesSaveRequest.Strategy.SAVE_ALL)
.future(future)
@ -240,7 +242,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.tenantId(tenantId)
.customerId(customerId)
.entityId(entityId)
.entries(sampleTelemetry)
.entries(sampleTimeseries)
.ttl(sampleTtl)
.strategy(TimeseriesSaveRequest.Strategy.LATEST_AND_WS)
.future(future)
@ -260,12 +262,12 @@ class DefaultTelemetrySubscriptionServiceTest {
entityView.setTenantId(tenantId);
entityView.setCustomerId(customerId);
entityView.setEntityId(entityId);
entityView.setKeys(new TelemetryEntityView(sampleTelemetry.stream().map(KvEntry::getKey).toList(), new AttributesEntityView()));
entityView.setKeys(new TelemetryEntityView(sampleTimeseries.stream().map(KvEntry::getKey).toList(), new AttributesEntityView()));
// mock that there is one entity view
given(tbEntityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId)).willReturn(immediateFuture(List.of(entityView)));
// mock that save latest call for entity view is successful
given(tsService.saveLatest(tenantId, entityView.getId(), sampleTelemetry)).willReturn(immediateFuture(listOfNNumbers(sampleTelemetry.size())));
given(tsService.saveLatest(tenantId, entityView.getId(), sampleTimeseries)).willReturn(immediateFuture(listOfNNumbers(sampleTimeseries.size())));
// mock TPI for entity view
given(partitionService.resolve(ServiceType.TB_CORE, tenantId, entityView.getId())).willReturn(tpi);
@ -273,7 +275,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.tenantId(tenantId)
.customerId(customerId)
.entityId(entityId)
.entries(sampleTelemetry)
.entries(sampleTimeseries)
.ttl(sampleTtl)
.strategy(new TimeseriesSaveRequest.Strategy(false, true, false))
.callback(emptyCallback)
@ -284,12 +286,12 @@ class DefaultTelemetrySubscriptionServiceTest {
// THEN
// should save latest to both the main entity and it's entity view
then(tsService).should().saveLatest(tenantId, entityId, sampleTelemetry);
then(tsService).should().saveLatest(tenantId, entityView.getId(), sampleTelemetry);
then(tsService).should().saveLatest(tenantId, entityId, sampleTimeseries);
then(tsService).should().saveLatest(tenantId, entityView.getId(), sampleTimeseries);
then(tsService).shouldHaveNoMoreInteractions();
// should send WS update only for entity view (WS update for the main entity is disabled in the save request)
then(subscriptionManagerService).should().onTimeSeriesUpdate(tenantId, entityView.getId(), sampleTelemetry, TbCallback.EMPTY);
then(subscriptionManagerService).should().onTimeSeriesUpdate(tenantId, entityView.getId(), sampleTimeseries, TbCallback.EMPTY);
then(subscriptionManagerService).shouldHaveNoMoreInteractions();
}
@ -300,7 +302,7 @@ class DefaultTelemetrySubscriptionServiceTest {
.tenantId(tenantId)
.customerId(customerId)
.entityId(entityId)
.entries(sampleTelemetry)
.entries(sampleTimeseries)
.ttl(sampleTtl)
.strategy(new TimeseriesSaveRequest.Strategy(true, false, false))
.callback(emptyCallback)
@ -311,7 +313,7 @@ class DefaultTelemetrySubscriptionServiceTest {
// THEN
// should save only time series for the main entity
then(tsService).should().saveWithoutLatest(tenantId, entityId, sampleTelemetry, sampleTtl);
then(tsService).should().saveWithoutLatest(tenantId, entityId, sampleTimeseries, sampleTtl);
then(tsService).shouldHaveNoMoreInteractions();
// should not send any WS updates
@ -319,14 +321,14 @@ class DefaultTelemetrySubscriptionServiceTest {
}
@ParameterizedTest
@MethodSource("booleanCombinations")
void shouldCallCorrectApiBasedOnBooleanFlagsInTheSaveRequest(boolean saveTimeseries, boolean saveLatest, boolean sendWsUpdate) {
@MethodSource("allCombinationsOfThreeBooleans")
void shouldCallCorrectSaveTimeseriesApiBasedOnBooleanFlagsInTheSaveRequest(boolean saveTimeseries, boolean saveLatest, boolean sendWsUpdate) {
// GIVEN
var request = TimeseriesSaveRequest.builder()
.tenantId(tenantId)
.customerId(customerId)
.entityId(entityId)
.entries(sampleTelemetry)
.entries(sampleTimeseries)
.ttl(sampleTtl)
.strategy(new TimeseriesSaveRequest.Strategy(saveTimeseries, saveLatest, sendWsUpdate))
.callback(emptyCallback)
@ -337,22 +339,22 @@ class DefaultTelemetrySubscriptionServiceTest {
// THEN
if (saveTimeseries && saveLatest) {
then(tsService).should().save(tenantId, entityId, sampleTelemetry, sampleTtl);
then(tsService).should().save(tenantId, entityId, sampleTimeseries, sampleTtl);
} else if (saveLatest) {
then(tsService).should().saveLatest(tenantId, entityId, sampleTelemetry);
then(tsService).should().saveLatest(tenantId, entityId, sampleTimeseries);
} else if (saveTimeseries) {
then(tsService).should().saveWithoutLatest(tenantId, entityId, sampleTelemetry, sampleTtl);
then(tsService).should().saveWithoutLatest(tenantId, entityId, sampleTimeseries, sampleTtl);
}
then(tsService).shouldHaveNoMoreInteractions();
if (sendWsUpdate) {
then(subscriptionManagerService).should().onTimeSeriesUpdate(tenantId, entityId, sampleTelemetry, TbCallback.EMPTY);
then(subscriptionManagerService).should().onTimeSeriesUpdate(tenantId, entityId, sampleTimeseries, TbCallback.EMPTY);
} else {
then(subscriptionManagerService).shouldHaveNoInteractions();
}
}
private static Stream<Arguments> booleanCombinations() {
private static Stream<Arguments> allCombinationsOfThreeBooleans() {
return Stream.of(
Arguments.of(true, true, true),
Arguments.of(true, true, false),
@ -365,7 +367,49 @@ class DefaultTelemetrySubscriptionServiceTest {
);
}
// used to emulate sequence numbers returned by save latest API
@ParameterizedTest
@MethodSource("allCombinationsOfTwoBooleans")
void shouldCallCorrectSaveAttributesApiBasedOnBooleanFlagsInTheSaveRequest(boolean saveAttributes, boolean sendWsUpdate) {
// GIVEN
var request = AttributesSaveRequest.builder()
.tenantId(tenantId)
.entityId(entityId)
.scope(AttributeScope.SERVER_SCOPE)
.entry(new DoubleDataEntry("temperature", 65.2))
.notifyDevice(false)
.strategy(new AttributesSaveRequest.Strategy(saveAttributes, sendWsUpdate))
.callback(emptyCallback)
.build();
lenient().when(attrService.save(tenantId, entityId, request.getScope(), request.getEntries())).thenReturn(immediateFuture(listOfNNumbers(request.getEntries().size())));
// WHEN
telemetryService.saveAttributes(request);
// THEN
if (saveAttributes) {
then(attrService).should().save(tenantId, entityId, request.getScope(), request.getEntries());
} else {
then(attrService).shouldHaveNoInteractions();
}
if (sendWsUpdate) {
then(subscriptionManagerService).should().onAttributesUpdate(tenantId, entityId, request.getScope().name(), request.getEntries(), request.isNotifyDevice(), TbCallback.EMPTY);
} else {
then(subscriptionManagerService).shouldHaveNoInteractions();
}
}
static Stream<Arguments> allCombinationsOfTwoBooleans() {
return Stream.of(
Arguments.of(true, true),
Arguments.of(true, false),
Arguments.of(false, true),
Arguments.of(false, false)
);
}
// used to emulate sequence numbers returned by save APIs
private static List<Long> listOfNNumbers(int N) {
return LongStream.range(0, N).boxed().toList();
}

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

@ -39,8 +39,11 @@
"type": "org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode",
"name": "Save Attributes",
"singletonMode": false,
"configurationVersion": 1,
"configurationVersion": 3,
"configuration": {
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "CLIENT_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,

17
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/AttributesSaveRequest.java

@ -40,8 +40,17 @@ public class AttributesSaveRequest {
private final AttributeScope scope;
private final List<AttributeKvEntry> entries;
private final boolean notifyDevice;
private final Strategy strategy;
private final FutureCallback<Void> callback;
public record Strategy(boolean saveAttributes, boolean sendWsUpdate) {
public static final Strategy PROCESS_ALL = new Strategy(true, true);
public static final Strategy WS_ONLY = new Strategy(false, true);
public static final Strategy SKIP_ALL = new Strategy(false, false);
}
public static Builder builder() {
return new Builder();
}
@ -53,6 +62,7 @@ public class AttributesSaveRequest {
private AttributeScope scope;
private List<AttributeKvEntry> entries;
private boolean notifyDevice = true;
private Strategy strategy = Strategy.PROCESS_ALL;
private FutureCallback<Void> callback;
Builder() {}
@ -100,6 +110,11 @@ public class AttributesSaveRequest {
return this;
}
public Builder strategy(Strategy strategy) {
this.strategy = strategy;
return this;
}
public Builder callback(FutureCallback<Void> callback) {
this.callback = callback;
return this;
@ -120,7 +135,7 @@ public class AttributesSaveRequest {
}
public AttributesSaveRequest build() {
return new AttributesSaveRequest(tenantId, entityId, scope, entries, notifyDevice, callback);
return new AttributesSaveRequest(tenantId, entityId, scope, entries, notifyDevice, strategy, callback);
}
}

46
rule-engine/rule-engine-api/src/test/java/org/thingsboard/rule/engine/api/AttributesSaveRequestTest.java

@ -0,0 +1,46 @@
/**
* Copyright © 2016-2025 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.
*/
package org.thingsboard.rule.engine.api;
import org.junit.jupiter.api.Test;
import static org.assertj.core.api.Assertions.assertThat;
class AttributesSaveRequestTest {
@Test
void testDefaultSaveStrategyIsProcessAll() {
var request = AttributesSaveRequest.builder().build();
assertThat(request.getStrategy()).isEqualTo(AttributesSaveRequest.Strategy.PROCESS_ALL);
}
@Test
void testProcessAllStrategy() {
assertThat(AttributesSaveRequest.Strategy.PROCESS_ALL).isEqualTo(new AttributesSaveRequest.Strategy(true, true));
}
@Test
void testWsOnlyStrategy() {
assertThat(AttributesSaveRequest.Strategy.WS_ONLY).isEqualTo(new AttributesSaveRequest.Strategy(false, true));
}
@Test
void testSkipAllStrategy() {
assertThat(AttributesSaveRequest.Strategy.SKIP_ALL).isEqualTo(new AttributesSaveRequest.Strategy(false, false));
}
}

109
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java

@ -23,6 +23,7 @@ import com.google.common.util.concurrent.MoreExecutors;
import com.google.gson.JsonParser;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.DonAsynchron;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.AttributesSaveRequest;
import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext;
@ -30,6 +31,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.util.TbNodeUtils;
import org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings;
import org.thingsboard.server.common.adaptor.JsonConverter;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.StringUtils;
@ -43,9 +45,14 @@ import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.UUID;
import java.util.function.Function;
import java.util.stream.Collectors;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.Advanced;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.Deduplicate;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.OnEveryMessage;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.WebSocketsOnly;
import static org.thingsboard.server.common.data.DataConstants.NOTIFY_DEVICE_METADATA_KEY;
import static org.thingsboard.server.common.data.DataConstants.SCOPE;
import static org.thingsboard.server.common.data.msg.TbMsgType.POST_ATTRIBUTES_REQUEST;
@ -55,13 +62,47 @@ import static org.thingsboard.server.common.data.msg.TbMsgType.POST_ATTRIBUTES_R
type = ComponentType.ACTION,
name = "save attributes",
configClazz = TbMsgAttributesNodeConfiguration.class,
version = 2,
nodeDescription = "Saves attributes data",
nodeDetails = "Saves entity attributes based on configurable scope parameter. Expects messages with 'POST_ATTRIBUTES_REQUEST' message type. " +
"If upsert(update/insert) operation is completed successfully rule node will send the incoming message via <b>Success</b> chain, otherwise, <b>Failure</b> chain is used. " +
"Additionally if checkbox <b>Send attributes updated notification</b> is set to true, rule node will put the \"Attributes Updated\" " +
"event for <b>SHARED_SCOPE</b> and <b>SERVER_SCOPE</b> attributes updates to the corresponding rule engine queue." +
"Performance checkbox 'Save attributes only if the value changes' will skip attributes overwrites for values with no changes (avoid concurrent writes because this check is not transactional; will not update 'Last updated time' for skipped attributes).",
version = 3,
nodeDescription = """
Saves attribute data with a configurable scope and according to configured processing strategies.
""",
nodeDetails = """
Node performs two <strong>actions:</strong>
<ul>
<li><strong>Attributes:</strong> save attribute data to a database.</li>
<li><strong>WebSockets:</strong> notify WebSockets subscriptions about attribute data updates.</li>
</ul>
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>Processing strategies</strong> are configured using <em>processing settings</em>, which support two modes:
<ul>
<li><strong>Basic</strong>
<ul>
<li><strong>On every message:</strong> applies the "On every message" strategy to all actions.</li>
<li><strong>Deduplicate:</strong> applies the "Deduplicate" strategy (with a specified interval) to all actions.</li>
<li><strong>WebSockets only:</strong> applies the "Skip" strategy to Attributes, and the "On every message" strategy to WebSockets.</li>
</ul>
</li>
<li><strong>Advanced:</strong> configure each action’s strategy independently.</li>
</ul>
Additionally:
<ul>
<li>If <b>Save attributes only if the value changes</b> is enabled, the rule node compares the received attribute value with the current stored value and skips the save operation if they match.</li>
<li>If <b>Send attributes updated notification</b> is enabled, the rule node will put the <b>Attributes Updated</b> event for <code>SHARED_SCOPE</code> and <code>SERVER_SCOPE</code> attribute updates to the queue named <code>Main</code>.</li>
<li>If <b>Force notification to the device</b> is enabled, then rule node will always notify device about <code>SHARED_SCOPE</code> attribute updates, regardless of the value of <code>notifyDevice</code> metadata property.</li>
</ul>
This node expects messages of type <code>POST_ATTRIBUTES_REQUEST</code>.
<br><br>
Output connections: <code>Success</code>, <code>Failure</code>.
""",
configDirective = "tbActionNodeAttributesConfig",
icon = "file_upload"
)
@ -73,9 +114,12 @@ public class TbMsgAttributesNode implements TbNode {
private TbMsgAttributesNodeConfiguration config;
private AttributesProcessingSettings processingSettings;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbMsgAttributesNodeConfiguration.class);
config = TbNodeUtils.convert(configuration, TbMsgAttributesNodeConfiguration.class);
processingSettings = config.getProcessingSettings();
}
@Override
@ -90,11 +134,20 @@ public class TbMsgAttributesNode implements TbNode {
ctx.tellSuccess(msg);
return;
}
AttributesSaveRequest.Strategy strategy = determineSaveStrategy(msg.getMetaDataTs(), msg.getOriginator().getId());
// short-circuit
if (!strategy.saveAttributes() && !strategy.sendWsUpdate()) {
ctx.tellSuccess(msg);
return;
}
AttributeScope scope = getScope(msg.getMetaData().getValue(SCOPE));
boolean sendAttributesUpdateNotification = checkSendNotification(scope);
if (!config.isUpdateAttributesOnlyOnValueChange()) {
saveAttr(newAttributes, ctx, msg, scope, sendAttributesUpdateNotification);
saveAttr(newAttributes, ctx, msg, scope, sendAttributesUpdateNotification, strategy);
return;
}
@ -104,13 +157,41 @@ public class TbMsgAttributesNode implements TbNode {
DonAsynchron.withCallback(findFuture,
currentAttributes -> {
List<AttributeKvEntry> attributesChanged = filterChangedAttr(currentAttributes, newAttributes);
saveAttr(attributesChanged, ctx, msg, scope, sendAttributesUpdateNotification);
saveAttr(attributesChanged, ctx, msg, scope, sendAttributesUpdateNotification, strategy);
},
throwable -> ctx.tellFailure(msg, throwable),
MoreExecutors.directExecutor());
}
void saveAttr(List<AttributeKvEntry> attributes, TbContext ctx, TbMsg msg, AttributeScope scope, boolean sendAttributesUpdateNotification) {
private AttributesSaveRequest.Strategy determineSaveStrategy(long ts, UUID originatorUuid) {
if (processingSettings instanceof OnEveryMessage) {
return AttributesSaveRequest.Strategy.PROCESS_ALL;
}
if (processingSettings instanceof WebSocketsOnly) {
return AttributesSaveRequest.Strategy.WS_ONLY;
}
if (processingSettings instanceof Deduplicate deduplicate) {
boolean isFirstMsgInInterval = deduplicate.getProcessingStrategy().shouldProcess(ts, originatorUuid);
return isFirstMsgInInterval ? AttributesSaveRequest.Strategy.PROCESS_ALL : AttributesSaveRequest.Strategy.SKIP_ALL;
}
if (processingSettings instanceof Advanced advanced) {
return new AttributesSaveRequest.Strategy(
advanced.attributes().shouldProcess(ts, originatorUuid),
advanced.webSockets().shouldProcess(ts, originatorUuid)
);
}
// should not happen
throw new IllegalArgumentException("Unknown processing settings type: " + processingSettings.getClass().getSimpleName());
}
private void saveAttr(
List<AttributeKvEntry> attributes,
TbContext ctx,
TbMsg msg,
AttributeScope scope,
boolean sendAttributesUpdateNotification,
AttributesSaveRequest.Strategy strategy
) {
if (attributes.isEmpty()) {
ctx.tellSuccess(msg);
return;
@ -124,11 +205,12 @@ public class TbMsgAttributesNode implements TbNode {
.scope(scope)
.entries(attributes)
.notifyDevice(config.isNotifyDevice() || checkNotifyDeviceMdValue(msg.getMetaData().getValue(NOTIFY_DEVICE_METADATA_KEY)))
.strategy(strategy)
.callback(callback)
.build());
}
List<AttributeKvEntry> filterChangedAttr(List<AttributeKvEntry> currentAttributes, List<AttributeKvEntry> newAttributes) {
private List<AttributeKvEntry> filterChangedAttr(List<AttributeKvEntry> currentAttributes, List<AttributeKvEntry> newAttributes) {
if (currentAttributes == null || currentAttributes.isEmpty()) {
return newAttributes;
}
@ -178,6 +260,9 @@ public class TbMsgAttributesNode implements TbNode {
hasChanges = fixEscapedBooleanConfigParameter(oldConfiguration, SEND_ATTRIBUTES_UPDATED_NOTIFICATION_KEY, hasChanges, false);
// update updateAttributesOnlyOnValueChange.
hasChanges = fixEscapedBooleanConfigParameter(oldConfiguration, UPDATE_ATTRIBUTES_ONLY_ON_VALUE_CHANGE_KEY, hasChanges, true);
case 2:
hasChanges = true;
((ObjectNode) oldConfiguration).set("processingSettings", JacksonUtil.valueToTree(new OnEveryMessage()));
break;
default:
break;

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNodeConfiguration.java

@ -15,13 +15,20 @@
*/
package org.thingsboard.rule.engine.telemetry;
import jakarta.validation.constraints.NotNull;
import lombok.Data;
import org.thingsboard.rule.engine.api.NodeConfiguration;
import org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings;
import org.thingsboard.server.common.data.DataConstants;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.OnEveryMessage;
@Data
public class TbMsgAttributesNodeConfiguration implements NodeConfiguration<TbMsgAttributesNodeConfiguration> {
@NotNull
private AttributesProcessingSettings processingSettings;
private String scope;
private boolean notifyDevice;
@ -31,6 +38,7 @@ public class TbMsgAttributesNodeConfiguration implements NodeConfiguration<TbMsg
@Override
public TbMsgAttributesNodeConfiguration defaultConfiguration() {
TbMsgAttributesNodeConfiguration configuration = new TbMsgAttributesNodeConfiguration();
configuration.setProcessingSettings(new OnEveryMessage());
configuration.setScope(DataConstants.SERVER_SCOPE);
configuration.setNotifyDevice(false);
configuration.setSendAttributesUpdatedNotification(false);
@ -38,4 +46,5 @@ public class TbMsgAttributesNodeConfiguration implements NodeConfiguration<TbMsg
configuration.setUpdateAttributesOnlyOnValueChange(true);
return configuration;
}
}

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

@ -27,6 +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.settings.TimeseriesProcessingSettings;
import org.thingsboard.rule.engine.telemetry.strategy.ProcessingStrategy;
import org.thingsboard.server.common.adaptor.JsonConverter;
import org.thingsboard.server.common.data.StringUtils;
@ -45,11 +46,10 @@ import java.util.Map;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
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.rule.engine.telemetry.settings.TimeseriesProcessingSettings.Advanced;
import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.Deduplicate;
import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.OnEveryMessage;
import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.WebSocketsOnly;
import static org.thingsboard.server.common.data.msg.TbMsgType.POST_TELEMETRY_REQUEST;
@Slf4j
@ -110,7 +110,7 @@ public class TbMsgTimeseriesNode implements TbNode {
private TbContext ctx;
private long tenantProfileDefaultStorageTtl;
private ProcessingSettings processingSettings;
private TimeseriesProcessingSettings processingSettings;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {

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

@ -15,23 +15,12 @@
*/
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;
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.ProcessingStrategy;
import org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings;
import java.util.Objects;
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.rule.engine.telemetry.settings.TimeseriesProcessingSettings.OnEveryMessage;
@Data
public class TbMsgTimeseriesNodeConfiguration implements NodeConfiguration<TbMsgTimeseriesNodeConfiguration> {
@ -39,7 +28,7 @@ public class TbMsgTimeseriesNodeConfiguration implements NodeConfiguration<TbMsg
private long defaultTTL;
private boolean useServerTs;
@NotNull
private TbMsgTimeseriesNodeConfiguration.ProcessingSettings processingSettings;
private TimeseriesProcessingSettings processingSettings;
@Override
public TbMsgTimeseriesNodeConfiguration defaultConfiguration() {
@ -50,49 +39,4 @@ public class TbMsgTimeseriesNodeConfiguration implements NodeConfiguration<TbMsg
return configuration;
}
@JsonTypeInfo(
use = JsonTypeInfo.Id.NAME,
include = JsonTypeInfo.As.PROPERTY,
property = "type"
)
@JsonSubTypes({
@JsonSubTypes.Type(value = OnEveryMessage.class, name = "ON_EVERY_MESSAGE"),
@JsonSubTypes.Type(value = WebSocketsOnly.class, name = "WEBSOCKETS_ONLY"),
@JsonSubTypes.Type(value = Deduplicate.class, name = "DEDUPLICATE"),
@JsonSubTypes.Type(value = Advanced.class, name = "ADVANCED")
})
sealed interface ProcessingSettings permits OnEveryMessage, Deduplicate, WebSocketsOnly, Advanced {
record OnEveryMessage() implements ProcessingSettings {}
record WebSocketsOnly() implements ProcessingSettings {}
@Getter
final class Deduplicate implements ProcessingSettings {
private final int deduplicationIntervalSecs;
@JsonIgnore
private final ProcessingStrategy processingStrategy;
@JsonCreator
Deduplicate(@JsonProperty("deduplicationIntervalSecs") int deduplicationIntervalSecs) {
this.deduplicationIntervalSecs = deduplicationIntervalSecs;
processingStrategy = ProcessingStrategy.deduplicate(deduplicationIntervalSecs);
}
}
record Advanced(ProcessingStrategy timeseries, ProcessingStrategy latest, ProcessingStrategy webSockets) implements ProcessingSettings {
public Advanced {
Objects.requireNonNull(timeseries);
Objects.requireNonNull(latest);
Objects.requireNonNull(webSockets);
}
}
}
}

51
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/settings/AttributesProcessingSettings.java

@ -0,0 +1,51 @@
/**
* Copyright © 2016-2025 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.
*/
package org.thingsboard.rule.engine.telemetry.settings;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import org.thingsboard.rule.engine.telemetry.strategy.ProcessingStrategy;
import java.util.Objects;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.Advanced;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.Deduplicate;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.OnEveryMessage;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.WebSocketsOnly;
@JsonTypeInfo(
use = JsonTypeInfo.Id.NAME,
include = JsonTypeInfo.As.PROPERTY,
property = "type"
)
@JsonSubTypes({
@JsonSubTypes.Type(value = OnEveryMessage.class, name = "ON_EVERY_MESSAGE"),
@JsonSubTypes.Type(value = WebSocketsOnly.class, name = "WEBSOCKETS_ONLY"),
@JsonSubTypes.Type(value = Deduplicate.class, name = "DEDUPLICATE"),
@JsonSubTypes.Type(value = Advanced.class, name = "ADVANCED")
})
public sealed interface AttributesProcessingSettings extends BaseProcessingSettings permits OnEveryMessage, Deduplicate, WebSocketsOnly, Advanced {
record Advanced(ProcessingStrategy attributes, ProcessingStrategy webSockets) implements AttributesProcessingSettings {
public Advanced {
Objects.requireNonNull(attributes);
Objects.requireNonNull(webSockets);
}
}
}

47
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/settings/BaseProcessingSettings.java

@ -0,0 +1,47 @@
/**
* Copyright © 2016-2025 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.
*/
package org.thingsboard.rule.engine.telemetry.settings;
import com.fasterxml.jackson.annotation.JsonCreator;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonProperty;
import lombok.Getter;
import org.thingsboard.rule.engine.telemetry.strategy.ProcessingStrategy;
sealed interface BaseProcessingSettings permits TimeseriesProcessingSettings, AttributesProcessingSettings {
record OnEveryMessage() implements TimeseriesProcessingSettings, AttributesProcessingSettings {}
record WebSocketsOnly() implements TimeseriesProcessingSettings, AttributesProcessingSettings {}
@Getter
final class Deduplicate implements TimeseriesProcessingSettings, AttributesProcessingSettings {
private final int deduplicationIntervalSecs;
@JsonIgnore
private final ProcessingStrategy processingStrategy;
@JsonCreator
public Deduplicate(@JsonProperty("deduplicationIntervalSecs") int deduplicationIntervalSecs) {
this.deduplicationIntervalSecs = deduplicationIntervalSecs;
this.processingStrategy = ProcessingStrategy.deduplicate(deduplicationIntervalSecs);
}
}
}

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

@ -0,0 +1,52 @@
/**
* Copyright © 2016-2025 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.
*/
package org.thingsboard.rule.engine.telemetry.settings;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import org.thingsboard.rule.engine.telemetry.strategy.ProcessingStrategy;
import java.util.Objects;
import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.Advanced;
import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.Deduplicate;
import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.OnEveryMessage;
import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.WebSocketsOnly;
@JsonTypeInfo(
use = JsonTypeInfo.Id.NAME,
include = JsonTypeInfo.As.PROPERTY,
property = "type"
)
@JsonSubTypes({
@JsonSubTypes.Type(value = OnEveryMessage.class, name = "ON_EVERY_MESSAGE"),
@JsonSubTypes.Type(value = WebSocketsOnly.class, name = "WEBSOCKETS_ONLY"),
@JsonSubTypes.Type(value = Deduplicate.class, name = "DEDUPLICATE"),
@JsonSubTypes.Type(value = Advanced.class, name = "ADVANCED")
})
public sealed interface TimeseriesProcessingSettings extends BaseProcessingSettings permits OnEveryMessage, Deduplicate, WebSocketsOnly, Advanced {
record Advanced(ProcessingStrategy timeseries, ProcessingStrategy latest, ProcessingStrategy webSockets) implements TimeseriesProcessingSettings {
public Advanced {
Objects.requireNonNull(timeseries);
Objects.requireNonNull(latest);
Objects.requireNonNull(webSockets);
}
}
}

29
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNodeConfigurationTest.java

@ -1,29 +0,0 @@
/**
* Copyright © 2016-2025 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.
*/
package org.thingsboard.rule.engine.telemetry;
import org.junit.jupiter.api.Test;
import static org.assertj.core.api.Assertions.assertThat;
class TbMsgAttributesNodeConfigurationTest {
@Test
void testDefaultConfig_givenUpdateAttributesOnlyOnValueChange_thenTrue_sinceVersion1() {
assertThat(new TbMsgAttributesNodeConfiguration().defaultConfiguration().isUpdateAttributesOnlyOnValueChange()).isTrue();
}
}

680
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNodeTest.java

@ -15,13 +15,15 @@
*/
package org.thingsboard.rule.engine.telemetry;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.extern.slf4j.Slf4j;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;
import org.mockito.Mock;
import org.mockito.Spy;
import org.mockito.junit.jupiter.MockitoExtension;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.AbstractRuleNodeUpgradeTest;
import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
@ -29,6 +31,7 @@ import org.thingsboard.rule.engine.api.TbContext;
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.telemetry.strategy.ProcessingStrategy;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
@ -42,185 +45,632 @@ import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.msg.TbMsgType;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.dao.service.ConstraintValidator;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.stream.Stream;
import static com.google.common.util.concurrent.Futures.immediateFuture;
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.anyBoolean;
import static org.mockito.ArgumentMatchers.anyList;
import static org.mockito.ArgumentMatchers.assertArg;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.willCallRealMethod;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.spy;
import static org.mockito.BDDMockito.given;
import static org.mockito.BDDMockito.then;
import static org.mockito.Mockito.clearInvocations;
import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
import static org.thingsboard.rule.engine.api.AttributesSaveRequest.Strategy;
import static org.thingsboard.rule.engine.api.AttributesSaveRequest.builder;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.Advanced;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.Deduplicate;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.OnEveryMessage;
import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.WebSocketsOnly;
import static org.thingsboard.server.common.data.DataConstants.NOTIFY_DEVICE_METADATA_KEY;
@Slf4j
@ExtendWith(MockitoExtension.class)
class TbMsgAttributesNodeTest extends AbstractRuleNodeUpgradeTest {
private TenantId tenantId;
private DeviceId deviceId;
private TbMsgAttributesNode node;
final TenantId tenantId = TenantId.fromUUID(UUID.fromString("6c18691e-4470-4766-9739-aface71d761f"));
final DeviceId deviceId = new DeviceId(UUID.fromString("b66159d7-c77e-45e8-bb41-a8f557f434c1"));
@Spy
TbMsgAttributesNode node;
TbMsgAttributesNodeConfiguration config;
@Mock
TbContext ctxMock;
@Mock
AttributesService attributesServiceMock;
@Mock
RuleEngineTelemetryService telemetryServiceMock;
@BeforeEach
void setUp() {
tenantId = new TenantId(UUID.fromString("6c18691e-4470-4766-9739-aface71d761f"));
deviceId = new DeviceId(UUID.fromString("b66159d7-c77e-45e8-bb41-a8f557f434c1"));
node = spy(TbMsgAttributesNode.class);
lenient().when(ctxMock.getTenantId()).thenReturn(tenantId);
lenient().when(ctxMock.getAttributesService()).thenReturn(attributesServiceMock);
lenient().when(ctxMock.getTelemetryService()).thenReturn(telemetryServiceMock);
config = new TbMsgAttributesNodeConfiguration().defaultConfiguration();
}
@Test
void testFilterChangedAttr_whenCurrentAttributesEmpty_thenReturnNewAttributes() {
List<AttributeKvEntry> newAttributes = new ArrayList<>();
void verifyDefaultConfig() {
assertThat(config.getProcessingSettings()).isInstanceOf(OnEveryMessage.class);
assertThat(config.getScope()).isEqualTo("SERVER_SCOPE");
assertThat(config.isNotifyDevice()).isFalse();
assertThat(config.isSendAttributesUpdatedNotification()).isFalse();
assertThat(config.isUpdateAttributesOnlyOnValueChange()).isTrue();
}
List<AttributeKvEntry> filtered = node.filterChangedAttr(Collections.emptyList(), newAttributes);
assertThat(filtered).isSameAs(newAttributes);
@Test
void givenProcessingSettingsAreNull_whenValidatingConstraints_thenThrowsException() {
// GIVEN
config.setProcessingSettings(null);
// WHEN-THEN
assertThatThrownBy(() -> ConstraintValidator.validateFields(config))
.isInstanceOf(DataValidationException.class)
.hasMessage("Validation error: processingSettings must not be null");
}
@Test
void testFilterChangedAttr_whenCurrentAttributesContainsInAnyOrderNewAttributes_thenReturnEmptyList() {
List<AttributeKvEntry> currentAttributes = List.of(
new BaseAttributeKvEntry(1694000000L, new StringDataEntry("address", "Peremohy ave 1")),
new BaseAttributeKvEntry(1694000000L, new BooleanDataEntry("valid", true)),
new BaseAttributeKvEntry(1694000000L, new LongDataEntry("counter", 100L)),
new BaseAttributeKvEntry(1694000000L, new DoubleDataEntry("temp", -18.35)),
new BaseAttributeKvEntry(1694000000L, new JsonDataEntry("json", "{\"warning\":\"out of paper\"}"))
);
List<AttributeKvEntry> newAttributes = new ArrayList<>(currentAttributes);
newAttributes.add(newAttributes.get(0));
newAttributes.remove(0);
assertThat(newAttributes).hasSize(currentAttributes.size());
assertThat(currentAttributes).isNotEmpty();
assertThat(newAttributes).containsExactlyInAnyOrderElementsOf(currentAttributes);
List<AttributeKvEntry> filtered = node.filterChangedAttr(currentAttributes, newAttributes);
assertThat(filtered).isEmpty(); //no changes
void givenOnEveryMessageProcessingSettingsAndSameMessageTwoTimes_whenOnMsg_thenPersistSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setUpdateAttributesOnlyOnValueChange(false);
config.setProcessingSettings(new OnEveryMessage());
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_ATTRIBUTES_REQUEST)
.originator(deviceId)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of(NOTIFY_DEVICE_METADATA_KEY, "false")))
.build();
// WHEN-THEN
var expectedSaveRequest = builder()
.tenantId(tenantId)
.entityId(msg.getOriginator())
.scope(AttributeScope.valueOf(config.getScope()))
.entry(new DoubleDataEntry("temperature", 22.3))
.notifyDevice(false)
.strategy(Strategy.PROCESS_ALL)
.build();
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(1)).saveAttributes(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest)
.usingRecursiveComparison()
.ignoringFields("callback", "entries.lastUpdateTs")
.isEqualTo(expectedSaveRequest)
));
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(2)).saveAttributes(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest)
.usingRecursiveComparison()
.ignoringFields("callback", "entries.lastUpdateTs")
.isEqualTo(expectedSaveRequest)
));
}
@Test
void givenDeduplicateProcessingSettingsAndSameMessageTwoTimes_whenOnMsg_thenPersistThisMessageOnlyFirstTime() throws TbNodeException {
// GIVEN
config.setUpdateAttributesOnlyOnValueChange(false);
config.setProcessingSettings(new Deduplicate(10));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_ATTRIBUTES_REQUEST)
.originator(deviceId)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of(NOTIFY_DEVICE_METADATA_KEY, "false")))
.build();
// WHEN-THEN
var expectedSaveRequest = builder()
.tenantId(tenantId)
.entityId(msg.getOriginator())
.scope(AttributeScope.valueOf(config.getScope()))
.entry(new DoubleDataEntry("temperature", 22.3))
.notifyDevice(false)
.strategy(Strategy.PROCESS_ALL)
.build();
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should().saveAttributes(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest)
.usingRecursiveComparison()
.ignoringFields("callback", "entries.lastUpdateTs")
.isEqualTo(expectedSaveRequest)
));
clearInvocations(telemetryServiceMock, ctxMock);
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(never()).saveAttributes(any());
}
@Test
void testFilterChangedAttr_whenCurrentAttributesContainsInAnyOrderNewAttributes_thenReturnExpectedList() {
void givenWebSocketsOnlyProcessingSettingsAndSameMessageTwoTimes_whenOnMsg_thenSendsOnlyWsUpdateTwoTimes() throws TbNodeException {
// GIVEN
config.setUpdateAttributesOnlyOnValueChange(false);
config.setProcessingSettings(new WebSocketsOnly());
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_ATTRIBUTES_REQUEST)
.originator(deviceId)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of(NOTIFY_DEVICE_METADATA_KEY, "false")))
.build();
// WHEN-THEN
var expectedSaveRequest = builder()
.tenantId(tenantId)
.entityId(msg.getOriginator())
.scope(AttributeScope.valueOf(config.getScope()))
.entry(new DoubleDataEntry("temperature", 22.3))
.notifyDevice(false)
.strategy(Strategy.WS_ONLY)
.build();
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(1)).saveAttributes(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest)
.usingRecursiveComparison()
.ignoringFields("callback", "entries.lastUpdateTs")
.isEqualTo(expectedSaveRequest)
));
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(2)).saveAttributes(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest)
.usingRecursiveComparison()
.ignoringFields("callback", "entries.lastUpdateTs")
.isEqualTo(expectedSaveRequest)
));
}
@Test
void givenAdvancedProcessingSettingsWithOnEveryMessageStrategiesForAllActionsAndSameMessageTwoTimes_whenOnMsg_thenPersistSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setUpdateAttributesOnlyOnValueChange(false);
config.setProcessingSettings(new Advanced(
ProcessingStrategy.onEveryMessage(),
ProcessingStrategy.onEveryMessage()
));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_ATTRIBUTES_REQUEST)
.originator(deviceId)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of(NOTIFY_DEVICE_METADATA_KEY, "false")))
.build();
// WHEN-THEN
var expectedSaveRequest = builder()
.tenantId(tenantId)
.entityId(msg.getOriginator())
.scope(AttributeScope.valueOf(config.getScope()))
.entry(new DoubleDataEntry("temperature", 22.3))
.notifyDevice(false)
.strategy(Strategy.PROCESS_ALL)
.build();
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(1)).saveAttributes(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest)
.usingRecursiveComparison()
.ignoringFields("callback", "entries.lastUpdateTs")
.isEqualTo(expectedSaveRequest)
));
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(times(2)).saveAttributes(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest)
.usingRecursiveComparison()
.ignoringFields("callback", "entries.lastUpdateTs")
.isEqualTo(expectedSaveRequest)
));
}
@Test
void givenAdvancedProcessingSettingsWithDifferentDeduplicateStrategyForEachAction_whenOnMsg_thenEvaluatesStrategiesForEachActionsIndependently() throws TbNodeException {
// GIVEN
config.setUpdateAttributesOnlyOnValueChange(false);
config.setProcessingSettings(new Advanced(
ProcessingStrategy.deduplicate(1),
ProcessingStrategy.deduplicate(2)
));
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_ATTRIBUTES_REQUEST)
.originator(deviceId)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of("ts", Long.toString(ts1))))
.build());
then(telemetryServiceMock).should().saveAttributes(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest.getStrategy()).isEqualTo(Strategy.PROCESS_ALL)
));
clearInvocations(telemetryServiceMock);
node.onMsg(ctxMock, TbMsg.newMsg()
.type(TbMsgType.POST_ATTRIBUTES_REQUEST)
.originator(deviceId)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of("ts", Long.toString(ts2))))
.build());
then(telemetryServiceMock).should().saveAttributes(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest.getStrategy()).isEqualTo(new Strategy(true, false))
));
clearInvocations(telemetryServiceMock);
node.onMsg(ctxMock, TbMsg.newMsg()
.type(TbMsgType.POST_ATTRIBUTES_REQUEST)
.originator(deviceId)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of("ts", Long.toString(ts3))))
.build());
then(telemetryServiceMock).should().saveAttributes(assertArg(
actualSaveRequest -> assertThat(actualSaveRequest.getStrategy()).isEqualTo(Strategy.PROCESS_ALL)
));
}
@Test
public void givenAdvancedProcessingSettingsWithSkipStrategiesForAllActionsAndSameMessageTwoTimes_whenOnMsg_thenSkipsSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setProcessingSettings(new Advanced(
ProcessingStrategy.skip(),
ProcessingStrategy.skip()
));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_ATTRIBUTES_REQUEST)
.originator(deviceId)
.data(JacksonUtil.newObjectNode().put("temperature", 22.3).toString())
.metaData(new TbMsgMetaData(Map.of(NOTIFY_DEVICE_METADATA_KEY, "false")))
.build();
// WHEN-THEN
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(never()).saveAttributes(any());
then(ctxMock).should(times(1)).tellSuccess(msg);
node.onMsg(ctxMock, msg);
then(telemetryServiceMock).should(never()).saveAttributes(any());
then(ctxMock).should(times(2)).tellSuccess(msg);
}
@Test
void givenVariousChangesToAttributes_whenUpdateOnlyOnValueChangeEnabled_thenShouldCorrectlyFilterChangedAttributes() throws TbNodeException {
// GIVEN
config.setUpdateAttributesOnlyOnValueChange(true);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
List<AttributeKvEntry> currentAttributes = List.of(
new BaseAttributeKvEntry(1694000000L, new StringDataEntry("address", "Peremohy ave 1")),
new BaseAttributeKvEntry(1694000000L, new BooleanDataEntry("valid", true)),
new BaseAttributeKvEntry(1694000000L, new LongDataEntry("counter", 100L)),
new BaseAttributeKvEntry(1694000000L, new DoubleDataEntry("temp", -18.35)),
new BaseAttributeKvEntry(1694000000L, new JsonDataEntry("json", "{\"warning\":\"out of paper\"}"))
);
List<AttributeKvEntry> newAttributes = List.of(
new BaseAttributeKvEntry(1694000999L, new JsonDataEntry("json", "{\"status\":\"OK\"}")), // value changed, reordered
new BaseAttributeKvEntry(1694000999L, new StringDataEntry("valid", "true")), //type changed
new BaseAttributeKvEntry(1694000999L, new LongDataEntry("counter", 101L)), //value changed
new BaseAttributeKvEntry(1694000999L, new DoubleDataEntry("temp", -18.35)),
new BaseAttributeKvEntry(1694000999L, new StringDataEntry("address", "Peremohy ave 1")) // reordered
new BaseAttributeKvEntry(123L, new StringDataEntry("address", "Prospect Beresteiskyi 1")),
new BaseAttributeKvEntry(123L, new BooleanDataEntry("valid", true)),
new BaseAttributeKvEntry(123L, new LongDataEntry("counter", 100L)),
new BaseAttributeKvEntry(123L, new DoubleDataEntry("temp", -18.35)),
new BaseAttributeKvEntry(123L, new JsonDataEntry("json", "{\"warning\":\"out of paper\"}"))
);
List<AttributeKvEntry> expected = List.of(
new BaseAttributeKvEntry(1694000999L, new StringDataEntry("valid", "true")),
new BaseAttributeKvEntry(1694000999L, new LongDataEntry("counter", 101L)),
new BaseAttributeKvEntry(1694000999L, new JsonDataEntry("json", "{\"status\":\"OK\"}"))
given(attributesServiceMock.find(eq(tenantId), eq(deviceId), eq(AttributeScope.valueOf(config.getScope())), anyList())).willReturn(immediateFuture(currentAttributes));
var data = JacksonUtil.newObjectNode()
.put("address", "Prospect Beresteiskyi 1") // no changes
.put("valid", "false") // type and value changed
.put("counter", 101L) // value changed
.put("temp", -18.35) // no changes
.put("json", "{\"warning\":\"out of paper\"}") // only type changed
.put("newKey", "newValue"); // new attribute
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_ATTRIBUTES_REQUEST)
.originator(deviceId)
.data(data.toString())
.metaData(TbMsgMetaData.EMPTY)
.build();
// WHEN
node.onMsg(ctxMock, msg);
// THEN
List<AttributeKvEntry> expectedChangedAttributes = List.of(
new BaseAttributeKvEntry(456L, new StringDataEntry("valid", "false")),
new BaseAttributeKvEntry(456L, new LongDataEntry("counter", 101L)),
new BaseAttributeKvEntry(456L, new StringDataEntry("json", "{\"warning\":\"out of paper\"}")),
new BaseAttributeKvEntry(456L, new StringDataEntry("newKey", "newValue"))
);
List<AttributeKvEntry> filtered = node.filterChangedAttr(currentAttributes, newAttributes);
assertThat(filtered).containsExactlyInAnyOrderElementsOf(expected);
then(telemetryServiceMock).should().saveAttributes(assertArg(request ->
assertThat(request.getEntries())
.usingRecursiveComparison()
.ignoringCollectionOrder()
.ignoringFields("lastUpdateTs")
.isEqualTo(expectedChangedAttributes)
));
}
// Notify device backward-compatibility test arguments
private static Stream<Arguments> givenNotifyDeviceMdValue_whenSaveAndNotify_thenVerifyExpectedArgumentForNotifyDeviceInSaveAndNotifyMethod() {
return Stream.of(
Arguments.of(null, true),
Arguments.of("null", false),
Arguments.of("true", true),
Arguments.of("false", false)
@Test
void givenNoChangesToAttributes_whenUpdateOnlyOnValueChangeEnabled_thenShouldNotCallSaveAndJustTellSuccess() throws TbNodeException {
// GIVEN
config.setUpdateAttributesOnlyOnValueChange(true);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
List<AttributeKvEntry> currentAttributes = List.of(
new BaseAttributeKvEntry(123L, new StringDataEntry("address", "Prospect Beresteiskyi 1")),
new BaseAttributeKvEntry(123L, new BooleanDataEntry("valid", true)),
new BaseAttributeKvEntry(123L, new LongDataEntry("counter", 100L))
);
given(attributesServiceMock.find(eq(tenantId), eq(deviceId), eq(AttributeScope.valueOf(config.getScope())), anyList())).willReturn(immediateFuture(currentAttributes));
var data = JacksonUtil.newObjectNode()
.put("address", "Prospect Beresteiskyi 1")
.put("valid", true)
.put("counter", 100L);
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_ATTRIBUTES_REQUEST)
.originator(deviceId)
.data(data.toString())
.metaData(TbMsgMetaData.EMPTY)
.build();
// WHEN
node.onMsg(ctxMock, msg);
// THEN
then(telemetryServiceMock).shouldHaveNoInteractions();
then(ctxMock).should().tellSuccess(msg);
}
// Notify device backward-compatibility test
@ParameterizedTest
@MethodSource
void givenNotifyDeviceMdValue_whenSaveAndNotify_thenVerifyExpectedArgumentForNotifyDeviceInSaveAndNotifyMethod(String mdValue, boolean expectedArgumentValue) throws TbNodeException {
var ctxMock = mock(TbContext.class);
var telemetryServiceMock = mock(RuleEngineTelemetryService.class);
ObjectNode defaultConfig = (ObjectNode) JacksonUtil.valueToTree(new TbMsgAttributesNodeConfiguration().defaultConfiguration());
defaultConfig.put("notifyDevice", false);
var tbNodeConfiguration = new TbNodeConfiguration(defaultConfig);
assertThat(defaultConfig.has("notifyDevice")).as("pre condition has notifyDevice").isTrue();
when(ctxMock.getTenantId()).thenReturn(tenantId);
when(ctxMock.getTelemetryService()).thenReturn(telemetryServiceMock);
willCallRealMethod().given(node).init(any(TbContext.class), any(TbNodeConfiguration.class));
willCallRealMethod().given(node).saveAttr(any(), eq(ctxMock), any(TbMsg.class), any(AttributeScope.class), anyBoolean());
node.init(ctxMock, tbNodeConfiguration);
TbMsgMetaData md = new TbMsgMetaData();
if (mdValue != null) {
md.putValue(NOTIFY_DEVICE_METADATA_KEY, mdValue);
}
// dummy list with one ts kv to pass the empty list check.
var testTbMsg = TbMsg.newMsg()
.type(TbMsgType.POST_TELEMETRY_REQUEST)
void givenVariousValuesForNotifyDeviceInMetadata_thenShouldCorrectlyParseValueFromMetadata(String mdValue, boolean expectedArgumentValue) throws TbNodeException {
// GIVEN
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
given(attributesServiceMock.find(tenantId, deviceId, AttributeScope.valueOf(config.getScope()), List.of("mode"))).willReturn(
immediateFuture(List.of(new BaseAttributeKvEntry(123L, new StringDataEntry("mode", "tilt"))))
);
var metadata = new TbMsgMetaData();
metadata.putValue(NOTIFY_DEVICE_METADATA_KEY, mdValue);
var msg = TbMsg.newMsg()
.type(TbMsgType.POST_ATTRIBUTES_REQUEST)
.originator(deviceId)
.copyMetaData(md)
.data(TbMsg.EMPTY_STRING)
.data(JacksonUtil.newObjectNode().put("mode", "vibration").toString())
.metaData(metadata)
.build();
List<AttributeKvEntry> testAttrList = List.of(new BaseAttributeKvEntry(0L, new StringDataEntry("testKey", "testValue")));
node.saveAttr(testAttrList, ctxMock, testTbMsg, AttributeScope.SHARED_SCOPE, false);
// WHEN
node.onMsg(ctxMock, msg);
verify(telemetryServiceMock, times(1)).saveAttributes(assertArg(request -> {
assertThat(request.getTenantId()).isEqualTo(tenantId);
assertThat(request.getEntityId()).isEqualTo(deviceId);
assertThat(request.getScope()).isEqualTo(AttributeScope.SHARED_SCOPE);
assertThat(request.getEntries()).isEqualTo(testAttrList);
assertThat(request.isNotifyDevice()).isEqualTo(expectedArgumentValue);
}));
// THEN
then(telemetryServiceMock).should().saveAttributes(assertArg(request -> assertThat(request.isNotifyDevice()).isEqualTo(expectedArgumentValue)));
}
// Notify device backward-compatibility test arguments
static Stream<Arguments> givenVariousValuesForNotifyDeviceInMetadata_thenShouldCorrectlyParseValueFromMetadata() {
return Stream.of(
Arguments.of(null, true),
Arguments.of("null", false),
Arguments.of("true", true),
Arguments.of("false", false)
);
}
// Rule nodes upgrade
private static Stream<Arguments> givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() {
static Stream<Arguments> givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() {
return Stream.of(
// default config for version 0
Arguments.of(0,
"{\"scope\":\"CLIENT_SCOPE\",\"notifyDevice\":\"false\",\"sendAttributesUpdatedNotification\":\"false\"}",
"""
{
"scope": "CLIENT_SCOPE",
"notifyDevice": "false",
"sendAttributesUpdatedNotification": "false"
}
""",
true,
"{\"scope\":\"CLIENT_SCOPE\",\"notifyDevice\":false,\"sendAttributesUpdatedNotification\":false,\"updateAttributesOnlyOnValueChange\":false}"),
"""
{
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "CLIENT_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": false
}
"""
),
// default config for version 1 with upgrade from version 0
Arguments.of(0,
"{\"scope\":\"CLIENT_SCOPE\",\"notifyDevice\":false,\"sendAttributesUpdatedNotification\":false,\"updateAttributesOnlyOnValueChange\":true}",
false,
"{\"scope\":\"CLIENT_SCOPE\",\"notifyDevice\":false,\"sendAttributesUpdatedNotification\":false,\"updateAttributesOnlyOnValueChange\":true}"),
"""
{
"scope": "CLIENT_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
""",
true,
"""
{
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "CLIENT_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
"""
),
// all flags are booleans
Arguments.of(1,
"{\"scope\":\"SHARED_SCOPE\",\"notifyDevice\":true,\"sendAttributesUpdatedNotification\":false,\"updateAttributesOnlyOnValueChange\":true}",
false,
"{\"scope\":\"SHARED_SCOPE\",\"notifyDevice\":true,\"sendAttributesUpdatedNotification\":false,\"updateAttributesOnlyOnValueChange\":true}"),
"""
{
"scope": "SHARED_SCOPE",
"notifyDevice": true,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
""",
true,
"""
{
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "SHARED_SCOPE",
"notifyDevice": true,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
"""
),
// no boolean flags set
Arguments.of(1,
"{\"scope\":\"CLIENT_SCOPE\"}",
"""
{
"scope": "CLIENT_SCOPE"
}
""",
true,
"{\"scope\":\"CLIENT_SCOPE\",\"notifyDevice\":true,\"sendAttributesUpdatedNotification\":false,\"updateAttributesOnlyOnValueChange\":true}"),
"""
{
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "CLIENT_SCOPE",
"notifyDevice": true,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
"""
),
// all flags are boolean strings
Arguments.of(1,
"{\"scope\":\"CLIENT_SCOPE\",\"notifyDevice\":\"false\",\"sendAttributesUpdatedNotification\":\"false\",\"updateAttributesOnlyOnValueChange\":\"true\"}",
"""
{
"scope": "CLIENT_SCOPE",
"notifyDevice": "false",
"sendAttributesUpdatedNotification": "false",
"updateAttributesOnlyOnValueChange": "true"
}
""",
true,
"{\"scope\":\"CLIENT_SCOPE\",\"notifyDevice\":false,\"sendAttributesUpdatedNotification\":false,\"updateAttributesOnlyOnValueChange\":true}"),
"""
{
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "CLIENT_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
"""
),
// at least one flag is boolean string
Arguments.of(1,
"{\"scope\":\"CLIENT_SCOPE\",\"notifyDevice\":\"false\",\"sendAttributesUpdatedNotification\":false,\"updateAttributesOnlyOnValueChange\":true}",
"""
{
"scope": "CLIENT_SCOPE",
"notifyDevice": "false",
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
""",
true,
"{\"scope\":\"CLIENT_SCOPE\",\"notifyDevice\":false,\"sendAttributesUpdatedNotification\":false,\"updateAttributesOnlyOnValueChange\":true}"),
"""
{
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "CLIENT_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
"""
),
// notify device flag is null
Arguments.of(1,
"{\"scope\":\"CLIENT_SCOPE\",\"notifyDevice\":\"null\",\"sendAttributesUpdatedNotification\":false,\"updateAttributesOnlyOnValueChange\":true}",
"""
{
"scope": "CLIENT_SCOPE",
"notifyDevice": "null",
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
""",
true,
"""
{
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "CLIENT_SCOPE",
"notifyDevice": true,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
"""
),
// default config for version 2
Arguments.of(2,
"""
{
"scope": "SERVER_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
""",
true,
"{\"scope\":\"CLIENT_SCOPE\",\"notifyDevice\":true,\"sendAttributesUpdatedNotification\":false,\"updateAttributesOnlyOnValueChange\":true}")
"""
{
"processingSettings": {
"type": "ON_EVERY_MESSAGE"
},
"scope": "SERVER_SCOPE",
"notifyDevice": false,
"sendAttributesUpdatedNotification": false,
"updateAttributesOnlyOnValueChange": true
}
"""
)
);
}

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

@ -73,6 +73,10 @@ 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.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.Advanced;
import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.Deduplicate;
import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.OnEveryMessage;
import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.WebSocketsOnly;
@ExtendWith(MockitoExtension.class)
public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
@ -110,7 +114,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
@Test
public void verifyDefaultConfig() {
assertThat(config.getDefaultTTL()).isEqualTo(0L);
assertThat(config.getProcessingSettings()).isInstanceOf(TbMsgTimeseriesNodeConfiguration.ProcessingSettings.OnEveryMessage.class);
assertThat(config.getProcessingSettings()).isInstanceOf(OnEveryMessage.class);
assertThat(config.isUseServerTs()).isFalse();
}
@ -223,7 +227,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
var timeseriesStrategy = ProcessingStrategy.onEveryMessage();
var latestStrategy = ProcessingStrategy.skip();
var webSockets = ProcessingStrategy.onEveryMessage();
var processingSettings = new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Advanced(timeseriesStrategy, latestStrategy, webSockets);
var processingSettings = new Advanced(timeseriesStrategy, latestStrategy, webSockets);
config.setProcessingSettings(processingSettings);
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -335,7 +339,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
@Test
public void givenOnEveryMessageProcessingSettingsAndSameMessageTwoTimes_whenOnMsg_thenPersistSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.OnEveryMessage());
config.setProcessingSettings(new OnEveryMessage());
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -370,7 +374,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
@Test
public void givenDeduplicateProcessingSettingsAndSameMessageTwoTimes_whenOnMsg_thenPersistThisMessageOnlyFirstTime() throws TbNodeException {
// GIVEN
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Deduplicate(10));
config.setProcessingSettings(new Deduplicate(10));
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -405,7 +409,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
@Test
public void givenWebSocketsOnlyProcessingSettingsAndSameMessageTwoTimes_whenOnMsg_thenSendsOnlyWsUpdateTwoTimes() throws TbNodeException {
// GIVEN
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.WebSocketsOnly());
config.setProcessingSettings(new WebSocketsOnly());
node.init(ctxMock, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
@ -440,7 +444,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
@Test
public void givenAdvancedProcessingSettingsWithOnEveryMessageStrategiesForAllActionsAndSameMessageTwoTimes_whenOnMsg_thenPersistSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Advanced(
config.setProcessingSettings(new Advanced(
ProcessingStrategy.onEveryMessage(),
ProcessingStrategy.onEveryMessage(),
ProcessingStrategy.onEveryMessage()
@ -479,7 +483,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
@Test
public void givenAdvancedProcessingSettingsWithDifferentDeduplicateStrategyForEachAction_whenOnMsg_thenEvaluatesStrategiesForEachActionsIndependently() throws TbNodeException {
// GIVEN
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Advanced(
config.setProcessingSettings(new Advanced(
ProcessingStrategy.deduplicate(1),
ProcessingStrategy.deduplicate(2),
ProcessingStrategy.deduplicate(3)
@ -530,7 +534,7 @@ public class TbMsgTimeseriesNodeTest extends AbstractRuleNodeUpgradeTest {
@Test
public void givenAdvancedProcessingSettingsWithSkipStrategiesForAllActionsAndSameMessageTwoTimes_whenOnMsg_thenSkipsSameMessageTwoTimes() throws TbNodeException {
// GIVEN
config.setProcessingSettings(new TbMsgTimeseriesNodeConfiguration.ProcessingSettings.Advanced(
config.setProcessingSettings(new Advanced(
ProcessingStrategy.skip(),
ProcessingStrategy.skip(),
ProcessingStrategy.skip()

Loading…
Cancel
Save