diff --git a/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java b/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java index e1c7f1198c..33084f98fa 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtils.java @@ -560,7 +560,7 @@ public class EdgeMsgConstructorUtils { .setEntityIdMSB(entityId.getId().getMostSignificantBits()) .setEntityIdLSB(entityId.getId().getLeastSignificantBits()) .setEntityType(entityId.getEntityType().name()); - long ts = getTs(entityData.getAsJsonObject()); + long ts = extractTs(entityData.getAsJsonObject()); switch (actionType) { case TIMESERIES_UPDATED: try { @@ -613,8 +613,8 @@ public class EdgeMsgConstructorUtils { return builder.build(); } - private static long getTs(JsonObject data) { - if (data.get("ts") != null && !data.get("ts").isJsonNull()) { + private static long extractTs(JsonObject data) { + if (data.has("ts") && data.get("ts").isJsonPrimitive()) { return data.getAsJsonPrimitive("ts").getAsLong(); } return System.currentTimeMillis(); @@ -740,7 +740,7 @@ public class EdgeMsgConstructorUtils { result.sort(Comparator.comparingLong(EdgeEvent::getSeqId)); return result; } catch (Exception e) { - log.warn("Can't merge downlink duplicates, edgeEvents [{}]", edgeEvents, e); + log.info("Can't merge downlink duplicates. Sending downlinks without merge. Original edgeEvents [{}]", edgeEvents, e); return edgeEvents; } } @@ -751,6 +751,9 @@ public class EdgeMsgConstructorUtils { } String bodyStr = JacksonUtil.toString(body); var jsonObject = JsonParser.parseString(bodyStr).getAsJsonObject(); + if (!jsonObject.has("ts")) { + return new AttrsTs(0L, List.of()); + } long ts = jsonObject.get("ts").getAsLong(); var kv = jsonObject.getAsJsonObject("kv"); List attrs = JsonConverter.convertToAttributes( @@ -761,22 +764,24 @@ public class EdgeMsgConstructorUtils { } private static JsonNode filterAttributesBody(JsonNode body, Map latestByKey) { - if (body == null || latestByKey == null || latestByKey.isEmpty()) { + if (body == null) { 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> it = kv.entrySet().iterator(); it.hasNext(); ) { - Map.Entry e = it.next(); - Long latestTs = latestByKey.get(e.getKey()); - if (latestTs == null || !latestTs.equals(ts)) { - it.remove(); + if (jsonObject.has("ts") && latestByKey != null && !latestByKey.isEmpty()) { + long ts = jsonObject.get("ts").getAsLong(); + JsonObject kv = jsonObject.getAsJsonObject("kv"); + for (Iterator> it = kv.entrySet().iterator(); it.hasNext(); ) { + Map.Entry e = it.next(); + Long latestTs = latestByKey.get(e.getKey()); + if (latestTs == null || !latestTs.equals(ts)) { + it.remove(); + } + } + if (kv.isEmpty()) { + return null; } - } - if (kv.isEmpty()) { - return null; } return JacksonUtil.toJsonNode(jsonObject.toString()); } diff --git a/application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java index 380dec8015..5c3f43b451 100644 --- a/application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/DeviceEdgeTest.java @@ -49,8 +49,6 @@ import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.ota.OtaPackageType; -import org.thingsboard.server.common.data.page.PageData; -import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; @@ -671,8 +669,7 @@ public class DeviceEdgeTest extends AbstractEdgeTest { .atMost(10, TimeUnit.SECONDS) .until(() -> { String urlTemplate = "/api/plugins/telemetry/DEVICE/" + device.getId() + "/keys/attributes/" + scope; - List actualKeys = doGetAsyncTyped(urlTemplate, new TypeReference<>() { - }); + List actualKeys = doGetAsyncTyped(urlTemplate, new TypeReference<>() {}); return actualKeys != null && !actualKeys.isEmpty() && actualKeys.contains(expectedKey); }); diff --git a/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java b/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java index ebb8305871..1e268e3bcb 100644 --- a/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java +++ b/application/src/test/java/org/thingsboard/server/service/edge/EdgeMsgConstructorUtilsTest.java @@ -230,10 +230,35 @@ public class EdgeMsgConstructorUtilsTest { Assertions.assertEquals(8, getIntValue(assetMergedAttrAD.getBody(), "d")); } + @Test + public void testMergeDownlinkDuplicates_attrBodyHasNoTs_returnOriginalList() { + UUID deviceId = UUID.randomUUID(); + TenantId tenantId = TenantId.fromUUID(UUID.randomUUID()); + + var deviceAttrUpdate1 = createEdgeEvent(tenantId, 1, EdgeEventActionType.ATTRIBUTES_UPDATED, + deviceId, EdgeEventType.DEVICE, createAttrBodyWithoutTs("{\"a\":1,\"b\":1,\"d\":1}")); + var deviceAttrUpdate2 = createEdgeEvent(tenantId, 2, EdgeEventActionType.ATTRIBUTES_UPDATED, + deviceId, EdgeEventType.DEVICE, createAttrBodyWithoutTs("{\"a\":2,\"b\":2,\"c\":2}")); + var deviceAttrUpdate3 = createEdgeEvent(tenantId, 3, EdgeEventActionType.ATTRIBUTES_UPDATED, + deviceId, EdgeEventType.DEVICE, createAttrBodyWithoutTs("{\"a\":3,\"d\":3}")); + + List input = List.of(deviceAttrUpdate1, deviceAttrUpdate2, deviceAttrUpdate3); + List merged = EdgeMsgConstructorUtils.mergeAndFilterDownlinkDuplicates(input); + + Assertions.assertEquals(3, merged.size()); + Assertions.assertEquals(deviceAttrUpdate1, merged.get(0)); + Assertions.assertEquals(deviceAttrUpdate2, merged.get(1)); + Assertions.assertEquals(deviceAttrUpdate3, merged.get(2)); + } + private Integer getIntValue(JsonNode body, String key) { return body.get("kv").get(key) != null ? body.get("kv").get(key).asInt() : null; } + private static JsonNode createAttrBodyWithoutTs(String kvJson) { + return JacksonUtil.toJsonNode("{\"kv\":" + kvJson + "}"); + } + private static JsonNode createAttrBody(long ts, String kvJson) { return JacksonUtil.toJsonNode("{\"ts\":" + ts + ",\"kv\":" + kvJson + "}"); }