|
|
|
@ -18,19 +18,18 @@ package org.thingsboard.server.actors.calculatedField; |
|
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
|
import lombok.SneakyThrows; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.jetbrains.annotations.NotNull; |
|
|
|
import org.thingsboard.common.util.DebugModeUtil; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
import org.thingsboard.server.actors.ActorSystemContext; |
|
|
|
import org.thingsboard.server.actors.TbActorCtx; |
|
|
|
import org.thingsboard.server.actors.shared.AbstractContextAwareMsgProcessor; |
|
|
|
import org.thingsboard.server.common.data.AttributeScope; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; |
|
|
|
import org.thingsboard.server.common.data.id.AssetProfileId; |
|
|
|
import org.thingsboard.server.common.data.id.CalculatedFieldId; |
|
|
|
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.page.PageDataIterable; |
|
|
|
import org.thingsboard.server.common.data.msg.TbMsgType; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.AttributeScopeProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.AttributeValueProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; |
|
|
|
@ -53,9 +52,7 @@ import java.util.List; |
|
|
|
import java.util.Map; |
|
|
|
import java.util.Set; |
|
|
|
import java.util.UUID; |
|
|
|
import java.util.concurrent.ExecutionException; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
import java.util.concurrent.TimeoutException; |
|
|
|
|
|
|
|
|
|
|
|
/** |
|
|
|
@ -111,9 +108,9 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM |
|
|
|
callback.onSuccess(CALLBACKS_PER_CF); |
|
|
|
} else { |
|
|
|
if (proto.getTsDataCount() > 0) { |
|
|
|
processArgumentValuesUpdate(ctx, cfIds, callback, mapToArguments(ctx, msg.getEntityId(), proto.getTsDataList())); |
|
|
|
processArgumentValuesUpdate(ctx, cfIds, callback, mapToArguments(ctx, msg.getEntityId(), proto.getTsDataList()), toTbMsgId(proto), toTbMsgType(proto)); |
|
|
|
} else if (proto.getAttrDataCount() > 0) { |
|
|
|
processArgumentValuesUpdate(ctx, cfIds, callback, mapToArguments(ctx, msg.getEntityId(), proto.getScope(), proto.getAttrDataList())); |
|
|
|
processArgumentValuesUpdate(ctx, cfIds, callback, mapToArguments(ctx, msg.getEntityId(), proto.getScope(), proto.getAttrDataList()), toTbMsgId(proto), toTbMsgType(proto)); |
|
|
|
} else { |
|
|
|
callback.onSuccess(CALLBACKS_PER_CF); |
|
|
|
} |
|
|
|
@ -136,27 +133,30 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM |
|
|
|
|
|
|
|
@SneakyThrows |
|
|
|
private void processTelemetry(CalculatedFieldCtx ctx, CalculatedFieldTelemetryMsgProto proto, List<CalculatedFieldId> cfIdList, MultipleTbCallback callback) { |
|
|
|
processArgumentValuesUpdate(ctx, cfIdList, callback, mapToArguments(ctx, proto.getTsDataList())); |
|
|
|
processArgumentValuesUpdate(ctx, cfIdList, callback, mapToArguments(ctx, proto.getTsDataList()), toTbMsgId(proto), toTbMsgType(proto)); |
|
|
|
} |
|
|
|
|
|
|
|
@SneakyThrows |
|
|
|
private void processAttributes(CalculatedFieldCtx ctx, CalculatedFieldTelemetryMsgProto proto, List<CalculatedFieldId> cfIdList, MultipleTbCallback callback) { |
|
|
|
processArgumentValuesUpdate(ctx, cfIdList, callback, mapToArguments(ctx, proto.getScope(), proto.getAttrDataList())); |
|
|
|
processArgumentValuesUpdate(ctx, cfIdList, callback, mapToArguments(ctx, proto.getScope(), proto.getAttrDataList()), toTbMsgId(proto), toTbMsgType(proto)); |
|
|
|
} |
|
|
|
|
|
|
|
@SneakyThrows |
|
|
|
private void processArgumentValuesUpdate(CalculatedFieldCtx ctx, List<CalculatedFieldId> cfIdList, MultipleTbCallback callback, |
|
|
|
Map<String, ArgumentEntry> newArgValues) { |
|
|
|
Map<String, ArgumentEntry> newArgValues, UUID tbMsgId, TbMsgType tbMsgType) { |
|
|
|
if (newArgValues.isEmpty()) { |
|
|
|
callback.onSuccess(CALLBACKS_PER_CF); |
|
|
|
} |
|
|
|
CalculatedFieldState state = getOrInitState(ctx); |
|
|
|
if (state.updateState(newArgValues)) { |
|
|
|
if (state.isReady()) { |
|
|
|
if (state.isReady() && ctx.isInitialized()) { |
|
|
|
CalculatedFieldResult calculationResult = state.performCalculation(ctx).get(5, TimeUnit.SECONDS); |
|
|
|
cfIdList = new ArrayList<>(cfIdList); |
|
|
|
cfIdList.add(ctx.getCfId()); |
|
|
|
cfService.pushMsgToRuleEngine(tenantId, entityId, calculationResult, cfIdList, callback); |
|
|
|
if (DebugModeUtil.isDebugAllAvailable(ctx.getCalculatedField())) { |
|
|
|
systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId, state.getArguments(), tbMsgId, tbMsgType, JacksonUtil.writeValueAsString(calculationResult.getResultMap()), null); |
|
|
|
} |
|
|
|
} else { |
|
|
|
callback.onSuccess(); // State was updated but no calculation performed;
|
|
|
|
} |
|
|
|
@ -183,13 +183,27 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM |
|
|
|
return state; |
|
|
|
} |
|
|
|
|
|
|
|
private UUID toTbMsgId(CalculatedFieldTelemetryMsgProto proto) { |
|
|
|
if (proto.getTbMsgIdMSB() != 0 && proto.getTbMsgIdLSB() != 0) { |
|
|
|
return new UUID(proto.getTbMsgIdMSB(), proto.getTbMsgIdLSB()); |
|
|
|
} |
|
|
|
return null; |
|
|
|
} |
|
|
|
|
|
|
|
private TbMsgType toTbMsgType(CalculatedFieldTelemetryMsgProto proto) { |
|
|
|
if (!proto.getTbMsgType().isEmpty()) { |
|
|
|
return TbMsgType.valueOf(proto.getTbMsgType()); |
|
|
|
} |
|
|
|
return null; |
|
|
|
} |
|
|
|
|
|
|
|
private Map<String, ArgumentEntry> mapToArguments(CalculatedFieldCtx ctx, List<TsKvProto> data) { |
|
|
|
return mapToArguments(ctx.getMainEntityArguments(), data); |
|
|
|
} |
|
|
|
|
|
|
|
private Map<String, ArgumentEntry> mapToArguments(CalculatedFieldCtx ctx, EntityId entityId, List<TsKvProto> data) { |
|
|
|
var argNames = ctx.getLinkedEntityArguments().get(entityId); |
|
|
|
if(argNames.isEmpty()) { |
|
|
|
if (argNames.isEmpty()) { |
|
|
|
return Collections.emptyMap(); |
|
|
|
} |
|
|
|
return mapToArguments(argNames, data); |
|
|
|
@ -221,7 +235,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM |
|
|
|
|
|
|
|
private Map<String, ArgumentEntry> mapToArguments(CalculatedFieldCtx ctx, EntityId entityId, AttributeScopeProto scope, List<AttributeValueProto> attrDataList) { |
|
|
|
var argNames = ctx.getLinkedEntityArguments().get(entityId); |
|
|
|
if(argNames.isEmpty()) { |
|
|
|
if (argNames.isEmpty()) { |
|
|
|
return Collections.emptyMap(); |
|
|
|
} |
|
|
|
return mapToArguments(argNames, scope, attrDataList); |
|
|
|
|