Browse Source

implemented logic to handle profile updates even if no cf exists by them

pull/12266/head
IrynaMatveieva 2 years ago
parent
commit
8f87caba14
  1. 11
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java
  2. 36
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

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

@ -542,6 +542,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
private void updateOrInitializeState(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map<String, ArgumentEntry> argumentValues, List<CalculatedFieldId> calculatedFieldIds) {
TenantId tenantId = calculatedFieldCtx.getTenantId();
CalculatedFieldId cfId = calculatedFieldCtx.getCfId();
Map<String, ArgumentEntry> argumentsMap = new HashMap<>(argumentValues);
TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, cfId);
if (tpi.isMyPartition()) {
CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId.getId(), entityId.getId());
@ -550,7 +551,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
CalculatedFieldEntityCtx calculatedFieldEntityCtx = ctx != null ? ctx : fetchCalculatedFieldEntityState(ctxId, calculatedFieldCtx.getCfType());
Consumer<CalculatedFieldState> performUpdateState = (state) -> {
if (state.updateState(argumentValues)) {
if (state.updateState(argumentsMap)) {
calculatedFieldEntityCtx.setState(state);
rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx));
Map<String, ArgumentEntry> arguments = state.getArguments();
@ -564,17 +565,17 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
CalculatedFieldState state = calculatedFieldEntityCtx.getState();
boolean allKeysPresent = argumentValues.keySet().containsAll(calculatedFieldCtx.getArguments().keySet());
boolean allKeysPresent = argumentsMap.keySet().containsAll(calculatedFieldCtx.getArguments().keySet());
boolean requiresTsRollingUpdate = calculatedFieldCtx.getArguments().values().stream()
.anyMatch(argument -> ArgumentType.TS_ROLLING.equals(argument.getType()) && state.getArguments().get(argument.getKey()) == null);
if (!allKeysPresent || requiresTsRollingUpdate) {
Map<String, Argument> missingArguments = calculatedFieldCtx.getArguments().entrySet().stream()
.filter(entry -> !argumentValues.containsKey(entry.getKey()) || (ArgumentType.TS_ROLLING.equals(entry.getValue().getType()) && state.getArguments().get(entry.getKey()) == null))
.filter(entry -> !argumentsMap.containsKey(entry.getKey()) || (ArgumentType.TS_ROLLING.equals(entry.getValue().getType()) && state.getArguments().get(entry.getKey()) == null))
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
fetchArguments(calculatedFieldCtx.getTenantId(), entityId, missingArguments, argumentValues::putAll)
fetchArguments(calculatedFieldCtx.getTenantId(), entityId, missingArguments, argumentsMap::putAll)
.addListener(() -> performUpdateState.accept(state),
calculatedFieldCallbackExecutor);
} else {
@ -583,7 +584,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
return calculatedFieldEntityCtx;
});
} else {
sendUpdateCalculatedFieldStateMsg(tenantId, cfId, entityId, calculatedFieldIds, argumentValues);
sendUpdateCalculatedFieldStateMsg(tenantId, cfId, entityId, calculatedFieldIds, argumentsMap);
}
}

36
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

