Browse Source

added calculated field execution service

pull/12092/head
IrynaMatveieva 2 years ago
parent
commit
c39a373038
  1. 25
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java
  2. 309
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java
  3. 16
      application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java
  4. 9
      application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java
  5. 21
      application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtxId.java
  6. 46
      application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldResult.java
  7. 9
      application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldState.java
  8. 258
      application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java
  9. 2
      application/src/main/java/org/thingsboard/server/service/entitiy/cf/RocksDBService.java
  10. 48
      application/src/main/java/org/thingsboard/server/service/entitiy/cf/ScriptCalculatedFieldState.java
  11. 37
      application/src/main/java/org/thingsboard/server/service/entitiy/cf/SimpleCalculatedFieldState.java
  12. 4
      application/src/main/java/org/thingsboard/server/service/entitiy/cf/TbCalculatedFieldService.java
  13. 10
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  14. 3
      application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java
  15. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java
  16. 32
      common/data/src/main/java/org/thingsboard/server/common/data/cf/Argument.java
  17. 8
      common/data/src/main/java/org/thingsboard/server/common/data/cf/BaseCalculatedFieldConfiguration.java
  18. 2
      common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldConfiguration.java
  19. 5
      dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java
  20. 2
      dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java
  21. 5
      dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java
  22. 3
      dao/src/test/java/org/thingsboard/server/dao/service/AssetServiceTest.java
  23. 3
      dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java
  24. 3
      dao/src/test/java/org/thingsboard/server/dao/service/CustomerServiceTest.java
  25. 3
      dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java

25
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java

@ -0,0 +1,25 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.cf;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.gen.transport.TransportProtos;
public interface CalculatedFieldExecutionService {
void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback);
}

309
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java

