From 31103e90e3e3cb5f440823872f2a0c85cdac6b10 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Fri, 22 Nov 2024 17:17:59 +0200 Subject: [PATCH] added implementation to handle update entity profiles events --- .../cf/CalculatedFieldExecutionService.java | 2 + ...efaultCalculatedFieldExecutionService.java | 74 ++++++++++++++++--- .../CalculatedFieldTbelScriptEngine.java | 14 +--- .../ctx/state/ScriptCalculatedFieldState.java | 13 +--- .../entitiy/EntityStateSourcingListener.java | 31 ++++---- .../queue/DefaultTbClusterService.java | 44 ++++++++++- .../queue/DefaultTbCoreConsumerService.java | 19 ++++- .../server/cluster/TbClusterService.java | 3 + .../server/dao/cf/CalculatedFieldService.java | 2 + common/proto/src/main/proto/queue.proto | 34 ++++++--- .../server/dao/asset/BaseAssetService.java | 4 +- .../dao/cf/BaseCalculatedFieldService.java | 14 +++- .../server/dao/cf/CalculatedFieldDao.java | 3 + .../dao/sql/cf/CalculatedFieldRepository.java | 3 + .../dao/sql/cf/JpaCalculatedFieldDao.java | 6 ++ 15 files changed, 202 insertions(+), 64 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java index d4e0d1da84..6b7b12655f 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java @@ -29,4 +29,6 @@ public interface CalculatedFieldExecutionService { void onTelemetryUpdate(TenantId tenantId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry); + void onEntityTypeChanged(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback); + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java index 6aa9997d1f..9edea50a05 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java @@ -27,12 +27,14 @@ import jakarta.annotation.PreDestroy; import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.math.NumberUtils; 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.script.api.tbel.TbelInvokeService; import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldLink; import org.thingsboard.server.common.data.cf.CalculatedFieldType; @@ -44,8 +46,14 @@ 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.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.BooleanDataEntry; +import org.thingsboard.server.common.data.kv.DoubleDataEntry; import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.page.PageDataIterable; import org.thingsboard.server.common.msg.TbMsg; @@ -210,7 +218,38 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas updateOrInitializeState(calculatedField, calculatedField.getEntityId(), updatedTelemetry); log.info("Successfully updated time series for calculatedFieldId: [{}]", calculatedFieldId); } catch (Exception e) { - log.trace("Failed to update time series for calculatedFieldId: [{}]", calculatedFieldId, e); + log.trace("Failed to update telemetry for calculatedFieldId: [{}]", calculatedFieldId, e); + } + } + + @Override + public void onEntityTypeChanged(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback) { + try { + TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); + EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); + EntityId oldProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getOldProfileIdMSB(), proto.getOldProfileIdLSB())); + EntityId newProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getNewProfileIdMSB(), proto.getNewProfileIdLSB())); + + log.info("Received EntityProfileUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); + + List cfIdsOfOldProfile = calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, oldProfileId); + cfIdsOfOldProfile.forEach(id -> states.remove(new CalculatedFieldCtxId(id.getId(), entityId.getId()))); + List ctxIdsToDelete = cfIdsOfOldProfile.stream().map(cfId -> JacksonUtil.writeValueAsString(new CalculatedFieldCtxId(cfId.getId(), entityId.getId()))).toList(); + rocksDBService.deleteAll(ctxIdsToDelete); + + calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, oldProfileId) + .forEach(cfId -> { + CalculatedFieldCtxId ctxId = new CalculatedFieldCtxId(cfId.getId(), entityId.getId()); + states.remove(ctxId); + rocksDBService.delete(JacksonUtil.writeValueAsString(ctxId)); + }); + + calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, newProfileId) + .stream() + .map(cfId -> calculatedFields.computeIfAbsent(cfId, id -> calculatedFieldService.findById(tenantId, id))) + .forEach(cf -> initializeStateForEntity(tenantId, cf, entityId, callback)); + } catch (Exception e) { + log.trace("Failed to process entity type update msg: [{}]", proto, e); } } @@ -271,10 +310,9 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas Map arguments = calculatedField.getConfiguration().getArguments(); Map argumentValues = new HashMap<>(); AtomicInteger remaining = new AtomicInteger(arguments.size()); - arguments.forEach((key, argument) -> Futures.addCallback(fetchArgumentValue(tenantId, argument), new FutureCallback<>() { + arguments.forEach((key, argument) -> Futures.addCallback(fetchArgumentValue(tenantId, argument, entityId), new FutureCallback<>() { @Override public void onSuccess(Optional result) { - // todo: should be rewritten implementation for default value argumentValues.put(key, result.orElse(null)); if (remaining.decrementAndGet() == 0) { updateOrInitializeState(calculatedField, entityId, argumentValues); @@ -289,20 +327,38 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas }, calculatedFieldCallbackExecutor)); } - private ListenableFuture> fetchArgumentValue(TenantId tenantId, Argument argument) { + private ListenableFuture> fetchArgumentValue(TenantId tenantId, Argument argument, EntityId targetEntityId) { + EntityId argumentEntityId = argument.getEntityId(); + EntityId entityId = EntityType.DEVICE_PROFILE.equals(argumentEntityId.getEntityType()) || EntityType.ASSET_PROFILE.equals(argumentEntityId.getEntityType()) ? targetEntityId : argumentEntityId; return switch (argument.getType()) { case "ATTRIBUTES" -> Futures.transform( - attributesService.find(tenantId, argument.getEntityId(), argument.getScope(), argument.getKey()), - result -> result.map(entry -> (KvEntry) entry), + attributesService.find(tenantId, entityId, argument.getScope(), argument.getKey()), + result -> result.or(() -> Optional.of( + new BaseAttributeKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument)) + )), MoreExecutors.directExecutor()); case "TIME_SERIES" -> Futures.transform( - timeseriesService.findLatest(tenantId, argument.getEntityId(), argument.getKey()), - result -> result.map(entry -> (KvEntry) entry), + timeseriesService.findLatest(tenantId, entityId, argument.getKey()), + result -> result.or(() -> Optional.of( + new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument)) + )), MoreExecutors.directExecutor()); default -> throw new IllegalArgumentException("Invalid argument type '" + argument.getType() + "'."); }; } + private KvEntry createDefaultKvEntry(Argument argument) { + String key = argument.getKey(); + String defaultValue = argument.getDefaultValue(); + if (NumberUtils.isParsable(defaultValue)) { + return new DoubleDataEntry(key, Double.parseDouble(defaultValue)); + } + if ("true".equalsIgnoreCase(defaultValue) || "false".equalsIgnoreCase(defaultValue)) { + return new BooleanDataEntry(key, Boolean.parseBoolean(defaultValue)); + } + return new StringDataEntry(key, defaultValue); + } + private void updateOrInitializeState(CalculatedField calculatedField, EntityId entityId, Map argumentValues) { CalculatedFieldCtxId ctxId = new CalculatedFieldCtxId(calculatedField.getUuidId(), entityId.getId()); CalculatedFieldCtx calculatedFieldCtx = states.computeIfAbsent(ctxId, ctx -> new CalculatedFieldCtx(ctxId, null)); @@ -322,7 +378,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas @Override public void onSuccess(CalculatedFieldResult result) { if (result != null) { - pushMsgToRuleEngine(calculatedField.getTenantId(), calculatedField.getEntityId(), result); + pushMsgToRuleEngine(calculatedField.getTenantId(), entityId, result); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldTbelScriptEngine.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldTbelScriptEngine.java index 7e8376be8e..5cc58b2a95 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldTbelScriptEngine.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldTbelScriptEngine.java @@ -19,7 +19,6 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; import lombok.extern.slf4j.Slf4j; -import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.script.api.ScriptType; import org.thingsboard.script.api.tbel.TbelInvokeService; import org.thingsboard.server.common.data.id.TenantId; @@ -55,18 +54,7 @@ public class CalculatedFieldTbelScriptEngine implements CalculatedFieldScriptEng @Override public ListenableFuture executeScriptAsync(Map arguments) { log.trace("execute script async, arguments {}", arguments); - Object[] args = new Object[arguments.size()]; - int index = 0; - for (KvEntry entry : arguments.values()) { - switch (entry.getDataType()) { - case BOOLEAN -> args[index] = entry.getBooleanValue().orElse(null); - case DOUBLE -> args[index] = entry.getDoubleValue().orElse(null); - case LONG -> args[index] = entry.getLongValue().orElse(null); - case JSON -> args[index] = entry.getJsonValue().map(JacksonUtil::toJsonNode).orElse(null); - default -> args[index] = entry.getValueAsString(); - } - index++; - } + Object[] args = arguments.values().stream().map(KvEntry::getValue).toArray(); return Futures.transformAsync(tbelInvokeService.invokeScript(tenantId, null, this.scriptId, args), o -> { try { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java index cbee8dc902..5c984a8d16 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java @@ -72,16 +72,9 @@ public class ScriptCalculatedFieldState implements CalculatedFieldState { return Futures.transform(resultFuture, result -> { Output output = calculatedFieldConfiguration.getOutput(); - Map resultMap = new HashMap<>(); - - if (result instanceof Map) { - Map map = JacksonUtil.convertValue(result, Map.class); - if (map != null) { - resultMap.putAll(map); - } - } else { - resultMap.put(output.getName(), JacksonUtil.convertValue(result, Object.class)); - } + Map resultMap = result instanceof Map + ? JacksonUtil.convertValue(result, Map.class) + : new HashMap<>(); CalculatedFieldResult calculatedFieldResult = new CalculatedFieldResult(); calculatedFieldResult.setType(output.getType()); diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java index 1b541bcd5d..154f5f4833 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/EntityStateSourcingListener.java @@ -31,6 +31,7 @@ import org.thingsboard.server.common.data.TbResource; import org.thingsboard.server.common.data.TbResourceInfo; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; +import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.edge.Edge; @@ -51,8 +52,6 @@ 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; @@ -67,8 +66,6 @@ public class EntityStateSourcingListener { private final TbClusterService tbClusterService; private final TenantService tenantService; - private final CalculatedFieldService calculatedFieldService; - private final DeviceProfileService deviceProfileService; @PostConstruct public void init() { @@ -88,7 +85,10 @@ public class EntityStateSourcingListener { ComponentLifecycleEvent lifecycleEvent = isCreated ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED; switch (entityType) { - case ASSET, ASSET_PROFILE, ENTITY_VIEW, NOTIFICATION_RULE -> { + case ASSET -> { + onAssetUpdate(event.getEntity(), event.getOldEntity()); + } + case ASSET_PROFILE, ENTITY_VIEW, NOTIFICATION_RULE -> { tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, lifecycleEvent); } case RULE_CHAIN -> { @@ -106,7 +106,7 @@ public class EntityStateSourcingListener { onTenantProfileUpdate(tenantProfile, lifecycleEvent); } case DEVICE -> { - onDeviceUpdate(tenantId, event.getEntity(), event.getOldEntity()); + onDeviceUpdate(event.getEntity(), event.getOldEntity()); } case DEVICE_PROFILE -> { DeviceProfile deviceProfile = (DeviceProfile) event.getEntity(); @@ -245,23 +245,24 @@ public class EntityStateSourcingListener { tbClusterService.broadcastEntityStateChangeEvent(tenantId, entityId, ComponentLifecycleEvent.DELETED); } - private void onDeviceUpdate(TenantId tenantId, Object entity, Object oldEntity) { + private void onDeviceUpdate(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); } + private void onAssetUpdate(Object entity, Object oldEntity) { + Asset asset = (Asset) entity; + Asset oldAsset = null; + if (oldEntity instanceof Asset) { + oldAsset = (Asset) oldEntity; + } + tbClusterService.onAssetUpdated(asset, oldAsset); + } + private void onEdgeEvent(TenantId tenantId, EntityId entityId, Object entity, ComponentLifecycleEvent lifecycleEvent) { if (entity instanceof Edge) { tbClusterService.onEdgeStateChangeEvent(new ComponentLifecycleMsg(tenantId, entityId, lifecycleEvent)); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java index f7c8230716..58a16d9c0a 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java @@ -68,6 +68,7 @@ import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.common.util.ProtoUtils; +import org.thingsboard.server.dao.cf.CalculatedFieldService; import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ComponentLifecycleMsgProto; @@ -149,6 +150,7 @@ public class DefaultTbClusterService implements TbClusterService { private final GatewayNotificationsService gatewayNotificationsService; private final EdgeService edgeService; private final TbTransactionalCache edgeIdServiceIdCache; + private final CalculatedFieldService calculatedFieldService; @Override public void pushMsgToCore(TenantId tenantId, EntityId entityId, ToCoreMsg msg, TbQueueCallback callback) { @@ -609,7 +611,11 @@ public class DefaultTbClusterService implements TbClusterService { if (deviceNameChanged) { gatewayNotificationsService.onDeviceUpdated(device, old); } - if (deviceNameChanged || !device.getType().equals(old.getType())) { + boolean deviceTypeChanged = !device.getType().equals(old.getType()); + if (deviceTypeChanged) { + handleProfileChange(device.getTenantId(), device.getId(), old.getDeviceProfileId(), device.getDeviceProfileId()); + } + if (deviceNameChanged || deviceTypeChanged) { pushMsgToCore(new DeviceNameOrTypeUpdateMsg(device.getTenantId(), device.getId(), device.getName(), device.getType()), null); } } @@ -618,6 +624,26 @@ public class DefaultTbClusterService implements TbClusterService { otaPackageStateService.update(device, old); } + @Override + public void onAssetUpdated(Asset asset, Asset old) { + var created = old == null; + broadcastEntityChangeToTransport(asset.getTenantId(), asset.getId(), asset, null); + if (old != null) { + boolean assetTypeChanged = !asset.getType().equals(old.getType()); + if (assetTypeChanged) { + handleProfileChange(asset.getTenantId(), asset.getId(), old.getAssetProfileId(), asset.getAssetProfileId()); + } + } + broadcastEntityStateChangeEvent(asset.getTenantId(), asset.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED); + } + + private void handleProfileChange(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) { + boolean cfExistsByProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, oldProfileId); + if (cfExistsByProfile) { + sendEntityTypeUpdatedEvent(tenantId, entityId, oldProfileId, newProfileId); + } + } + @Override public void sendNotificationMsgToEdge(TenantId tenantId, EdgeId edgeId, EntityId entityId, String body, EdgeEventType type, EdgeEventActionType action, EdgeId originatorEdgeId) { if (!edgesEnabled) { @@ -775,4 +801,20 @@ public class DefaultTbClusterService implements TbClusterService { pushMsgToCore(tenantId, calculatedFieldId, ToCoreMsg.newBuilder().setCalculatedFieldMsg(msg).build(), null); } + private void sendEntityTypeUpdatedEvent(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) { + TransportProtos.EntityProfileUpdateMsgProto.Builder builder = TransportProtos.EntityProfileUpdateMsgProto.newBuilder(); + builder.setTenantIdMSB(tenantId.getId().getMostSignificantBits()); + builder.setTenantIdLSB(tenantId.getId().getLeastSignificantBits()); + builder.setEntityType(entityId.getEntityType().name()); + builder.setEntityIdMSB(entityId.getId().getMostSignificantBits()); + builder.setEntityIdLSB(entityId.getId().getLeastSignificantBits()); + builder.setEntityProfileType(newProfileId.getEntityType().name()); + builder.setOldProfileIdMSB(oldProfileId.getId().getMostSignificantBits()); + builder.setOldProfileIdLSB(oldProfileId.getId().getLeastSignificantBits()); + builder.setNewProfileIdMSB(newProfileId.getId().getMostSignificantBits()); + builder.setNewProfileIdLSB(newProfileId.getId().getLeastSignificantBits()); + TransportProtos.EntityProfileUpdateMsgProto msg = builder.build(); + pushMsgToCore(tenantId, entityId, ToCoreMsg.newBuilder().setEntityProfileUpdateMsg(msg).build(), null); + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index 2d3323e097..3caaab6613 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -40,6 +40,7 @@ import org.thingsboard.server.common.data.event.Event; import org.thingsboard.server.common.data.event.LifecycleEvent; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; @@ -157,6 +158,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService> firmwareStatesConsumer; private volatile ListeningExecutorService deviceActivityEventsExecutor; + private volatile ListeningExecutorService calculatedFieldsExecutor; public DefaultTbCoreConsumerService(TbCoreQueueFactory tbCoreQueueFactory, ActorSystemContext actorContext, @@ -202,6 +204,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService, CoreQueueConfig>builder() .queueKey(new QueueKey(ServiceType.TB_CORE)) @@ -315,6 +318,8 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService future = deviceActivityEventsExecutor.submit(() -> calculatedFieldExecutionService.onCalculatedFieldMsg(calculatedFieldMsg, callback)); + ListenableFuture future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onCalculatedFieldMsg(calculatedFieldMsg, callback)); DonAsynchron.withCallback(future, __ -> callback.onSuccess(), t -> { @@ -677,6 +682,18 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onEntityTypeChanged(profileUpdateMsg, callback)); + DonAsynchron.withCallback(future, + __ -> callback.onSuccess(), + t -> { + log.warn("[{}] Failed to process device type updated message for device [{}]", tenantId.getId(), entityId.getId(), t); + callback.onFailure(t); + }); + } + private void forwardToNotificationSchedulerService(TransportProtos.NotificationSchedulerServiceMsg msg, TbCallback callback) { TenantId tenantId = toTenantId(msg.getTenantIdMSB(), msg.getTenantIdLSB()); NotificationRequestId notificationRequestId = new NotificationRequestId(new UUID(msg.getRequestIdMSB(), msg.getRequestIdLSB())); diff --git a/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java b/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java index f173005107..131b50a52c 100644 --- a/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java +++ b/common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java @@ -21,6 +21,7 @@ import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.TbResourceInfo; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantProfile; +import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventType; @@ -97,6 +98,8 @@ public interface TbClusterService extends TbQueueClusterService { void onDeviceAssignedToTenant(TenantId oldTenantId, Device device); + void onAssetUpdated(Asset asset, Asset old); + void onResourceChange(TbResourceInfo resource, TbQueueCallback callback); void onResourceDeleted(TbResourceInfo resource, TbQueueCallback callback); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java index e44ff0ba22..1e64fdac60 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldService.java @@ -36,6 +36,8 @@ public interface CalculatedFieldService extends EntityDaoService { ListenableFuture findCalculatedFieldByIdAsync(TenantId tenantId, CalculatedFieldId calculatedFieldId); + List findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId); + List findAllCalculatedFields(); PageData findAllCalculatedFields(PageLink pageLink); diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 95e4fa8601..bf49a68c80 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -771,6 +771,29 @@ message DeviceInactivityProto { int64 lastInactivityTime = 5; } +message CalculatedFieldMsgProto { + int64 tenantIdMSB = 1; + int64 tenantIdLSB = 2; + int64 calculatedFieldIdMSB = 3; + int64 calculatedFieldIdLSB = 4; + bool added = 5; + bool updated = 6; + bool deleted = 7; +} + +message EntityProfileUpdateMsgProto { + int64 tenantIdMSB = 1; + int64 tenantIdLSB = 2; + string entityType = 3; + int64 entityIdMSB = 4; + int64 entityIdLSB = 5; + string entityProfileType = 6; + int64 oldProfileIdMSB = 7; + int64 oldProfileIdLSB = 8; + int64 newProfileIdMSB = 9; + int64 newProfileIdLSB = 10; +} + //Used to report session state to tb-Service and persist this state in the cache on the tb-Service level. message SubscriptionInfoProto { int64 lastActivityTime = 1; @@ -1267,16 +1290,6 @@ message ToDeviceActorNotificationMsgProto { DeviceDeleteMsgProto deviceDeleteMsg = 8; } -message CalculatedFieldMsgProto { - int64 tenantIdMSB = 1; - int64 tenantIdLSB = 2; - int64 calculatedFieldIdMSB = 3; - int64 calculatedFieldIdLSB = 4; - bool added = 5; - bool updated = 6; - bool deleted = 7; -} - /** TB Core to Version Control Service */ @@ -1513,6 +1526,7 @@ message ToCoreMsg { DeviceDisconnectProto deviceDisconnectMsg = 51; DeviceInactivityProto deviceInactivityMsg = 52; CalculatedFieldMsgProto calculatedFieldMsg = 53; + EntityProfileUpdateMsgProto entityProfileUpdateMsg = 54; } /* High priority messages with low latency are handled by ThingsBoard Core Service separately */ diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java index 2135792174..7242cfde68 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java @@ -177,8 +177,8 @@ public class BaseAssetService extends AbstractCachedEntityService findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId) { + log.trace("Executing findCalculatedFieldIdsByEntityId [{}]", entityId); + validateId(entityId.getId(), id -> INCORRECT_ENTITY_ID + id); + return calculatedFieldDao.findCalculatedFieldIdsByEntityId(tenantId, entityId); + } + @Override public List findAllCalculatedFields() { log.trace("Executing findAll"); @@ -133,7 +141,7 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements public int deleteAllCalculatedFieldsByEntityId(TenantId tenantId, EntityId entityId) { log.trace("Executing deleteAllCalculatedFieldsByEntityId, tenantId [{}], entityId [{}]", tenantId, entityId); validateId(tenantId, id -> INCORRECT_TENANT_ID + id); - validateId(entityId.getId(), id -> "Incorrect entityId " + id); + validateId(entityId.getId(), id -> INCORRECT_ENTITY_ID + id); List calculatedFields = calculatedFieldDao.removeAllByEntityId(tenantId, entityId); return calculatedFields.size(); } @@ -212,7 +220,7 @@ public class BaseCalculatedFieldService extends AbstractEntityService implements @Override public boolean existsCalculatedFieldByEntityId(TenantId tenantId, EntityId entityId) { return calculatedFieldDao.existsByEntityId(tenantId, entityId); - }; + } @Override public Optional> findEntity(TenantId tenantId, EntityId entityId) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java b/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java index a6b7c2dea1..5b3bcc2750 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/cf/CalculatedFieldDao.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.cf; import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; @@ -28,6 +29,8 @@ public interface CalculatedFieldDao extends Dao { List findAllByTenantId(TenantId tenantId); + List findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId); + List findAll(); PageData findAll(PageLink pageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java index 333057e8c5..9aa0aee428 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/CalculatedFieldRepository.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.sql.cf; import org.springframework.data.jpa.repository.JpaRepository; +import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.dao.model.sql.CalculatedFieldEntity; import java.util.List; @@ -25,6 +26,8 @@ public interface CalculatedFieldRepository extends JpaRepository findCalculatedFieldIdsByTenantIdAndEntityId(UUID tenantId, UUID entityId); + List findAllByTenantId(UUID tenantId); List removeAllByTenantIdAndEntityId(UUID tenantId, UUID entityId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java index 737a089a15..e3762f6157 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/cf/JpaCalculatedFieldDao.java @@ -22,6 +22,7 @@ import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; @@ -49,6 +50,11 @@ public class JpaCalculatedFieldDao extends JpaAbstractDao findCalculatedFieldIdsByEntityId(TenantId tenantId, EntityId entityId) { + return calculatedFieldRepository.findCalculatedFieldIdsByTenantIdAndEntityId(tenantId.getId(), entityId.getId()); + } + @Override public List findAll() { return DaoUtil.convertDataList(calculatedFieldRepository.findAll());