Browse Source

Merge pull request #14603 from volodymyr-babak/merge-downlink-duplicates

Edge Events - added merge and filter duplicates
pull/14630/head
Viacheslav Klimov 9 months ago
committed by GitHub
parent
commit
8581dbd332
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 133
      application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java
  2. 3
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  3. 1
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java
  4. 25
      application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java
  5. 13
      application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java
  6. 12
      application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java
  7. 4
      application/src/test/java/org/thingsboard/server/edge/TelemetryEdgeTest.java
  8. 98
      application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java

133
application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java

@ -21,6 +21,7 @@ import com.google.gson.Gson;
import com.google.gson.JsonArray; import com.google.gson.JsonArray;
import com.google.gson.JsonElement; import com.google.gson.JsonElement;
import com.google.gson.JsonObject; import com.google.gson.JsonObject;
import com.google.gson.JsonParser;
import com.google.gson.JsonPrimitive; import com.google.gson.JsonPrimitive;
import com.google.gson.reflect.TypeToken; import com.google.gson.reflect.TypeToken;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
@ -52,6 +53,7 @@ import org.thingsboard.server.common.data.asset.AssetProfile;
import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.domain.DomainInfo; import org.thingsboard.server.common.data.domain.DomainInfo;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.id.AiModelId; import org.thingsboard.server.common.data.id.AiModelId;
import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetId;
@ -76,6 +78,7 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.id.WidgetTypeId; import org.thingsboard.server.common.data.id.WidgetTypeId;
import org.thingsboard.server.common.data.id.WidgetsBundleId; import org.thingsboard.server.common.data.id.WidgetsBundleId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.notification.rule.NotificationRule; import org.thingsboard.server.common.data.notification.rule.NotificationRule;
import org.thingsboard.server.common.data.notification.targets.NotificationTarget; import org.thingsboard.server.common.data.notification.targets.NotificationTarget;
import org.thingsboard.server.common.data.notification.template.NotificationTemplate; import org.thingsboard.server.common.data.notification.template.NotificationTemplate;
@ -88,6 +91,7 @@ import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.UserCredentials; import org.thingsboard.server.common.data.security.UserCredentials;
import org.thingsboard.server.common.data.widget.WidgetTypeDetails; import org.thingsboard.server.common.data.widget.WidgetTypeDetails;
import org.thingsboard.server.common.data.widget.WidgetsBundle; import org.thingsboard.server.common.data.widget.WidgetsBundle;
import org.thingsboard.server.common.transport.util.JsonUtils;
import org.thingsboard.server.gen.edge.v1.AiModelUpdateMsg; import org.thingsboard.server.gen.edge.v1.AiModelUpdateMsg;
import org.thingsboard.server.gen.edge.v1.AlarmCommentUpdateMsg; import org.thingsboard.server.gen.edge.v1.AlarmCommentUpdateMsg;
import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg;
@ -127,11 +131,16 @@ import org.thingsboard.server.gen.edge.v1.WidgetTypeUpdateMsg;
import org.thingsboard.server.gen.edge.v1.WidgetsBundleUpdateMsg; import org.thingsboard.server.gen.edge.v1.WidgetsBundleUpdateMsg;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Iterator; import java.util.Iterator;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Set; import java.util.Set;
import java.util.UUID; import java.util.UUID;
import java.util.stream.Collectors;
@Slf4j @Slf4j
public class EdgeMsgConstructorUtils { public class EdgeMsgConstructorUtils {
@ -670,4 +679,128 @@ public class EdgeMsgConstructorUtils {
.setIdLSB(aiModelId.getId().getLeastSignificantBits()).build(); .setIdLSB(aiModelId.getId().getLeastSignificantBits()).build();
} }
public static List<EdgeEvent> mergeAndFilterDownlinkDuplicates(List<EdgeEvent> edgeEvents) {
try {
edgeEvents = removeDownlinkDuplicates(edgeEvents);
List<AttrUpdateMsg> attrUpdateMsgs = new ArrayList<>();
for (EdgeEvent edgeEvent : edgeEvents) {
if (EdgeEventActionType.ATTRIBUTES_UPDATED.equals(edgeEvent.getAction())) {
attrUpdateMsgs.add(new AttrUpdateMsg(edgeEvent.getEntityId(), edgeEvent.getBody()));
}
}
Map<UUID, Map<String, Long>> latestTsByEntityAndKey = computeLatestTsByEntityAndKey(attrUpdateMsgs);
List<EdgeEvent> result = new ArrayList<>();
for (EdgeEvent edgeEvent : edgeEvents) {
if (!EdgeEventActionType.ATTRIBUTES_UPDATED.equals(edgeEvent.getAction())) {
result.add(edgeEvent);
continue;
}
Map<String, Long> latestByKey = latestTsByEntityAndKey.get(edgeEvent.getEntityId());
JsonNode filteredBody = filterAttributesBody(edgeEvent.getBody(), latestByKey);
if (filteredBody == null) {
continue;
}
result.add(createFilteredEdgeEvent(edgeEvent, filteredBody));
}
result.sort(Comparator.comparingLong(EdgeEvent::getSeqId));
return result;
} catch (Exception e) {
log.warn("Can't merge downlink duplicates, edgeEvents [{}]", edgeEvents, e);
return edgeEvents;
}
}
private static AttrsTs extractAttributes(JsonNode body) {
if (body == null) {
return new AttrsTs(0L, List.of());
}
String bodyStr = JacksonUtil.toString(body);
var jsonObject = JsonParser.parseString(bodyStr).getAsJsonObject();
long ts = jsonObject.get("ts").getAsLong();
var kv = jsonObject.getAsJsonObject("kv");
List<AttributeKvEntry> attrs = JsonConverter.convertToAttributes(
JsonUtils.getJsonObject(
JsonConverter.convertToAttributesProto(kv).getKvList()
), ts);
return new AttrsTs(ts, attrs);
}
private static JsonNode filterAttributesBody(JsonNode body, Map<String, Long> latestByKey) {
if (body == null || latestByKey == null || latestByKey.isEmpty()) {
return null;
}
String bodyStr = JacksonUtil.toString(body);
JsonObject jsonObject = JsonParser.parseString(bodyStr).getAsJsonObject();
long ts = jsonObject.get("ts").getAsLong();
JsonObject kv = jsonObject.getAsJsonObject("kv");
for (Iterator<Map.Entry<String, JsonElement>> it = kv.entrySet().iterator(); it.hasNext(); ) {
Map.Entry<String, JsonElement> e = it.next();
Long latestTs = latestByKey.get(e.getKey());
if (latestTs == null || !latestTs.equals(ts)) {
it.remove();
}
}
if (kv.isEmpty()) {
return null;
}
return JacksonUtil.toJsonNode(jsonObject.toString());
}
private static Map<UUID, Map<String, Long>> computeLatestTsByEntityAndKey(List<AttrUpdateMsg> attrUpdateMsgs) {
Map<UUID, Map<String, Long>> latestTsByEntityAndKey = new HashMap<>();
for (AttrUpdateMsg attrUpdateMsg : attrUpdateMsgs) {
UUID entityId = attrUpdateMsg.entityId();
AttrsTs attrsTs = extractAttributes(attrUpdateMsg.body());
Map<String, Long> map = latestTsByEntityAndKey.computeIfAbsent(entityId, id -> new HashMap<>());
long ts = attrsTs.ts();
for (AttributeKvEntry attr : attrsTs.attrs()) {
map.merge(attr.getKey(), ts, Math::max);
}
}
return latestTsByEntityAndKey;
}
private static EdgeEvent createFilteredEdgeEvent(EdgeEvent edgeEvent, JsonNode filteredBody) {
EdgeEvent filtered = new EdgeEvent();
filtered.setSeqId(edgeEvent.getSeqId());
filtered.setTenantId(edgeEvent.getTenantId());
filtered.setEdgeId(edgeEvent.getEdgeId());
filtered.setAction(edgeEvent.getAction());
filtered.setEntityId(edgeEvent.getEntityId());
filtered.setUid(edgeEvent.getUid());
filtered.setType(edgeEvent.getType());
filtered.setBody(filteredBody);
return filtered;
}
private static List<EdgeEvent> removeDownlinkDuplicates(List<EdgeEvent> edgeEvents) {
Set<EventKey> seen = new HashSet<>();
return edgeEvents.stream()
.filter(e -> seen.add(new EventKey(
e.getTenantId(),
e.getAction(),
e.getEntityId(),
e.getType().name(),
(e.getBody() != null ? e.getBody().toString() : "null"))))
.collect(Collectors.toList());
}
private record EventKey(TenantId tenantId,
EdgeEventActionType action,
UUID entityId,
String type,
String body) {
}
private record AttrsTs(long ts, List<AttributeKvEntry> attrs) {
}
private record AttrUpdateMsg(UUID entityId, JsonNode body) {
}
} }

