|
|
@ -36,6 +36,7 @@ import org.thingsboard.server.common.data.id.EntityId; |
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
import org.thingsboard.server.common.data.kv.StringDataEntry; |
|
|
import org.thingsboard.server.common.data.kv.StringDataEntry; |
|
|
import org.thingsboard.server.common.data.msg.TbMsgType; |
|
|
import org.thingsboard.server.common.data.msg.TbMsgType; |
|
|
|
|
|
import org.thingsboard.server.common.data.util.CollectionsUtil; |
|
|
import org.thingsboard.server.common.msg.CalculatedFieldStatePartitionRestoreMsg; |
|
|
import org.thingsboard.server.common.msg.CalculatedFieldStatePartitionRestoreMsg; |
|
|
import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg; |
|
|
import org.thingsboard.server.common.msg.cf.CalculatedFieldPartitionChangeMsg; |
|
|
import org.thingsboard.server.common.msg.queue.ServiceType; |
|
|
import org.thingsboard.server.common.msg.queue.ServiceType; |
|
|
@ -56,6 +57,8 @@ import org.thingsboard.server.service.cf.ctx.state.aggregation.RelatedEntitiesAg |
|
|
import org.thingsboard.server.service.cf.ctx.state.alarm.AlarmCalculatedFieldState; |
|
|
import org.thingsboard.server.service.cf.ctx.state.alarm.AlarmCalculatedFieldState; |
|
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry; |
|
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry; |
|
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; |
|
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; |
|
|
|
|
|
import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationArgumentEntry; |
|
|
|
|
|
import org.thingsboard.server.service.cf.ctx.state.propagation.PropagationCalculatedFieldState; |
|
|
|
|
|
|
|
|
import java.util.ArrayList; |
|
|
import java.util.ArrayList; |
|
|
import java.util.Collection; |
|
|
import java.util.Collection; |
|
|
@ -71,6 +74,7 @@ import java.util.concurrent.TimeUnit; |
|
|
import java.util.stream.Collectors; |
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
|
import static org.thingsboard.server.common.data.DataConstants.REEVALUATION_MSG; |
|
|
import static org.thingsboard.server.common.data.DataConstants.REEVALUATION_MSG; |
|
|
|
|
|
import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; |
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createStateByType; |
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createStateByType; |
|
|
|
|
|
|
|
|
/** |
|
|
/** |
|
|
@ -226,10 +230,9 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM |
|
|
|
|
|
|
|
|
private void handleRelationUpdate(CalculatedFieldRelationActionMsg msg) throws CalculatedFieldException { |
|
|
private void handleRelationUpdate(CalculatedFieldRelationActionMsg msg) throws CalculatedFieldException { |
|
|
CalculatedFieldCtx ctx = msg.getCalculatedField(); |
|
|
CalculatedFieldCtx ctx = msg.getCalculatedField(); |
|
|
var callback = new MultipleTbCallback(CALLBACKS_PER_CF, msg.getCallback()); |
|
|
|
|
|
var state = states.get(ctx.getCfId()); |
|
|
var state = states.get(ctx.getCfId()); |
|
|
try { |
|
|
try { |
|
|
Map<String, ArgumentEntry> updatedArgs = new HashMap<>(); |
|
|
Map<String, ArgumentEntry> updatedArgs = null; |
|
|
if (state == null) { |
|
|
if (state == null) { |
|
|
state = createState(ctx); |
|
|
state = createState(ctx); |
|
|
} else { |
|
|
} else { |
|
|
@ -237,11 +240,19 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM |
|
|
Map<String, ArgumentEntry> fetchedArgs = cfService.fetchArgsFromDb(tenantId, msg.getRelatedEntityId(), ctx.getArguments()); |
|
|
Map<String, ArgumentEntry> fetchedArgs = cfService.fetchArgsFromDb(tenantId, msg.getRelatedEntityId(), ctx.getArguments()); |
|
|
updatedArgs = relatedEntitiesAggState.updateEntityData(setEntityIdToSingleEntityArguments(msg.getRelatedEntityId(), fetchedArgs)); |
|
|
updatedArgs = relatedEntitiesAggState.updateEntityData(setEntityIdToSingleEntityArguments(msg.getRelatedEntityId(), fetchedArgs)); |
|
|
} |
|
|
} |
|
|
|
|
|
if (state instanceof PropagationCalculatedFieldState propagationState) { |
|
|
|
|
|
PropagationArgumentEntry entry = new PropagationArgumentEntry(); |
|
|
|
|
|
entry.setAdded(msg.getRelatedEntityId()); |
|
|
|
|
|
updatedArgs = propagationState.update(Map.of(PROPAGATION_CONFIG_ARGUMENT, entry), ctx); |
|
|
|
|
|
} |
|
|
|
|
|
if (CollectionsUtil.isEmpty(updatedArgs)) { |
|
|
|
|
|
msg.getCallback().onSuccess(); |
|
|
|
|
|
return; |
|
|
|
|
|
} |
|
|
state.checkStateSize(new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId), ctx.getMaxStateSize()); |
|
|
state.checkStateSize(new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId), ctx.getMaxStateSize()); |
|
|
} |
|
|
} |
|
|
if (state.isSizeOk()) { |
|
|
if (state.isSizeOk()) { |
|
|
processStateIfReady(state, updatedArgs, ctx, Collections.singletonList(ctx.getCfId()), null, null, callback); |
|
|
processStateIfReady(state, updatedArgs, ctx, Collections.singletonList(ctx.getCfId()), null, null, msg.getCallback()); |
|
|
} else { |
|
|
} else { |
|
|
throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).errorMessage(ctx.getSizeExceedsLimitMessage()).build(); |
|
|
throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).errorMessage(ctx.getSizeExceedsLimitMessage()).build(); |
|
|
} |
|
|
} |
|
|
@ -272,9 +283,14 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM |
|
|
} else { |
|
|
} else { |
|
|
throw new RuntimeException(ctx.getSizeExceedsLimitMessage()); |
|
|
throw new RuntimeException(ctx.getSizeExceedsLimitMessage()); |
|
|
} |
|
|
} |
|
|
} else { |
|
|
return; |
|
|
msg.getCallback().onSuccess(); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
if (state instanceof PropagationCalculatedFieldState propagationState) { |
|
|
|
|
|
PropagationArgumentEntry entry = new PropagationArgumentEntry(); |
|
|
|
|
|
entry.setRemoved(msg.getRelatedEntityId()); |
|
|
|
|
|
propagationState.update(Map.of(PROPAGATION_CONFIG_ARGUMENT, entry), ctx); |
|
|
|
|
|
} |
|
|
|
|
|
msg.getCallback().onSuccess(); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
public void process(EntityCalculatedFieldTelemetryMsg msg) throws CalculatedFieldException { |
|
|
public void process(EntityCalculatedFieldTelemetryMsg msg) throws CalculatedFieldException { |
|
|
|