Browse Source

review fixes

pull/14225/head
IrynaMatveieva 9 months ago
parent
commit
ce92740f1f
  1. 12
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
  2. 4
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java
  3. 6
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
  4. 2
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesImmediateOutputStrategy.java
  5. 4
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesOutput.java
  6. 6
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesOutputStrategy.java
  7. 2
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesRuleChainOutputStrategy.java
  8. 37
      dao/src/main/java/org/thingsboard/server/dao/util/KvUtils.java
  9. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java
  10. 4
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java
  11. 60
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/TelemetryUtil.java

12
application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java

@ -37,7 +37,7 @@ import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.ArgumentType;
import org.thingsboard.server.common.data.cf.configuration.AttributeImmediateOutputStrategy;
import org.thingsboard.server.common.data.cf.configuration.AttributesImmediateOutputStrategy;
import org.thingsboard.server.common.data.cf.configuration.OutputStrategy;
import org.thingsboard.server.common.data.cf.configuration.OutputType;
import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration;
@ -79,12 +79,12 @@ import java.util.concurrent.ExecutionException;
import java.util.function.Predicate;
import java.util.stream.Collectors;
import static org.thingsboard.rule.engine.util.TelemetryUtil.filterChangedAttr;
import static org.thingsboard.rule.engine.util.TelemetryUtil.toTsKvEntryList;
import static org.thingsboard.server.common.data.cf.CalculatedFieldType.PROPAGATION;
import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT;
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY;
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY;
import static org.thingsboard.server.dao.util.KvUtils.filterChangedAttr;
import static org.thingsboard.server.dao.util.KvUtils.toTsKvEntryList;
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultAttributeEntry;
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry;
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformSingleValueArgument;
@ -379,7 +379,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
}
private void saveAttributes(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, AttributeScope scope, List<CalculatedFieldId> cfIds, SettableFuture<Void> future) {
if (!(outputStrategy instanceof AttributeImmediateOutputStrategy attOutputStrategy)) {
if (!(outputStrategy instanceof AttributesImmediateOutputStrategy attOutputStrategy)) {
future.setException(new IllegalArgumentException("Only AttributeImmediateOutputStrategy is supported."));
} else {
AttributesSaveRequest.Strategy strategy = new Strategy(attOutputStrategy.isSaveAttribute(), attOutputStrategy.isSendWsUpdate(), attOutputStrategy.isProcessCfs());
@ -413,7 +413,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
List<AttributeKvEntry> entries,
AttributesSaveRequest.Strategy strategy,
SettableFuture<Void> future) {
tsSubService.saveAttributesInternal(AttributesSaveRequest.builder()
tsSubService.saveAttributes(AttributesSaveRequest.builder()
.tenantId(tenantId)
.entityId(entityId)
.scope(scope)
@ -452,7 +452,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
if (cfIds != null && !cfIds.isEmpty()) {
builder.previousCalculatedFieldIds(cfIds);
}
tsSubService.saveTimeseriesInternal(builder.build());
tsSubService.saveTimeseries(builder.build());
}
}

4
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java

@ -39,12 +39,8 @@ public interface CalculatedFieldProcessingService {
Map<String, ArgumentEntry> fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map<String, Argument> arguments);
void processImmediately(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback);
void processResult(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback);
void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback);
void pushMsgToLinks(CalculatedFieldTelemetryMsg msg, List<CalculatedFieldEntityCtxId> linkedCalculatedFields, TbCallback callback);
}

6
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java

@ -143,8 +143,7 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF
}
}
@Override
public void processImmediately(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback) {
private void processImmediately(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback) {
if (result instanceof TelemetryCalculatedFieldResult telemetryResult) {
saveTelemetryResult(tenantId, entityId, telemetryResult, cfIds, callback);
return;
@ -157,8 +156,7 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF
callback.onSuccess();
}
@Override
public void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback) {
private void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback) {
if (result instanceof PropagationCalculatedFieldResult propagationResult) {
handlePropagationResults(propagationResult, callback,
(entity, res, cb) -> sendMsgToRuleEngine(tenantId, entityId, cb, res.toTbMsg(entity, cfIds)));

2
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeImmediateOutputStrategy.java → common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesImmediateOutputStrategy.java

@ -22,7 +22,7 @@ import lombok.NoArgsConstructor;
@Data
@AllArgsConstructor
@NoArgsConstructor
public class AttributeImmediateOutputStrategy implements AttributeOutputStrategy {
public class AttributesImmediateOutputStrategy implements AttributesOutputStrategy {
private boolean updateAttributesOnlyOnValueChange;

4
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesOutput.java

@ -25,10 +25,10 @@ public class AttributesOutput implements Output {
private AttributeScope scope;
private Integer decimalsByDefault;
private AttributeOutputStrategy strategy;
private AttributesOutputStrategy strategy;
public AttributesOutput() {
this.strategy = new AttributeRuleChainOutputStrategy();
this.strategy = new AttributesRuleChainOutputStrategy();
}
@Override

6
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeOutputStrategy.java → common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesOutputStrategy.java

@ -26,8 +26,8 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo;
property = "type"
)
@JsonSubTypes({
@JsonSubTypes.Type(value = AttributeImmediateOutputStrategy.class, name = "IMMEDIATE"),
@JsonSubTypes.Type(value = AttributeRuleChainOutputStrategy.class, name = "RULE_CHAIN"),
@JsonSubTypes.Type(value = AttributesImmediateOutputStrategy.class, name = "IMMEDIATE"),
@JsonSubTypes.Type(value = AttributesRuleChainOutputStrategy.class, name = "RULE_CHAIN"),
})
public interface AttributeOutputStrategy extends OutputStrategy {
public interface AttributesOutputStrategy extends OutputStrategy {
}

2
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeRuleChainOutputStrategy.java → common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesRuleChainOutputStrategy.java

@ -20,7 +20,7 @@ import lombok.NoArgsConstructor;
@Data
@NoArgsConstructor
public class AttributeRuleChainOutputStrategy implements AttributeOutputStrategy {
public class AttributesRuleChainOutputStrategy implements AttributesOutputStrategy {
@Override
public OutputStrategyType getType() {

37
dao/src/main/java/org/thingsboard/server/dao/util/KvUtils.java

@ -19,13 +19,21 @@ import com.fasterxml.jackson.databind.JsonNode;
import com.github.benmanes.caffeine.cache.Cache;
import com.github.benmanes.caffeine.cache.Caffeine;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.dao.exception.IncorrectParameterException;
import org.thingsboard.server.dao.service.NoXssValidator;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.TimeUnit;
import java.util.function.Function;
import java.util.stream.Collectors;
public class KvUtils {
@ -74,4 +82,33 @@ public class KvUtils {
}
}
}
public static List<TsKvEntry> toTsKvEntryList(Map<Long, List<KvEntry>> tsKvMap) {
List<TsKvEntry> tsKvEntryList = new ArrayList<>();
for (Map.Entry<Long, List<KvEntry>> tsKvEntry : tsKvMap.entrySet()) {
for (KvEntry kvEntry : tsKvEntry.getValue()) {
tsKvEntryList.add(new BasicTsKvEntry(tsKvEntry.getKey(), kvEntry));
}
}
return tsKvEntryList;
}
public static List<AttributeKvEntry> filterChangedAttr(List<AttributeKvEntry> currentAttributes, List<AttributeKvEntry> newAttributes) {
if (currentAttributes == null || currentAttributes.isEmpty()) {
return newAttributes;
}
Map<String, AttributeKvEntry> currentAttrMap = currentAttributes.stream()
.collect(Collectors.toMap(AttributeKvEntry::getKey, Function.identity(), (existing, replacement) -> existing));
return newAttributes.stream()
.filter(item -> {
AttributeKvEntry cacheAttr = currentAttrMap.get(item.getKey());
return cacheAttr == null
|| !Objects.equals(item.getValue(), cacheAttr.getValue()) //JSON and String can be equals by value, but different by type
|| !Objects.equals(item.getDataType(), cacheAttr.getDataType());
})
.collect(Collectors.toList());
}
}

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

@ -48,10 +48,10 @@ import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessin
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.rule.engine.util.TelemetryUtil.filterChangedAttr;
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;
import static org.thingsboard.server.dao.util.KvUtils.filterChangedAttr;
@RuleNode(
type = ComponentType.ACTION,

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

@ -31,7 +31,6 @@ 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;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.plugin.ComponentType;
@ -39,7 +38,6 @@ import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileCon
import org.thingsboard.server.common.data.util.TbPair;
import org.thingsboard.server.common.msg.TbMsg;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.UUID;
@ -49,8 +47,8 @@ import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessin
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.rule.engine.util.TelemetryUtil.toTsKvEntryList;
import static org.thingsboard.server.common.data.msg.TbMsgType.POST_TELEMETRY_REQUEST;
import static org.thingsboard.server.dao.util.KvUtils.toTsKvEntryList;
@RuleNode(
type = ComponentType.ACTION,

60
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/TelemetryUtil.java

@ -1,60 +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.util;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.function.Function;
import java.util.stream.Collectors;
public class TelemetryUtil {
public static List<TsKvEntry> toTsKvEntryList(Map<Long, List<KvEntry>> tsKvMap) {
List<TsKvEntry> tsKvEntryList = new ArrayList<>();
for (Map.Entry<Long, List<KvEntry>> tsKvEntry : tsKvMap.entrySet()) {
for (KvEntry kvEntry : tsKvEntry.getValue()) {
tsKvEntryList.add(new BasicTsKvEntry(tsKvEntry.getKey(), kvEntry));
}
}
return tsKvEntryList;
}
public static List<AttributeKvEntry> filterChangedAttr(List<AttributeKvEntry> currentAttributes, List<AttributeKvEntry> newAttributes) {
if (currentAttributes == null || currentAttributes.isEmpty()) {
return newAttributes;
}
Map<String, AttributeKvEntry> currentAttrMap = currentAttributes.stream()
.collect(Collectors.toMap(AttributeKvEntry::getKey, Function.identity(), (existing, replacement) -> existing));
return newAttributes.stream()
.filter(item -> {
AttributeKvEntry cacheAttr = currentAttrMap.get(item.getKey());
return cacheAttr == null
|| !Objects.equals(item.getValue(), cacheAttr.getValue()) //JSON and String can be equals by value, but different by type
|| !Objects.equals(item.getDataType(), cacheAttr.getDataType());
})
.collect(Collectors.toList());
}
}
Loading…
Cancel
Save