3
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java

@ -654,7 +654,8 @@ public abstract class EdgeGrpcSession implements Closeable {
protected List<DownlinkMsg> convertToDownlinkMsgsPack(List<EdgeEvent> edgeEvents) { protected List<DownlinkMsg> convertToDownlinkMsgsPack(List<EdgeEvent> edgeEvents) {
List<DownlinkMsg> result = new ArrayList<>(); List<DownlinkMsg> result = new ArrayList<>();
for (EdgeEvent edgeEvent : edgeEvents) { List<EdgeEvent> filtered = EdgeMsgConstructorUtils.mergeAndFilterDownlinkDuplicates(edgeEvents);
for (EdgeEvent edgeEvent : filtered) {
log.trace("[{}][{}] converting edge event to downlink msg [{}]", tenantId, edge.getId(), edgeEvent); log.trace("[{}][{}] converting edge event to downlink msg [{}]", tenantId, edge.getId(), edgeEvent);
DownlinkMsg downlinkMsg = null; DownlinkMsg downlinkMsg = null;
try { try {

1
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeSyncCursor.java

@ -92,6 +92,7 @@ public class EdgeSyncCursor {
fetchers.add(new NotificationTargetEdgeEventFetcher(ctx.getNotificationTargetService())); fetchers.add(new NotificationTargetEdgeEventFetcher(ctx.getNotificationTargetService()));
fetchers.add(new NotificationRuleEdgeEventFetcher(ctx.getNotificationRuleService())); fetchers.add(new NotificationRuleEdgeEventFetcher(ctx.getNotificationRuleService()));
fetchers.add(new OtaPackagesEdgeEventFetcher(ctx.getOtaPackageService())); fetchers.add(new OtaPackagesEdgeEventFetcher(ctx.getOtaPackageService()));
// sync device profiles twice to update software and hardware fields
fetchers.add(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService())); fetchers.add(new DeviceProfilesEdgeEventFetcher(ctx.getDeviceProfileService()));
fetchers.add(new TenantResourcesEdgeEventFetcher(ctx.getResourceService())); fetchers.add(new TenantResourcesEdgeEventFetcher(ctx.getResourceService()));
} }

25
application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java

@ -22,6 +22,7 @@ import com.google.protobuf.AbstractMessage;
import com.google.protobuf.InvalidProtocolBufferException; import com.google.protobuf.InvalidProtocolBufferException;
import com.google.protobuf.MessageLite; import com.google.protobuf.MessageLite;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.awaitility.Awaitility;
import org.junit.After; import org.junit.After;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Before; import org.junit.Before;
@ -102,6 +103,7 @@ import org.thingsboard.server.gen.edge.v1.UserUpdateMsg;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.Map;
import java.util.Optional; import java.util.Optional;
import java.util.TreeMap; import java.util.TreeMap;
import java.util.UUID; import java.util.UUID;
@ -738,4 +740,27 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest {
return rpc; return rpc;
} }
protected void verifyEdgeDisconnected() {
verifyEdgeActiveFlag(false);
}
protected void verifyEdgeConnected() {
verifyEdgeActiveFlag(true);
}
private void verifyEdgeActiveFlag(boolean value) {
Awaitility.await()
.atMost(TIMEOUT, TimeUnit.SECONDS)
.until(() -> {
List<Map<String, Object>> values = doGetAsyncTyped("/api/plugins/telemetry/EDGE/" + edge.getId() +
"/values/attributes/SERVER_SCOPE", new TypeReference<>() {});
Optional<Map<String, Object>> activeAttrOpt = values.stream().filter(att -> att.get("key").equals("active")).findFirst();
if (activeAttrOpt.isEmpty()) {
return false;
}
Map<String, Object> activeAttr = activeAttrOpt.get();
return Boolean.toString(value).equals(activeAttr.get("value").toString());
});
}
} }

13
application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java

@ -893,18 +893,7 @@ public class DeviceEdgeTest extends AbstractEdgeTest {
ObjectNode attributes = JacksonUtil.newObjectNode(); ObjectNode attributes = JacksonUtil.newObjectNode();
attributes.put("active", true); attributes.put("active", true);
doPost("/api/plugins/telemetry/EDGE/" + edge.getId() + "/attributes/" + DataConstants.SERVER_SCOPE, attributes); doPost("/api/plugins/telemetry/EDGE/" + edge.getId() + "/attributes/" + DataConstants.SERVER_SCOPE, attributes);
Awaitility.await() verifyEdgeConnected();
.atMost(TIMEOUT, TimeUnit.SECONDS)
.until(() -> {
List<Map<String, Object>> values = doGetAsyncTyped("/api/plugins/telemetry/EDGE/" + edge.getId() +
"/values/attributes/SERVER_SCOPE", new TypeReference<>() {});
Optional<Map<String, Object>> activeAttrOpt = values.stream().filter(att -> att.get("key").equals("active")).findFirst();
if (activeAttrOpt.isEmpty()) {
return false;
}
Map<String, Object> activeAttr = activeAttrOpt.get();
return "true".equals(activeAttr.get("value").toString());
});
} }
} }