@ -0,0 +1,309 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.cf;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.ListeningExecutorService;
import com.google.common.util.concurrent.MoreExecutors;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.cf.Argument;
import org.thingsboard.server.common.data.cf.BaseCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.AssetProfileId;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.page.PageDataIterable;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.cf.CalculatedFieldService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.entitiy.cf.CalculatedFieldCtx;
import org.thingsboard.server.service.entitiy.cf.CalculatedFieldCtxId;
import org.thingsboard.server.service.entitiy.cf.CalculatedFieldState;
import org.thingsboard.server.service.entitiy.cf.RocksDBService;
import org.thingsboard.server.service.entitiy.cf.ScriptCalculatedFieldState;
import org.thingsboard.server.service.entitiy.cf.SimpleCalculatedFieldState;
import org.thingsboard.server.service.partition.AbstractPartitionBasedService;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
@TbCoreComponent
@Service
@Slf4j
@RequiredArgsConstructor
public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBasedService<CalculatedFieldId> implements CalculatedFieldExecutionService {
private final CalculatedFieldService calculatedFieldService;
private final AssetService assetService;
private final DeviceService deviceService;
private final AttributesService attributesService;
private final TimeseriesService timeseriesService;
private final RocksDBService rocksDBService;
private ListeningExecutorService calculatedFieldExecutor;
private ListeningExecutorService calculatedFieldCallbackExecutor;
private final ConcurrentMap<CalculatedFieldId, CalculatedField> calculatedFields = new ConcurrentHashMap<>();
private final ConcurrentMap<CalculatedFieldId, List<CalculatedFieldLink>> calculatedFieldLinks = new ConcurrentHashMap<>();
private final ConcurrentMap<CalculatedFieldCtxId, CalculatedFieldCtx> states = new ConcurrentHashMap<>();
@Value("${calculatedField.initFetchPackSize:50000}")
@Getter
private int initFetchPackSize;
@PostConstruct
public void init() {
super.init();
calculatedFieldExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(
Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field"));
calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(
Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field-callback"));
scheduledExecutor.submit(this::fetchCalculatedFields);
}
@PreDestroy
public void stop() {
if (calculatedFieldExecutor != null) {
calculatedFieldExecutor.shutdownNow();
}
if (calculatedFieldCallbackExecutor != null) {
calculatedFieldCallbackExecutor.shutdownNow();
}
}
@Override
protected String getServiceName() {
return "Calculated Field Execution";
}
@Override
protected String getSchedulerExecutorName() {
return "calculated-field-scheduled";
}
@Override
protected Map<TopicPartitionInfo, List<ListenableFuture<?>>> onAddedPartitions(Set<TopicPartitionInfo> addedPartitions) {
// TODO: implementation for cluster mode
return Map.of();
}
@Override
protected void cleanupEntityOnPartitionRemoval(CalculatedFieldId entityId) {
// TODO: implementation for cluster mode
}
@Override
public void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback) {
try {
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB()));
CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB()));
log.info("Received CalculatedFieldMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId);
if (proto.getDeleted()) {
log.warn("Executing onCalculatedFieldDelete, calculatedFieldId=[{}]", calculatedFieldId);
onCalculatedFieldDelete(calculatedFieldId, callback);
callback.onSuccess();
}
CalculatedField cf = calculatedFieldService.findById(tenantId, calculatedFieldId);
if (proto.getUpdated()) {
log.info("Executing onCalculatedFieldUpdate, calculatedFieldId=[{}]", calculatedFieldId);
boolean shouldReinit = onCalculatedFieldUpdate(cf, callback);
if (!shouldReinit) {
return;
}
}
List<CalculatedFieldLink> links = calculatedFieldService.findAllCalculatedFieldLinksById(tenantId, calculatedFieldId);
if (cf != null) {
EntityId entityId = cf.getEntityId();
calculatedFields.put(calculatedFieldId, cf);
calculatedFieldLinks.put(calculatedFieldId, links);
switch (entityId.getEntityType()) {
case ASSET, DEVICE -> {
log.info("Initializing state for entity: tenantId=[{}], entityId=[{}]", tenantId, entityId);
initializeStateForEntity(tenantId, cf, entityId, callback);
}
case ASSET_PROFILE -> {
log.info("Initializing state for all assets in profile: tenantId=[{}], assetProfileId=[{}]", tenantId, entityId);
PageDataIterable<AssetId> assetIds = new PageDataIterable<>(pageLink ->
assetService.findAssetIdsByTenantIdAndAssetProfileId(tenantId, (AssetProfileId) entityId, pageLink), initFetchPackSize);
assetIds.forEach(assetId -> initializeStateForEntity(tenantId, cf, assetId, callback));
}
case DEVICE_PROFILE -> {
log.info("Initializing state for all devices in profile: tenantId=[{}], deviceProfileId=[{}]", tenantId, entityId);
PageDataIterable<DeviceId> deviceIds = new PageDataIterable<>(pageLink ->
deviceService.findDeviceIdsByTenantIdAndDeviceProfileId(tenantId, (DeviceProfileId) entityId, pageLink), initFetchPackSize);
deviceIds.forEach(deviceId -> initializeStateForEntity(tenantId, cf, deviceId, callback));
}
default ->
throw new IllegalArgumentException("Entity type '" + calculatedFieldId.getEntityType() + "' does not support calculated fields.");
}
} else {
//Calculated field was probably deleted while message was in queue;
log.warn("Calculated field not found, possibly deleted: {}", calculatedFieldId);
callback.onSuccess();
}
callback.onSuccess();
log.info("Successfully processed calculated field message for calculatedFieldId: [{}]", calculatedFieldId);
} catch (Exception e) {
log.trace("Failed to process calculated field msg: [{}]", proto, e);
callback.onFailure(e);
}
}
private boolean onCalculatedFieldUpdate(CalculatedField newCalculatedField, TbCallback callback) {
CalculatedField oldCalculatedField = calculatedFields.get(newCalculatedField.getId());
boolean shouldReinit = true;
if (hasSignificantChanges(oldCalculatedField, newCalculatedField)) {
onCalculatedFieldDelete(newCalculatedField.getId(), callback);
} else {
calculatedFields.put(newCalculatedField.getId(), newCalculatedField);
callback.onSuccess();
shouldReinit = false;
}
return shouldReinit;
}
private void onCalculatedFieldDelete(CalculatedFieldId calculatedFieldId, TbCallback callback) {
try {
calculatedFieldLinks.remove(calculatedFieldId);
calculatedFields.remove(calculatedFieldId);
states.keySet().removeIf(ctxId -> calculatedFields.keySet().stream().noneMatch(id -> ctxId.cfId().equals(id.getId())));
List<String> statesToRemove = states.keySet().stream()
.filter(ctxId -> !calculatedFields.containsKey(new CalculatedFieldId(ctxId.cfId())))
.map(JacksonUtil::writeValueAsString)
.toList();
rocksDBService.deleteAll(statesToRemove);
} catch (Exception e) {
log.trace("Failed to delete calculated field: [{}]", calculatedFieldId, e);
callback.onFailure(e);
}
}
private boolean hasSignificantChanges(CalculatedField oldCalculatedField, CalculatedField newCalculatedField) {
if (oldCalculatedField == null) {
return true;
}
boolean entityIdChanged = !oldCalculatedField.getEntityId().equals(newCalculatedField.getEntityId());
boolean typeChanged = !oldCalculatedField.getType().equals(newCalculatedField.getType());
CalculatedFieldConfiguration oldConfig = oldCalculatedField.getConfiguration();
CalculatedFieldConfiguration newConfig = newCalculatedField.getConfiguration();
boolean argumentsChanged = !oldConfig.getArguments().equals(newConfig.getArguments());
boolean outputTypeChanged = !oldConfig.getOutput().getType().equals(newConfig.getOutput().getType());
boolean outputExpressionChanged = !oldConfig.getOutput().getExpression().equals(newConfig.getOutput().getExpression());
return entityIdChanged || typeChanged || argumentsChanged || outputTypeChanged || outputExpressionChanged;
}
private void fetchCalculatedFields() {
PageDataIterable<CalculatedField> cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize);
cfs.forEach(cf -> calculatedFields.putIfAbsent(cf.getId(), cf));
PageDataIterable<CalculatedFieldLink> cfls = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFieldLinks, initFetchPackSize);
cfls.forEach(link -> calculatedFieldLinks.computeIfAbsent(link.getCalculatedFieldId(), id -> new ArrayList<>()).add(link));
rocksDBService.getAll().forEach((ctxId, ctx) -> states.put(JacksonUtil.fromString(ctxId, CalculatedFieldCtxId.class), JacksonUtil.fromString(ctx, CalculatedFieldCtx.class)));
states.keySet().removeIf(ctxId -> calculatedFields.keySet().stream().noneMatch(id -> ctxId.cfId().equals(id.getId())));
}
private void initializeStateForEntity(TenantId tenantId, CalculatedField calculatedField, EntityId entityId, TbCallback callback) {
Map<String, Argument> arguments = calculatedField.getConfiguration().getArguments();
Map<String, String> argumentValues = new HashMap<>();
arguments.forEach((key, argument) -> Futures.addCallback(fetchArgumentValue(tenantId, argument), new FutureCallback<>() {
@Override
public void onSuccess(Optional<? extends KvEntry> result) {
String value = result.map(KvEntry::getValueAsString).orElse(argument.getDefaultValue());
argumentValues.put(key, value);
}
@Override
public void onFailure(Throwable t) {
log.warn("Failed to initialize state for entity: [{}]", entityId, t);
callback.onFailure(t);
}
}, calculatedFieldCallbackExecutor));
updateOrInitializeState(calculatedField, entityId, argumentValues);
}
private ListenableFuture<Optional<? extends KvEntry>> fetchArgumentValue(TenantId tenantId, Argument argument) {
return switch (argument.getType()) {
case "ATTRIBUTES" -> Futures.transform(
attributesService.find(tenantId, argument.getEntityId(), AttributeScope.SERVER_SCOPE, argument.getKey()),
result -> result.map(entry -> (KvEntry) entry),
MoreExecutors.directExecutor());
case "TIME_SERIES" -> Futures.transform(
timeseriesService.findLatest(tenantId, argument.getEntityId(), argument.getKey()),
result -> result.map(entry -> (KvEntry) entry),
MoreExecutors.directExecutor());
default -> throw new IllegalArgumentException("Invalid argument type '" + argument.getType() + "'.");
};
}
private void updateOrInitializeState(CalculatedField calculatedField, EntityId entityId, Map<String, String> argumentValues) {
CalculatedFieldCtxId ctxId = new CalculatedFieldCtxId(calculatedField.getUuidId(), entityId.getId());
CalculatedFieldCtx calculatedFieldCtx = states.computeIfAbsent(ctxId, ctx -> new CalculatedFieldCtx(ctxId, null));
CalculatedFieldState state = calculatedFieldCtx.getState();
if (state == null) {
state = createStateByType(calculatedField.getType());
}
state.initState(argumentValues);
calculatedFieldCtx.setState(state);
states.put(ctxId, calculatedFieldCtx);
rocksDBService.put(JacksonUtil.writeValueAsString(ctxId), JacksonUtil.writeValueAsString(calculatedFieldCtx));
state.performCalculation(calculatedField.getConfiguration());
}
private CalculatedFieldState createStateByType(CalculatedFieldType calculatedFieldType) {
return switch (calculatedFieldType) {
case SIMPLE -> new SimpleCalculatedFieldState();
case SCRIPT -> new ScriptCalculatedFieldState();
};
}
}