@ -68,7 +68,6 @@ 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;
@ -150,7 +149,6 @@ public class DefaultTbClusterService implements TbClusterService {
private final GatewayNotificationsService gatewayNotificationsService;
private final EdgeService edgeService;
private final TbTransactionalCache<EdgeId, String> edgeIdServiceIdCache;
private final CalculatedFieldService calculatedFieldService;
@Override
public void pushMsgToCore(TenantId tenantId, EntityId entityId, ToCoreMsg msg, TbQueueCallback callback) {
@ -393,7 +391,7 @@ public class DefaultTbClusterService implements TbClusterService {
public void onDeviceDeleted(TenantId tenantId, Device device, TbQueueCallback callback) {
DeviceId deviceId = device.getId();
gatewayNotificationsService.onDeviceDeleted(device);
handleEntityDelete(tenantId, deviceId, device.getDeviceProfileId());
sendProfileEntityEvent(tenantId, deviceId, device.getDeviceProfileId(), false, true);
broadcastEntityDeleteToTransport(tenantId, deviceId, device.getName(), callback);
sendDeviceStateServiceEvent(tenantId, deviceId, false, false, true);
broadcastEntityStateChangeEvent(tenantId, deviceId, ComponentLifecycleEvent.DELETED);
@ -402,17 +400,10 @@ public class DefaultTbClusterService implements TbClusterService {
@Override
public void onAssetDeleted(TenantId tenantId, Asset asset, TbQueueCallback callback) {
AssetId assetId = asset.getId();
handleEntityDelete(tenantId, assetId, asset.getAssetProfileId());
sendProfileEntityEvent(tenantId, assetId, asset.getAssetProfileId(), false, true);
broadcastEntityStateChangeEvent(tenantId, assetId, ComponentLifecycleEvent.DELETED);
}
private void handleEntityDelete(TenantId tenantId, EntityId entityId, EntityId profileId) {
boolean cfExistsByProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, profileId);
if (cfExistsByProfile) {
sendProfileEntityEvent(tenantId, entityId, profileId, false, true);
}
}
@Override
public void onDeviceAssignedToTenant(TenantId oldTenantId, Device device) {
onDeviceDeleted(oldTenantId, device, null);
@ -633,13 +624,13 @@ public class DefaultTbClusterService implements TbClusterService {
}
boolean deviceTypeChanged = !device.getType().equals(old.getType());
if (deviceTypeChanged) {
handleProfileUpdate(device.getTenantId(), device.getId(), old.getDeviceProfileId(), device.getDeviceProfileId());
sendEntityProfileUpdatedEvent(device.getTenantId(), device.getId(), old.getDeviceProfileId(), device.getDeviceProfileId());
}
if (deviceNameChanged || deviceTypeChanged) {
pushMsgToCore(new DeviceNameOrTypeUpdateMsg(device.getTenantId(), device.getId(), device.getName(), device.getType()), null);
}
} else {
handleEntityCreate(device.getTenantId(), device.getId(), device.getDeviceProfileId());
sendProfileEntityEvent(device.getTenantId(), device.getId(), device.getDeviceProfileId(), true, false);
}
broadcastEntityStateChangeEvent(device.getTenantId(), device.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED);
sendDeviceStateServiceEvent(device.getTenantId(), device.getId(), created, !created, false);
@ -653,29 +644,14 @@ public class DefaultTbClusterService implements TbClusterService {
if (old != null) {
boolean assetTypeChanged = !asset.getType().equals(old.getType());
if (assetTypeChanged) {
handleProfileUpdate(asset.getTenantId(), asset.getId(), old.getAssetProfileId(), asset.getAssetProfileId());
sendEntityProfileUpdatedEvent(asset.getTenantId(), asset.getId(), old.getAssetProfileId(), asset.getAssetProfileId());
}
} else {
handleEntityCreate(asset.getTenantId(), asset.getId(), asset.getAssetProfileId());
sendProfileEntityEvent(asset.getTenantId(), asset.getId(), asset.getAssetProfileId(), true, false);
}
broadcastEntityStateChangeEvent(asset.getTenantId(), asset.getId(), created ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED);
}
private void handleProfileUpdate(TenantId tenantId, EntityId entityId, EntityId oldProfileId, EntityId newProfileId) {
boolean cfExistsByOldProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, oldProfileId);
boolean cfExistsByNewProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, newProfileId);
if (cfExistsByOldProfile || cfExistsByNewProfile) {
sendEntityProfileUpdatedEvent(tenantId, entityId, oldProfileId, newProfileId);
}
}
private void handleEntityCreate(TenantId tenantId, EntityId entityId, EntityId profileId) {
boolean cfExistsByProfile = calculatedFieldService.existsCalculatedFieldByEntityId(tenantId, profileId);
if (cfExistsByProfile) {
sendProfileEntityEvent(tenantId, entityId, profileId, true, false);
}
}
@Override
public void sendNotificationMsgToEdge(TenantId tenantId, EdgeId edgeId, EntityId entityId, String body, EdgeEventType type, EdgeEventActionType action, EdgeId originatorEdgeId) {
if (!edgesEnabled) {

Loading…
Cancel
Save