12
application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java

@ -151,14 +151,20 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest {
// delete profile when edge is offline // delete profile when edge is offline
edgeImitator.disconnect(); edgeImitator.disconnect();
verifyEdgeDisconnected();
doDelete("/api/deviceProfile/" + deviceProfile.getUuidId()) doDelete("/api/deviceProfile/" + deviceProfile.getUuidId())
.andExpect(status().isOk()); .andExpect(status().isOk());
edgeImitator.connect();
// 25 sync message // 25 sync message
// + 2 RuleChain and RuleChainMetadata // + 1 RuleChain Added
// + 1 delete DeviceProfile // + 1 RuleChainMetadata Added
// + 1 DeviceProfile Delete
edgeImitator.expectMessageAmount(SYNC_MESSAGE_COUNT + 3); edgeImitator.expectMessageAmount(SYNC_MESSAGE_COUNT + 3);
edgeImitator.connect();
verifyEdgeConnected();
Assert.assertTrue(edgeImitator.waitForMessages()); Assert.assertTrue(edgeImitator.waitForMessages());
latestMessage = edgeImitator.getLatestMessage(); latestMessage = edgeImitator.getLatestMessage();

4
application/src/test/java/org/thingsboard/server/edge/TelemetryEdgeTest.java

@ -223,8 +223,10 @@ public class TelemetryEdgeTest extends AbstractEdgeTest {
device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData); device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData);
edgeEventService.saveAsync(failedEdgeEvent).get(); edgeEventService.saveAsync(failedEdgeEvent).get();
// add unique body to avoid merge and filter by device id in edge service (mergeAndFilterDownlinkDuplicates)
JsonNode uniqueBody = JacksonUtil.toJsonNode("{\"idx\":" + idx + "}");
EdgeEvent successEdgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.UPDATED, EdgeEvent successEdgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.UPDATED,
device.getId().getId(), EdgeEventType.DEVICE, null); device.getId().getId(), EdgeEventType.DEVICE, uniqueBody);
edgeEventService.saveAsync(successEdgeEvent).get(); edgeEventService.saveAsync(successEdgeEvent).get();
} }