16
application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java

@ -51,6 +51,8 @@ import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg;
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.msg.rule.engine.DeviceCredentialsUpdateNotificationMsg;
import org.thingsboard.server.dao.cf.CalculatedFieldService;
import org.thingsboard.server.dao.device.DeviceProfileService;
import org.thingsboard.server.dao.eventsourcing.ActionEntityEvent;
import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent;
import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent;
@ -65,6 +67,8 @@ public class EntityStateSourcingListener {
private final TbClusterService tbClusterService;
private final TenantService tenantService;
private final CalculatedFieldService calculatedFieldService;
private final DeviceProfileService deviceProfileService;
@PostConstruct
public void init() {
@ -102,7 +106,7 @@ public class EntityStateSourcingListener {
onTenantProfileUpdate(tenantProfile, lifecycleEvent);
}
case DEVICE -> {
onDeviceUpdate(event.getEntity(), event.getOldEntity());
onDeviceUpdate(tenantId, event.getEntity(), event.getOldEntity());
}
case DEVICE_PROFILE -> {
DeviceProfile deviceProfile = (DeviceProfile) event.getEntity();
@ -241,11 +245,19 @@ public class EntityStateSourcingListener {
tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED);
}
private void onDeviceUpdate(Object entity, Object oldEntity) {
private void onDeviceUpdate(TenantId tenantId, Object entity, Object oldEntity) {
Device device = (Device) entity;
Device oldDevice = null;
if (oldEntity instanceof Device) {
oldDevice = (Device) oldEntity;
// TODO: move verification of device type to cluster service
if (!oldDevice.getType().equals(device.getType())) {
DeviceProfile profile = deviceProfileService.findDeviceProfileByName(tenantId, device.getType());
boolean cfExistsByProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, profile.getId());
if (cfExistsByProfile) {
// TODO: send device type updated msg to core
}
}
}
tbClusterService.onDeviceUpdated(device, oldDevice);
}

9
application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtx.java

@ -22,16 +22,15 @@ import org.thingsboard.server.common.data.id.EntityId;
@Data
public class CalculatedFieldCtx {
private CalculatedFieldId calculatedFieldId;
private EntityId entityId;
private CalculatedFieldCtxId id;
private CalculatedFieldState state;
public CalculatedFieldCtx() {
}
public CalculatedFieldCtx(CalculatedFieldId calculatedFieldId, EntityId entityId, CalculatedFieldState state) {
this.calculatedFieldId = calculatedFieldId;
this.entityId = entityId;
public CalculatedFieldCtx(CalculatedFieldCtxId id, CalculatedFieldState state) {
this.id = id;
this.state = state;
}
}

21
application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldCtxId.java

@ -0,0 +1,21 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.entitiy.cf;
import java.util.UUID;
public record CalculatedFieldCtxId(UUID cfId, UUID entityId) {
}

46
application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldResult.java

@ -0,0 +1,46 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.entitiy.cf;
import lombok.Data;
import org.thingsboard.server.common.data.AttributeScope;
@Data
public class CalculatedFieldResult {
private String name;
private String type;
private AttributeScope scope;
private String value;
public static CalculatedFieldResult createAttributesResult(String name, AttributeScope scope, String value) {
CalculatedFieldResult result = new CalculatedFieldResult();
result.name = name;
result.type = "ATTRIBUTES";
result.scope = scope;
result.value = value;
return result;
}
public static CalculatedFieldResult createTimeSeriesResult(String name, String value) {
CalculatedFieldResult result = new CalculatedFieldResult();
result.name = name;
result.type = "TIME_SERIES";
result.value = value;
return result;
}
}

9
application/src/main/java/org/thingsboard/server/service/entitiy/cf/CalculatedFieldState.java

@ -18,6 +18,7 @@ package org.thingsboard.server.service.entitiy.cf;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import org.thingsboard.server.common.data.cf.BaseCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
@ -36,6 +37,12 @@ public interface CalculatedFieldState {
@JsonIgnore
CalculatedFieldType getType();
void performCalculation(Map<String, String> argumentValues, CalculatedFieldConfiguration calculatedFieldConfiguration, boolean initialCalculation);
default boolean isValid(Map<String, String> arguments, CalculatedFieldConfiguration calculatedFieldConfiguration) {
return arguments.keySet().containsAll(calculatedFieldConfiguration.getArguments().keySet());
}
void initState(Map<String, String> argumentValues);
CalculatedFieldResult performCalculation(CalculatedFieldConfiguration calculatedFieldConfiguration);
}

258
application/src/main/java/org/thingsboard/server/service/entitiy/cf/DefaultTbCalculatedFieldService.java

@ -15,65 +15,28 @@
*/
package org.thingsboard.server.service.entitiy.cf;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.ListeningExecutorService;
import com.google.common.util.concurrent.ListeningScheduledExecutorService;
import com.google.common.util.concurrent.MoreExecutors;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.HasTenantId;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.cf.BaseCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.CalculatedFieldLink;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.AssetProfileId;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.HasId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.page.PageDataIterable;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.cf.CalculatedFieldService;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.entitiy.AbstractTbEntityService;
import org.thingsboard.server.service.security.model.SecurityUser;
import org.thingsboard.server.service.security.permission.Operation;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import static org.thingsboard.server.dao.service.Validator.validateEntityId;
@ -84,109 +47,6 @@ import static org.thingsboard.server.dao.service.Validator.validateEntityId;
public class DefaultTbCalculatedFieldService extends AbstractTbEntityService implements TbCalculatedFieldService {
private final CalculatedFieldService calculatedFieldService;
private final AttributesService attributesService;
private final TimeseriesService timeseriesService;
private final RocksDBService rocksDBService;
private ListeningScheduledExecutorService scheduledExecutor;
private ListeningExecutorService calculatedFieldExecutor;
private ListeningExecutorService calculatedFieldCallbackExecutor;
private final ConcurrentMap<CalculatedFieldId, CalculatedField> calculatedFields = new ConcurrentHashMap<>();
private final ConcurrentMap<CalculatedFieldId, List<CalculatedFieldLink>> calculatedFieldLinks = new ConcurrentHashMap<>();
private final ConcurrentMap<String, CalculatedFieldCtx> states = new ConcurrentHashMap<>();
@Value("${calculatedField.initFetchPackSize:50000}")
@Getter
private int initFetchPackSize;
@Value("10")
@Getter
private int defaultCalculatedFieldCheckIntervalInSec;
@PostConstruct
public void init() {
// from AbstractPartitionBasedService
scheduledExecutor = MoreExecutors.listeningDecorator(Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("calculated-field-scheduled")));
///
calculatedFieldExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(
Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field"));
calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(
Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field-callback"));
scheduledExecutor.scheduleWithFixedDelay(this::fetchCalculatedFields, 0, defaultCalculatedFieldCheckIntervalInSec, TimeUnit.SECONDS);
}
@PreDestroy
public void stop() {
// from AbstractPartitionBasedService
if (scheduledExecutor != null) {
scheduledExecutor.shutdown();
}
///
if (calculatedFieldExecutor != null) {
calculatedFieldExecutor.shutdownNow();
}
if (calculatedFieldCallbackExecutor != null) {
calculatedFieldCallbackExecutor.shutdownNow();
}
}
@Override
public void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback) {
try {
TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB()));
CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB()));
log.info("Received CalculatedFieldMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId);
if (proto.getDeleted()) {
log.warn("Executing onCalculatedFieldDelete, calculatedFieldId=[{}]", calculatedFieldId);
onCalculatedFieldDelete(calculatedFieldId, callback);
callback.onSuccess();
}
CalculatedField cf = calculatedFieldService.findById(tenantId, calculatedFieldId);
if (proto.getUpdated()) {
log.info("Executing onCalculatedFieldUpdate, calculatedFieldId=[{}]", calculatedFieldId);
boolean shouldReinit = onCalculatedFieldUpdate(cf, callback);
if (!shouldReinit) {
return;
}
}
List<CalculatedFieldLink> links = calculatedFieldService.findAllCalculatedFieldLinksById(tenantId, calculatedFieldId);
if (cf != null) {
EntityId entityId = cf.getEntityId();
calculatedFields.put(calculatedFieldId, cf);
calculatedFieldLinks.put(calculatedFieldId, links);
switch (entityId.getEntityType()) {
case ASSET, DEVICE -> {
log.info("Initializing state for entity: tenantId=[{}], entityId=[{}]", tenantId, entityId);
initializeStateForEntity(tenantId, cf, entityId, callback);
}
case ASSET_PROFILE -> {
log.info("Initializing state for all assets in profile: tenantId=[{}], assetProfileId=[{}]", tenantId, entityId);
PageDataIterable<AssetId> assetIds = new PageDataIterable<>(pageLink ->
assetService.findAssetIdsByTenantIdAndAssetProfileId(tenantId, (AssetProfileId) entityId, pageLink), initFetchPackSize);
assetIds.forEach(assetId -> initializeStateForEntity(tenantId, cf, assetId, callback));
}
case DEVICE_PROFILE -> {
log.info("Initializing state for all devices in profile: tenantId=[{}], deviceProfileId=[{}]", tenantId, entityId);
PageDataIterable<DeviceId> deviceIds = new PageDataIterable<>(pageLink ->
deviceService.findDeviceIdsByTenantIdAndDeviceProfileId(tenantId, (DeviceProfileId) entityId, pageLink), initFetchPackSize);
deviceIds.forEach(deviceId -> initializeStateForEntity(tenantId, cf, deviceId, callback));
}
default ->
throw new IllegalArgumentException("Entity type '" + calculatedFieldId.getEntityType() + "' does not support calculated fields.");
}
} else {
//Calculated field was probably deleted while message was in queue;
log.warn("Calculated field not found, possibly deleted: {}", calculatedFieldId);
callback.onSuccess();
}
callback.onSuccess();
log.info("Successfully processed calculated field message for calculatedFieldId: [{}]", calculatedFieldId);
} catch (Exception e) {
log.trace("Failed to process calculated field msg: [{}]", proto, e);
callback.onFailure(e);
}
}
@Override
public CalculatedField save(CalculatedField calculatedField, SecurityUser user) throws ThingsboardException {
@ -224,58 +84,6 @@ public class DefaultTbCalculatedFieldService extends AbstractTbEntityService imp
}
}
private void onCalculatedFieldDelete(CalculatedFieldId calculatedFieldId, TbCallback callback) {
try {
calculatedFieldLinks.remove(calculatedFieldId);
calculatedFields.remove(calculatedFieldId);
states.keySet().removeIf(ctxId -> ctxId.startsWith(calculatedFieldId.getId().toString()));
List<String> statesToRemove = states.keySet().stream()
.filter(key -> key.startsWith(calculatedFieldId.getId().toString()))
.collect(Collectors.toList());
rocksDBService.deleteAll(statesToRemove);
} catch (Exception e) {
log.trace("Failed to delete calculated field: [{}]", calculatedFieldId, e);
callback.onFailure(e);
}
}
private boolean onCalculatedFieldUpdate(CalculatedField newCalculatedField, TbCallback callback) {
CalculatedField oldCalculatedField = calculatedFields.get(newCalculatedField.getId());
boolean shouldReinit = true;
if (hasSignificantChanges(oldCalculatedField, newCalculatedField)) {
onCalculatedFieldDelete(newCalculatedField.getId(), callback);
} else {
calculatedFields.put(newCalculatedField.getId(), newCalculatedField);
callback.onSuccess();
shouldReinit = false;
}
return shouldReinit;
}
private boolean hasSignificantChanges(CalculatedField oldCalculatedField, CalculatedField newCalculatedField) {
if (oldCalculatedField == null) {
return true;
}
boolean entityIdChanged = !oldCalculatedField.getEntityId().equals(newCalculatedField.getEntityId());
boolean typeChanged = !oldCalculatedField.getType().equals(newCalculatedField.getType());
CalculatedFieldConfiguration oldConfig = oldCalculatedField.getConfiguration();
CalculatedFieldConfiguration newConfig = newCalculatedField.getConfiguration();
boolean argumentsChanged = !oldConfig.getArguments().equals(newConfig.getArguments());
boolean outputTypeChanged = !oldConfig.getOutput().getType().equals(newConfig.getOutput().getType());
boolean outputExpressionChanged = !oldConfig.getOutput().getExpression().equals(newConfig.getOutput().getExpression());
return entityIdChanged || typeChanged || argumentsChanged || outputTypeChanged || outputExpressionChanged;
}
private void fetchCalculatedFields() {
PageDataIterable<CalculatedField> cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize);
cfs.forEach(cf -> calculatedFields.putIfAbsent(cf.getId(), cf));
PageDataIterable<CalculatedFieldLink> cfls = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFieldLinks, initFetchPackSize);
cfls.forEach(link -> calculatedFieldLinks.computeIfAbsent(link.getCalculatedFieldId(), id -> new ArrayList<>()).add(link));
rocksDBService.getAll().forEach((ctxId, ctx) -> states.put(ctxId, JacksonUtil.fromString(ctx, CalculatedFieldCtx.class)));
states.keySet().removeIf(ctxId -> calculatedFields.keySet().stream().noneMatch(id -> ctxId.startsWith(id.toString())));
}
private void checkEntityExistence(TenantId tenantId, EntityId entityId) {
switch (entityId.getEntityType()) {
case ASSET, DEVICE, ASSET_PROFILE, DEVICE_PROFILE ->
@ -305,70 +113,4 @@ public class DefaultTbCalculatedFieldService extends AbstractTbEntityService imp
};
}
private void initializeStateForEntity(TenantId tenantId, CalculatedField calculatedField, EntityId entityId, TbCallback callback) {
Map<String, BaseCalculatedFieldConfiguration.Argument> arguments = calculatedField.getConfiguration().getArguments();
Map<String, String> argumentValues = new HashMap<>();
arguments.forEach((key, argument) -> Futures.addCallback(fetchArgumentValue(tenantId, argument), new FutureCallback<>() {
@Override
public void onSuccess(Optional<? extends KvEntry> result) {
String value = result.map(KvEntry::getValueAsString).orElse(argument.getDefaultValue());
argumentValues.put(key, value);
}
@Override
public void onFailure(Throwable t) {
log.warn("Failed to initialize state for entity: [{}]", entityId, t);
callback.onFailure(t);
}
}, calculatedFieldCallbackExecutor));
updateOrInitializeState(calculatedField, entityId, argumentValues);
}
private ListenableFuture<Optional<? extends KvEntry>> fetchArgumentValue(TenantId tenantId, BaseCalculatedFieldConfiguration.Argument argument) {
return switch (argument.getType()) {
case "ATTRIBUTES" -> Futures.transform(
attributesService.find(tenantId, argument.getEntityId(), AttributeScope.SERVER_SCOPE, argument.getKey()),
result -> result.map(entry -> (KvEntry) entry),
MoreExecutors.directExecutor());
case "TIME_SERIES" -> Futures.transform(
timeseriesService.findLatest(tenantId, argument.getEntityId(), argument.getKey()),
result -> result.map(entry -> (KvEntry) entry),
MoreExecutors.directExecutor());
default -> throw new IllegalArgumentException("Invalid argument type '" + argument.getType() + "'.");
};
}
private void updateOrInitializeState(CalculatedField calculatedField, EntityId entityId, Map<String, String> argumentValues) {
String ctxId = generateCtxId(calculatedField.getId(), entityId);
CalculatedFieldCtx calculatedFieldCtx = states.computeIfAbsent(ctxId,
ctx -> new CalculatedFieldCtx(calculatedField.getId(), calculatedField.getEntityId(), null));
CalculatedFieldState state = calculatedFieldCtx.getState();
if (state != null) {
state.performCalculation(argumentValues, calculatedField.getConfiguration(), false);
} else {
CalculatedFieldState newState = createStateByType(calculatedField.getType());
newState.performCalculation(argumentValues, calculatedField.getConfiguration(), true);
state = newState;
}
calculatedFieldCtx.setState(state);
states.put(ctxId, calculatedFieldCtx);
rocksDBService.put(ctxId, Objects.requireNonNull(JacksonUtil.writeValueAsString(calculatedFieldCtx)));
}
private CalculatedFieldState createStateByType(CalculatedFieldType calculatedFieldType) {
return switch (calculatedFieldType) {
case SIMPLE -> new SimpleCalculatedFieldState();
default ->
throw new IllegalArgumentException("Invalid calculated field type '" + calculatedFieldType + "'.");
};
}
private String generateCtxId(CalculatedFieldId calculatedFieldId, EntityId entityId) {
return calculatedFieldId.getId() + "_" + entityId.getId();
}
}

