|
|
@ -21,6 +21,7 @@ import com.google.gson.JsonElement; |
|
|
import com.google.gson.JsonParser; |
|
|
import com.google.gson.JsonParser; |
|
|
import com.google.gson.JsonPrimitive; |
|
|
import com.google.gson.JsonPrimitive; |
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
|
|
import org.jetbrains.annotations.NotNull; |
|
|
import org.thingsboard.common.util.CollectionsUtil; |
|
|
import org.thingsboard.common.util.CollectionsUtil; |
|
|
import org.thingsboard.common.util.DonAsynchron; |
|
|
import org.thingsboard.common.util.DonAsynchron; |
|
|
import org.thingsboard.rule.engine.api.EmptyNodeConfiguration; |
|
|
import org.thingsboard.rule.engine.api.EmptyNodeConfiguration; |
|
|
@ -90,26 +91,7 @@ public class TbCopyAttributesToEntityViewNode implements TbNode { |
|
|
long startTime = entityView.getStartTimeMs(); |
|
|
long startTime = entityView.getStartTimeMs(); |
|
|
long endTime = entityView.getEndTimeMs(); |
|
|
long endTime = entityView.getEndTimeMs(); |
|
|
if ((endTime != 0 && endTime > now && startTime < now) || (endTime == 0 && startTime < now)) { |
|
|
if ((endTime != 0 && endTime > now && startTime < now) || (endTime == 0 && startTime < now)) { |
|
|
if (DataConstants.ATTRIBUTES_UPDATED.equals(msg.getType()) || |
|
|
if (DataConstants.ATTRIBUTES_DELETED.equals(msg.getType())) { |
|
|
DataConstants.ACTIVITY_EVENT.equals(msg.getType()) || |
|
|
|
|
|
DataConstants.INACTIVITY_EVENT.equals(msg.getType()) || |
|
|
|
|
|
SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msg.getType())) { |
|
|
|
|
|
Set<AttributeKvEntry> attributes = JsonConverter.convertToAttributes(new JsonParser().parse(msg.getData())); |
|
|
|
|
|
List<AttributeKvEntry> filteredAttributes = |
|
|
|
|
|
attributes.stream().filter(attr -> attributeContainsInEntityView(scope, attr.getKey(), entityView)).collect(Collectors.toList()); |
|
|
|
|
|
ctx.getTelemetryService().saveAndNotify(ctx.getTenantId(), entityView.getId(), scope, filteredAttributes, |
|
|
|
|
|
new FutureCallback<Void>() { |
|
|
|
|
|
@Override |
|
|
|
|
|
public void onSuccess(@Nullable Void result) { |
|
|
|
|
|
transformAndTellNext(ctx, msg, entityView); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
ctx.tellFailure(msg, t); |
|
|
|
|
|
} |
|
|
|
|
|
}); |
|
|
|
|
|
} else if (DataConstants.ATTRIBUTES_DELETED.equals(msg.getType())) { |
|
|
|
|
|
List<String> attributes = new ArrayList<>(); |
|
|
List<String> attributes = new ArrayList<>(); |
|
|
for (JsonElement element : new JsonParser().parse(msg.getData()).getAsJsonObject().get("attributes").getAsJsonArray()) { |
|
|
for (JsonElement element : new JsonParser().parse(msg.getData()).getAsJsonObject().get("attributes").getAsJsonArray()) { |
|
|
if (element.isJsonPrimitive()) { |
|
|
if (element.isJsonPrimitive()) { |
|
|
@ -122,9 +104,14 @@ public class TbCopyAttributesToEntityViewNode implements TbNode { |
|
|
List<String> filteredAttributes = |
|
|
List<String> filteredAttributes = |
|
|
attributes.stream().filter(attr -> attributeContainsInEntityView(scope, attr, entityView)).collect(Collectors.toList()); |
|
|
attributes.stream().filter(attr -> attributeContainsInEntityView(scope, attr, entityView)).collect(Collectors.toList()); |
|
|
if (!filteredAttributes.isEmpty()) { |
|
|
if (!filteredAttributes.isEmpty()) { |
|
|
ctx.getAttributesService().removeAll(ctx.getTenantId(), entityView.getId(), scope, filteredAttributes); |
|
|
ctx.getTelemetryService().deleteAndNotify(ctx.getTenantId(), entityView.getId(), scope, filteredAttributes, getFutureCallback(ctx, msg, entityView)); |
|
|
transformAndTellNext(ctx, msg, entityView); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
} else { |
|
|
|
|
|
Set<AttributeKvEntry> attributes = JsonConverter.convertToAttributes(new JsonParser().parse(msg.getData())); |
|
|
|
|
|
List<AttributeKvEntry> filteredAttributes = |
|
|
|
|
|
attributes.stream().filter(attr -> attributeContainsInEntityView(scope, attr.getKey(), entityView)).collect(Collectors.toList()); |
|
|
|
|
|
ctx.getTelemetryService().saveAndNotify(ctx.getTenantId(), entityView.getId(), scope, filteredAttributes, |
|
|
|
|
|
getFutureCallback(ctx, msg, entityView)); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
@ -139,6 +126,21 @@ public class TbCopyAttributesToEntityViewNode implements TbNode { |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@NotNull |
|
|
|
|
|
private FutureCallback<Void> getFutureCallback(TbContext ctx, TbMsg msg, EntityView entityView) { |
|
|
|
|
|
return new FutureCallback<Void>() { |
|
|
|
|
|
@Override |
|
|
|
|
|
public void onSuccess(@Nullable Void result) { |
|
|
|
|
|
transformAndTellNext(ctx, msg, entityView); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
ctx.tellFailure(msg, t); |
|
|
|
|
|
} |
|
|
|
|
|
}; |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
private void transformAndTellNext(TbContext ctx, TbMsg msg, EntityView entityView) { |
|
|
private void transformAndTellNext(TbContext ctx, TbMsg msg, EntityView entityView) { |
|
|
ctx.enqueueForTellNext(ctx.newMsg(msg.getQueueName(), msg.getType(), entityView.getId(), msg.getCustomerId(), msg.getMetaData(), msg.getData()), SUCCESS); |
|
|
ctx.enqueueForTellNext(ctx.newMsg(msg.getQueueName(), msg.getType(), entityView.getId(), msg.getCustomerId(), msg.getMetaData(), msg.getData()), SUCCESS); |
|
|
} |
|
|
} |
|
|
|