diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java index c255ce93ea..a7a1282dc0 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java @@ -173,6 +173,8 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { edgeImitator = new EdgeImitator("localhost", 7070, edge.getRoutingKey(), edge.getSecret()); edgeImitator.expectMessageAmount(9); edgeImitator.connect(); + + testReceivedInitialData(); } @After @@ -188,9 +190,7 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { } @Test - public void test() throws Exception { - testReceivedInitialData(); - + public void generalTest() throws Exception { testDevices(); testAssets(); @@ -218,6 +218,55 @@ abstract public class BaseEdgeTest extends AbstractControllerTest { testRpcCall(); } + @Test + public void testTimeseriesWithFailures() throws Exception { + log.info("Testing timeseries with failures"); + + int numberOfTimeseriesToSend = 1000; + + edgeImitator.setRandomFailuresOnTimeseriesDownlink(true); + // imitator will generate failure in 5% of cases + edgeImitator.setFailureProbability(5.0); + + edgeImitator.expectMessageAmount(numberOfTimeseriesToSend); + Device device = findDeviceByName("Edge Device 1"); + for (int idx = 1; idx <= numberOfTimeseriesToSend; idx++) { + String timeseriesData = "{\"data\":{\"idx\":" + idx + "},\"ts\":" + System.currentTimeMillis() + "}"; + JsonNode timeseriesEntityData = mapper.readTree(timeseriesData); + EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.TIMESERIES_UPDATED, + device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData); + edgeEventService.saveAsync(edgeEvent); + clusterService.onEdgeEventUpdate(tenantId, edge.getId()); + } + + Assert.assertTrue(edgeImitator.waitForMessages(60)); + + List allTelemetryMsgs = edgeImitator.findAllMessagesByType(EntityDataProto.class); + Assert.assertEquals(numberOfTimeseriesToSend, allTelemetryMsgs.size()); + + for (int idx = 1; idx <= numberOfTimeseriesToSend; idx++) { + Assert.assertTrue(isIdxExistsInTheDownlinkList(idx, allTelemetryMsgs)); + } + + edgeImitator.setRandomFailuresOnTimeseriesDownlink(false); + log.info("Timeseries with failures tested successfully"); + } + + private boolean isIdxExistsInTheDownlinkList(int idx, List allTelemetryMsgs) { + for (EntityDataProto proto : allTelemetryMsgs) { + TransportProtos.PostTelemetryMsg postTelemetryMsg = proto.getPostTelemetryMsg(); + Assert.assertEquals(1, postTelemetryMsg.getTsKvListCount()); + TransportProtos.TsKvListProto tsKvListProto = postTelemetryMsg.getTsKvList(0); + Assert.assertEquals(1, tsKvListProto.getKvCount()); + TransportProtos.KeyValueProto keyValueProto = tsKvListProto.getKv(0); + Assert.assertEquals("idx", keyValueProto.getKey()); + if (keyValueProto.getLongV() == idx) { + return true; + } + } + return false; + } + private Device findDeviceByName(String deviceName) throws Exception { List edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?", new TypeReference>() { diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java index 864f2e6a52..4988f7d1cf 100644 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java +++ b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java @@ -21,11 +21,11 @@ import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; import com.google.protobuf.AbstractMessage; import lombok.Getter; +import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.checkerframework.checker.nullness.qual.Nullable; import org.thingsboard.edge.rpc.EdgeGrpcClient; import org.thingsboard.edge.rpc.EdgeRpcClient; -import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.gen.edge.AlarmUpdateMsg; import org.thingsboard.server.gen.edge.AssetUpdateMsg; import org.thingsboard.server.gen.edge.CustomerUpdateMsg; @@ -55,6 +55,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Optional; import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -74,6 +75,11 @@ public class EdgeImitator { private CountDownLatch responsesLatch; private List> ignoredTypes; + @Setter + private boolean randomFailuresOnTimeseriesDownlink = false; + @Setter + private double failureProbability = 0.0; + @Getter private EdgeConfiguration configuration; @Getter @@ -208,7 +214,15 @@ public class EdgeImitator { } if (downlinkMsg.getEntityDataCount() > 0) { for (EntityDataProto entityData : downlinkMsg.getEntityDataList()) { - result.add(saveDownlinkMsg(entityData)); + if (randomFailuresOnTimeseriesDownlink) { + if (getRandomBoolean()) { + result.add(Futures.immediateFailedFuture(new RuntimeException("Random failure"))); + } else { + result.add(saveDownlinkMsg(entityData)); + } + } else { + result.add(saveDownlinkMsg(entityData)); + } } } if (downlinkMsg.getEntityViewUpdateMsgCount() > 0) { @@ -254,6 +268,11 @@ public class EdgeImitator { return Futures.allAsList(result); } + private boolean getRandomBoolean() { + double randomValue = ThreadLocalRandom.current().nextDouble() * 100; + return randomValue <= this.failureProbability; + } + private ListenableFuture saveDownlinkMsg(AbstractMessage message) { if (!ignoredTypes.contains(message.getClass())) { try { @@ -271,8 +290,8 @@ public class EdgeImitator { return waitForMessages(5); } - public boolean waitForMessages(int timeout) throws InterruptedException { - return messagesLatch.await(timeout, TimeUnit.SECONDS); + public boolean waitForMessages(int timeoutInSeconds) throws InterruptedException { + return messagesLatch.await(timeoutInSeconds, TimeUnit.SECONDS); } public void expectMessageAmount(int messageAmount) {