2
application/src/main/java/org/thingsboard/server/service/entitiy/cf/RocksDBService.java

@ -21,6 +21,7 @@ import org.rocksdb.RocksDBException;
import org.rocksdb.RocksIterator;
import org.rocksdb.WriteBatch;
import org.rocksdb.WriteOptions;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.stereotype.Service;
import org.thingsboard.server.utils.RocksDBConfig;
@ -31,6 +32,7 @@ import java.util.Map;
@Service
@Slf4j
@ConditionalOnExpression("'${service.type:null}'=='monolith'")
public class RocksDBService {
private final RocksDB db;

48
application/src/main/java/org/thingsboard/server/service/entitiy/cf/ScriptCalculatedFieldState.java

@ -0,0 +1,48 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.entitiy.cf;
import lombok.Data;
import org.thingsboard.script.api.tbel.TbelInvokeService;
import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import java.util.HashMap;
import java.util.Map;
@Data
public class ScriptCalculatedFieldState implements CalculatedFieldState {
private TbelInvokeService tbelInvokeService;
private Map<String, String> arguments = new HashMap<>();
@Override
public CalculatedFieldType getType() {
return CalculatedFieldType.SCRIPT;
}
@Override
public void initState(Map<String, String> argumentValues) {
}
@Override
public CalculatedFieldResult performCalculation(CalculatedFieldConfiguration calculatedFieldConfiguration) {
return null;
}
}

37
application/src/main/java/org/thingsboard/server/service/entitiy/cf/SimpleCalculatedFieldState.java

@ -16,6 +16,8 @@
package org.thingsboard.server.service.entitiy.cf;
import lombok.Data;
import net.objecthunter.exp4j.Expression;
import net.objecthunter.exp4j.ExpressionBuilder;
import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
@ -26,8 +28,8 @@ import java.util.Map;
public class SimpleCalculatedFieldState implements CalculatedFieldState {
// TODO: use value object(TsKv) instead of string
Map<String, String> arguments = new HashMap<>();
String result;
private Map<String, String> arguments = new HashMap<>();
private String outputResult;
@Override
public CalculatedFieldType getType() {
@ -35,15 +37,30 @@ public class SimpleCalculatedFieldState implements CalculatedFieldState {
}
@Override
public void performCalculation(Map<String, String> argumentValues, CalculatedFieldConfiguration calculatedFieldConfiguration, boolean initialCalculation) {
if (initialCalculation) {
// todo: perform initial calculation
this.arguments = argumentValues;
} else {
// todo: perform calculation based on previous data
this.arguments.putAll(argumentValues);
public void initState(Map<String, String> argumentValues) {
this.arguments = argumentValues;
}
@Override
public CalculatedFieldResult performCalculation(CalculatedFieldConfiguration calculatedFieldConfiguration) {
if (isValid(arguments, calculatedFieldConfiguration)) {
String expression = calculatedFieldConfiguration.getOutput().getExpression();
ThreadLocal<Expression> customExpression = new ThreadLocal<>();
var expr = customExpression.get();
if (expr == null) {
expr = new ExpressionBuilder(expression)
.implicitMultiplication(true)
.variables(arguments.keySet())
.build();
customExpression.set(expr);
}
Map<String, Double> variables = new HashMap<>();
arguments.forEach((k, v) -> variables.put(k, Double.parseDouble(v)));
expr.setVariables(variables);
double result = expr.evaluate();
this.outputResult = Double.toString(result);
}
this.result = "result";
return null;
}
}

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

@ -18,14 +18,10 @@ package org.thingsboard.server.service.entitiy.cf;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.service.security.model.SecurityUser;
public interface TbCalculatedFieldService {
void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback);
CalculatedField save(CalculatedField calculatedField, SecurityUser user) throws ThingsboardException;
CalculatedField findById(CalculatedFieldId calculatedFieldId, SecurityUser user);

10
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -86,7 +86,7 @@ import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.entitiy.cf.TbCalculatedFieldService;
import org.thingsboard.server.service.cf.CalculatedFieldExecutionService;
import org.thingsboard.server.service.notification.NotificationSchedulerService;
import org.thingsboard.server.service.ota.OtaPackageStateService;
import org.thingsboard.server.service.profile.TbAssetProfileCache;
@ -150,7 +150,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
private final TbImageService imageService;
private final RuleEngineCallService ruleEngineCallService;
private final TbCoreConsumerStats stats;
private final TbCalculatedFieldService calculatedFieldService;
private final CalculatedFieldExecutionService calculatedFieldExecutionService;
private MainQueueConsumerManager<TbProtoQueueMsg<ToCoreMsg>, CoreQueueConfig> mainConsumer;
private QueueConsumerManager<TbProtoQueueMsg<ToUsageStatsServiceMsg>> usageStatsConsumer;
@ -179,7 +179,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
NotificationRuleProcessor notificationRuleProcessor,
TbImageService imageService,
RuleEngineCallService ruleEngineCallService,
TbCalculatedFieldService calculatedFieldService) {
CalculatedFieldExecutionService calculatedFieldExecutionService) {
super(actorContext, tenantProfileCache, deviceProfileCache, assetProfileCache, apiUsageStateService, partitionService,
eventPublisher, jwtSettingsService);
this.stateService = stateService;
@ -195,7 +195,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
this.imageService = imageService;
this.ruleEngineCallService = ruleEngineCallService;
this.queueFactory = tbCoreQueueFactory;
this.calculatedFieldService = calculatedFieldService;
this.calculatedFieldExecutionService = calculatedFieldExecutionService;
}
@PostConstruct
@ -668,7 +668,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
private void forwardToCalculatedFieldService(TransportProtos.CalculatedFieldMsgProto calculatedFieldMsg, TbCallback callback) {
var tenantId = toTenantId(calculatedFieldMsg.getTenantIdMSB(), calculatedFieldMsg.getTenantIdLSB());
var calculatedFieldId = new CalculatedFieldId(new UUID(calculatedFieldMsg.getCalculatedFieldIdMSB(), calculatedFieldMsg.getCalculatedFieldIdLSB()));
ListenableFuture<?> future = deviceActivityEventsExecutor.submit(() -> calculatedFieldService.onCalculatedFieldMsg(calculatedFieldMsg, callback));
ListenableFuture<?> future = deviceActivityEventsExecutor.submit(() -> calculatedFieldExecutionService.onCalculatedFieldMsg(calculatedFieldMsg, callback));
DonAsynchron.withCallback(future,
__ -> callback.onSuccess(),
t -> {

3
application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java

@ -21,6 +21,7 @@ import org.junit.Test;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.cf.Argument;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
@ -136,7 +137,7 @@ public class CalculatedFieldControllerTest extends AbstractControllerTest {
private CalculatedFieldConfiguration getCalculatedFieldConfig(EntityId referencedEntityId) {
SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration();
SimpleCalculatedFieldConfiguration.Argument argument = new SimpleCalculatedFieldConfiguration.Argument();
Argument argument = new Argument();
argument.setEntityId(referencedEntityId);
argument.setType("TIME_SERIES");
argument.setKey("temperature");

2
common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java

@ -60,4 +60,6 @@ public interface CalculatedFieldService extends EntityDaoService {
boolean referencedInAnyCalculatedField(TenantId tenantId, EntityId referencedEntityId);
boolean existsCalculatedFieldByEntityId(TenantId tenantId, EntityId entityId);
}

32
common/data/src/main/java/org/thingsboard/server/common/data/cf/Argument.java

@ -0,0 +1,32 @@
/**
* Copyright © 2016-2024 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.data.cf;
import lombok.Data;
import org.thingsboard.server.common.data.id.EntityId;
@Data
public class Argument {
private EntityId entityId;
private String key;
private String type;
private String defaultValue;
private int limit;
private long timeWindow;
}

8
common/data/src/main/java/org/thingsboard/server/common/data/cf/BaseCalculatedFieldConfiguration.java

@ -105,14 +105,6 @@ public abstract class BaseCalculatedFieldConfiguration implements CalculatedFiel
return configNode;
}
@Data
public static class Argument {
private EntityId entityId;
private String key;
private String type;
private String defaultValue;
}
@Data
public static class Output {
private String name;

2
common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldConfiguration.java

@ -39,7 +39,7 @@ public interface CalculatedFieldConfiguration {
@JsonIgnore
CalculatedFieldType getType();
Map<String, BaseCalculatedFieldConfiguration.Argument> getArguments();
Map<String, Argument> getArguments();
BaseCalculatedFieldConfiguration.Output getOutput();

5
dao/src/main/java/org/thingsboard/server/dao/cf/BaseCalculatedFieldService.java

@ -195,6 +195,11 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements
.anyMatch(referencedEntities -> referencedEntities.contains(referencedEntityId));
}
@Override
public boolean existsCalculatedFieldByEntityId(TenantId tenantId, EntityId entityId) {
return calculatedFieldDao.existsByEntityId(tenantId, entityId);
};
@Override
public Optional<HasId<?>> findEntity(TenantId tenantId, EntityId entityId) {
return Optional.ofNullable(findById(tenantId, new CalculatedFieldId(entityId.getId())));

2
dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java

@ -34,4 +34,6 @@ public interface CalculatedFieldDao extends Dao<CalculatedField> {
List<CalculatedField> removeAllByEntityId(TenantId tenantId, EntityId entityId);
boolean existsByEntityId(TenantId tenantId, EntityId entityId);
}

5
dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java

@ -66,6 +66,11 @@ public class JpaCalculatedFieldDao extends JpaAbstractDao<CalculatedFieldEntity,
return DaoUtil.convertDataList(calculatedFieldRepository.removeAllByTenantIdAndEntityId(tenantId.getId(), entityId.getId()));
}
@Override
public boolean existsByEntityId(TenantId tenantId, EntityId entityId) {
return calculatedFieldRepository.existsByTenantIdAndEntityId(tenantId.getId(), entityId.getId());
}
@Override
protected Class<CalculatedFieldEntity> getEntityClass() {
return CalculatedFieldEntity.class;

3
dao/src/test/java/org/thingsboard/server/dao/service/AssetServiceTest.java

@ -30,6 +30,7 @@ import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.asset.AssetInfo;
import org.thingsboard.server.common.data.asset.AssetProfile;
import org.thingsboard.server.common.data.cf.Argument;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration;
@ -880,7 +881,7 @@ public class AssetServiceTest extends AbstractServiceTest {
SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration();
SimpleCalculatedFieldConfiguration.Argument argument = new SimpleCalculatedFieldConfiguration.Argument();
Argument argument = new Argument();
argument.setEntityId(savedAsset.getId());
argument.setType("TIME_SERIES");
argument.setKey("temperature");

3
dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java

@ -23,6 +23,7 @@ import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.cf.Argument;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
@ -149,7 +150,7 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest {
private CalculatedFieldConfiguration getCalculatedFieldConfig(EntityId referencedEntityId) {
SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration();
SimpleCalculatedFieldConfiguration.Argument argument = new SimpleCalculatedFieldConfiguration.Argument();
Argument argument = new Argument();
argument.setEntityId(referencedEntityId);
argument.setType("TIME_SERIES");
argument.setKey("temperature");

3
dao/src/test/java/org/thingsboard/server/dao/service/CustomerServiceTest.java

@ -31,6 +31,7 @@ import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.cf.Argument;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration;
@ -375,7 +376,7 @@ public class CustomerServiceTest extends AbstractServiceTest {
SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration();
SimpleCalculatedFieldConfiguration.Argument argument = new SimpleCalculatedFieldConfiguration.Argument();
Argument argument = new Argument();
argument.setEntityId(savedCustomer.getId());
argument.setType("TIME_SERIES");
argument.setKey("temperature");

3
dao/src/test/java/org/thingsboard/server/dao/service/DeviceServiceTest.java

@ -39,6 +39,7 @@ import org.thingsboard.server.common.data.OtaPackageInfo;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.cf.Argument;
import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.SimpleCalculatedFieldConfiguration;
@ -1218,7 +1219,7 @@ public class DeviceServiceTest extends AbstractServiceTest {
SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration();
SimpleCalculatedFieldConfiguration.Argument argument = new SimpleCalculatedFieldConfiguration.Argument();
Argument argument = new Argument();
argument.setEntityId(device.getId());
argument.setType("TIME_SERIES");
argument.setKey("temperature");

Loading…
Cancel
Save