98
application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java

@ -15,10 +15,12 @@
*/ */
package org.thingsboard.server.service.edge; package org.thingsboard.server.service.edge;
import com.fasterxml.jackson.databind.JsonNode;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.DisplayName;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.TestInstance; import org.junit.jupiter.api.TestInstance;
import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.MethodSource; import org.junit.jupiter.params.provider.MethodSource;
@ -36,6 +38,10 @@ import org.thingsboard.rule.engine.rest.TbSendRestApiCallReplyNode;
import org.thingsboard.rule.engine.telemetry.TbCalculatedFieldsNode; import org.thingsboard.rule.engine.telemetry.TbCalculatedFieldsNode;
import org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode; import org.thingsboard.rule.engine.telemetry.TbMsgAttributesNode;
import org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode; import org.thingsboard.rule.engine.telemetry.TbMsgTimeseriesNode;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.rule.RuleChainMetaData;
import org.thingsboard.server.common.data.rule.RuleNode; import org.thingsboard.server.common.data.rule.RuleNode;
import org.thingsboard.server.gen.edge.v1.EdgeVersion; import org.thingsboard.server.gen.edge.v1.EdgeVersion;
@ -44,6 +50,7 @@ import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
import java.lang.reflect.Constructor; import java.lang.reflect.Constructor;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.Optional;
import java.util.UUID;
import java.util.stream.Stream; import java.util.stream.Stream;
import static org.thingsboard.server.service.edge.EdgeMsgConstructorUtils.EXCLUDED_NODES_BY_EDGE_VERSION; import static org.thingsboard.server.service.edge.EdgeMsgConstructorUtils.EXCLUDED_NODES_BY_EDGE_VERSION;
@ -155,4 +162,95 @@ public class EdgeMsgConstructorUtilsTest {
String.format("For EdgeVersion '%s', ruleNode '%s' should not be included.", edgeVersion, ruleNode.getType())); String.format("For EdgeVersion '%s', ruleNode '%s' should not be included.", edgeVersion, ruleNode.getType()));
} }
@Test
@DisplayName("mergeDownlinkDuplicates: latest per attribute key is retained and duplicates removed")
public void testMergeDownlinkDuplicates() {
UUID deviceId = UUID.randomUUID();
UUID assetId = UUID.randomUUID();
TenantId tenantId = TenantId.fromUUID(UUID.randomUUID());
var deviceAttrUpdate1 = createEdgeEvent(tenantId, 1, EdgeEventActionType.ATTRIBUTES_UPDATED,
deviceId, EdgeEventType.DEVICE, createAttrBody(1_000L, "{\"a\":1,\"b\":1,\"d\":1}"));
var deviceAttrUpdate2 = createEdgeEvent(tenantId, 2, EdgeEventActionType.ATTRIBUTES_UPDATED,
deviceId, EdgeEventType.DEVICE, createAttrBody(2_000L, "{\"a\":2,\"b\":2,\"c\":2}"));
var deviceAttrUpdate3 = createEdgeEvent(tenantId, 3, EdgeEventActionType.ATTRIBUTES_UPDATED,
deviceId, EdgeEventType.DEVICE, createAttrBody(3_000L, "{\"a\":3,\"d\":3}"));
var deviceUpdate = createEdgeEvent(tenantId, 4, EdgeEventActionType.UPDATED,
deviceId, EdgeEventType.DEVICE, null);
var deviceUpdateDup = createEdgeEvent(tenantId, 5, EdgeEventActionType.UPDATED,
deviceId, EdgeEventType.DEVICE, null);
var assetAttrUpdate1 = createEdgeEvent(tenantId, 6, EdgeEventActionType.ATTRIBUTES_UPDATED,
assetId, EdgeEventType.ASSET, createAttrBody(6_000L, "{\"a\":6,\"d\":6}"));
var assetAttrUpdate2 = createEdgeEvent(tenantId, 7, EdgeEventActionType.ATTRIBUTES_UPDATED,
assetId, EdgeEventType.ASSET, createAttrBody(7_000L, "{\"a\":7,\"b\":7,\"c\":7}"));
var assetAttrUpdate3 = createEdgeEvent(tenantId, 8, EdgeEventActionType.ATTRIBUTES_UPDATED,
assetId, EdgeEventType.ASSET, createAttrBody(8_000L, "{\"a\":8,\"d\":8}"));
List<EdgeEvent> input = List.of(deviceAttrUpdate1, deviceAttrUpdate2, deviceAttrUpdate3,
deviceUpdate, deviceUpdateDup,
assetAttrUpdate1, assetAttrUpdate2, assetAttrUpdate3);
List<EdgeEvent> merged = EdgeMsgConstructorUtils.mergeAndFilterDownlinkDuplicates(input);
Assertions.assertEquals(5, merged.size());
EdgeEvent deviceMergedAttrBC = merged.get(0);
Assertions.assertEquals(2, deviceMergedAttrBC.getSeqId());
Assertions.assertEquals(deviceId, deviceMergedAttrBC.getEntityId());
Assertions.assertEquals(2_000L, deviceMergedAttrBC.getBody().get("ts").asLong());
Assertions.assertEquals(2, getIntValue(deviceMergedAttrBC.getBody(), "b"));
Assertions.assertEquals(2, getIntValue(deviceMergedAttrBC.getBody(), "c"));
Assertions.assertNull(getIntValue(deviceMergedAttrBC.getBody(), "a"));
EdgeEvent deviceMergedAttrAD = merged.get(1);
Assertions.assertEquals(3, deviceMergedAttrAD.getSeqId());
Assertions.assertEquals(deviceId, deviceMergedAttrAD.getEntityId());
Assertions.assertEquals(3_000L, deviceMergedAttrAD.getBody().get("ts").asLong());
Assertions.assertEquals(3, getIntValue(deviceMergedAttrAD.getBody(), "a"));
Assertions.assertEquals(3, getIntValue(deviceMergedAttrAD.getBody(), "d"));
EdgeEvent mergedDeviceUpdate = merged.get(2);
Assertions.assertEquals(4, mergedDeviceUpdate.getSeqId());
Assertions.assertEquals(EdgeEventActionType.UPDATED, mergedDeviceUpdate.getAction());
EdgeEvent assetMergedAttrBC = merged.get(3);
Assertions.assertEquals(7, assetMergedAttrBC.getSeqId());
Assertions.assertEquals(assetId, assetMergedAttrBC.getEntityId());
Assertions.assertEquals(7_000L, assetMergedAttrBC.getBody().get("ts").asLong());
Assertions.assertEquals(7, getIntValue(assetMergedAttrBC.getBody(), "b"));
Assertions.assertEquals(7, getIntValue(assetMergedAttrBC.getBody(), "c"));
Assertions.assertNull(getIntValue(assetMergedAttrBC.getBody(), "a"));
EdgeEvent assetMergedAttrAD = merged.get(4);
Assertions.assertEquals(8, assetMergedAttrAD.getSeqId());
Assertions.assertEquals(assetId, assetMergedAttrAD.getEntityId());
Assertions.assertEquals(8_000L, assetMergedAttrAD.getBody().get("ts").asLong());
Assertions.assertEquals(8, getIntValue(assetMergedAttrAD.getBody(), "a"));
Assertions.assertEquals(8, getIntValue(assetMergedAttrAD.getBody(), "d"));
}
private Integer getIntValue(JsonNode body, String key) {
return body.get("kv").get(key) != null ? body.get("kv").get(key).asInt() : null;
}
private static JsonNode createAttrBody(long ts, String kvJson) {
return JacksonUtil.toJsonNode("{\"ts\":" + ts + ",\"kv\":" + kvJson + "}");
}
private static EdgeEvent createEdgeEvent(TenantId tenantId,
long seqId,
EdgeEventActionType action,
UUID entityId,
EdgeEventType type,
JsonNode body) {
EdgeEvent edgeEvent = new EdgeEvent();
edgeEvent.setSeqId(seqId);
edgeEvent.setTenantId(tenantId);
edgeEvent.setAction(action);
edgeEvent.setEntityId(entityId);
edgeEvent.setType(type);
edgeEvent.setBody(body);
return edgeEvent;
}
} }

Loading…
Cancel
Save