|
|
|
@ -40,7 +40,6 @@ import org.thingsboard.server.gen.transport.TransportProtos.AttributeScopeProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.AttributeValueProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.ToCalculatedFieldMsg; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos.TsKvProto; |
|
|
|
import org.thingsboard.server.queue.TbQueueCallback; |
|
|
|
import org.thingsboard.server.queue.TbQueueMsgMetadata; |
|
|
|
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
|
|
|
@ -86,7 +85,10 @@ public class DefaultCalculatedFieldQueueService implements CalculatedFieldQueueS |
|
|
|
public void pushRequestToQueue(TimeseriesSaveRequest request, TimeseriesSaveResult result, FutureCallback<Void> callback) { |
|
|
|
var tenantId = request.getTenantId(); |
|
|
|
var entityId = request.getEntityId(); |
|
|
|
checkEntityAndPushToQueue(tenantId, entityId, cf -> cf.matches(request.getEntries()), cf -> cf.linkMatches(entityId, request.getEntries()), |
|
|
|
var entries = request.getEntries(); |
|
|
|
checkEntityAndPushToQueue(tenantId, entityId, |
|
|
|
cf -> cf.matches(entries), |
|
|
|
cf -> cf.linkMatches(entityId, entries), |
|
|
|
() -> toCalculatedFieldTelemetryMsgProto(request, result), callback); |
|
|
|
} |
|
|
|
|
|
|
|
@ -99,7 +101,11 @@ public class DefaultCalculatedFieldQueueService implements CalculatedFieldQueueS |
|
|
|
public void pushRequestToQueue(AttributesSaveRequest request, List<Long> result, FutureCallback<Void> callback) { |
|
|
|
var tenantId = request.getTenantId(); |
|
|
|
var entityId = request.getEntityId(); |
|
|
|
checkEntityAndPushToQueue(tenantId, entityId, cf -> cf.matches(request.getEntries(), request.getScope()), cf -> cf.linkMatches(entityId, request.getEntries(), request.getScope()), |
|
|
|
var entries = request.getEntries(); |
|
|
|
var scope = request.getScope(); |
|
|
|
checkEntityAndPushToQueue(tenantId, entityId, |
|
|
|
cf -> cf.matches(entries, scope), |
|
|
|
cf -> cf.linkMatches(entityId, entries, scope), |
|
|
|
() -> toCalculatedFieldTelemetryMsgProto(request, result), callback); |
|
|
|
} |
|
|
|
|
|
|
|
@ -112,7 +118,10 @@ public class DefaultCalculatedFieldQueueService implements CalculatedFieldQueueS |
|
|
|
public void pushRequestToQueue(AttributesDeleteRequest request, List<String> result, FutureCallback<Void> callback) { |
|
|
|
var tenantId = request.getTenantId(); |
|
|
|
var entityId = request.getEntityId(); |
|
|
|
checkEntityAndPushToQueue(tenantId, entityId, cf -> cf.matchesKeys(result, request.getScope()), cf -> cf.linkMatchesAttrKeys(entityId, result, request.getScope()), |
|
|
|
var scope = request.getScope(); |
|
|
|
checkEntityAndPushToQueue(tenantId, entityId, |
|
|
|
cf -> cf.matchesKeys(result, scope), |
|
|
|
cf -> cf.linkMatchesAttrKeys(entityId, result, scope), |
|
|
|
() -> toCalculatedFieldTelemetryMsgProto(request, result), callback); |
|
|
|
} |
|
|
|
|
|
|
|
@ -120,8 +129,9 @@ public class DefaultCalculatedFieldQueueService implements CalculatedFieldQueueS |
|
|
|
public void pushRequestToQueue(TimeseriesDeleteRequest request, List<String> result, FutureCallback<Void> callback) { |
|
|
|
var tenantId = request.getTenantId(); |
|
|
|
var entityId = request.getEntityId(); |
|
|
|
|
|
|
|
checkEntityAndPushToQueue(tenantId, entityId, cf -> cf.matchesKeys(result), cf -> cf.linkMatchesTsKeys(entityId, result), |
|
|
|
checkEntityAndPushToQueue(tenantId, entityId, |
|
|
|
cf -> cf.matchesKeys(result), |
|
|
|
cf -> cf.linkMatchesTsKeys(entityId, result), |
|
|
|
() -> toCalculatedFieldTelemetryMsgProto(request, result), callback); |
|
|
|
} |
|
|
|
|
|
|
|
@ -183,23 +193,20 @@ public class DefaultCalculatedFieldQueueService implements CalculatedFieldQueueS |
|
|
|
} |
|
|
|
|
|
|
|
private ToCalculatedFieldMsg toCalculatedFieldTelemetryMsgProto(TimeseriesSaveRequest request, TimeseriesSaveResult result) { |
|
|
|
ToCalculatedFieldMsg.Builder msg = ToCalculatedFieldMsg.newBuilder(); |
|
|
|
|
|
|
|
CalculatedFieldTelemetryMsgProto.Builder telemetryMsg = buildTelemetryMsgProto(request.getTenantId(), request.getEntityId(), request.getPreviousCalculatedFieldIds(), request.getTbMsgId(), request.getTbMsgType()); |
|
|
|
|
|
|
|
List<TsKvEntry> entries = request.getEntries(); |
|
|
|
List<Long> versions = result != null ? result.getVersions() : Collections.emptyList(); |
|
|
|
|
|
|
|
for (int i = 0; i < entries.size(); i++) { |
|
|
|
TsKvProto.Builder tsProtoBuilder = toTsKvProto(entries.get(i)).toBuilder(); |
|
|
|
TsKvEntry tsKvEntry = entries.get(i); |
|
|
|
if (result != null) { |
|
|
|
tsProtoBuilder.setVersion(versions.get(i)); |
|
|
|
tsKvEntry.setVersion(versions.get(i)); |
|
|
|
} |
|
|
|
telemetryMsg.addTsData(tsProtoBuilder.build()); |
|
|
|
telemetryMsg.addTsData(toTsKvProto(tsKvEntry)); |
|
|
|
} |
|
|
|
|
|
|
|
msg.setTelemetryMsg(telemetryMsg.build()); |
|
|
|
return msg.build(); |
|
|
|
return ToCalculatedFieldMsg.newBuilder().setTelemetryMsg(telemetryMsg).build(); |
|
|
|
} |
|
|
|
|
|
|
|
private ToCalculatedFieldMsg toCalculatedFieldTelemetryMsgProto(AttributesSaveRequest request, List<Long> versions) { |
|
|
|
|