diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 9864d89048..205948290c 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -86,6 +86,8 @@ public final class EdgeGrpcSession implements Closeable { private static final ReentrantLock downlinkMsgLock = new ReentrantLock(); + private static final int MAX_DOWNLINK_ATTEMPTS = 10; // max number of attemps to send downlink message if edge connected + private static final String QUEUE_START_TS_ATTR_KEY = "queueStartTs"; private final UUID sessionId; @@ -252,9 +254,9 @@ public final class EdgeGrpcSession implements Closeable { try { if (msg.getSuccess()) { sessionState.getPendingMsgsMap().remove(msg.getDownlinkMsgId()); - log.debug("[{}] Msg has been processed successfully! {}", edge.getRoutingKey(), msg); + log.debug("[{}] Msg has been processed successfully!Msd Id: [{}], Msg: {}", edge.getRoutingKey(), msg.getDownlinkMsgId(), msg); } else { - log.error("[{}] Msg processing failed! Error msg: {}", edge.getRoutingKey(), msg.getErrorMsg()); + log.error("[{}] Msg processing failed! Msd Id: [{}], Error msg: {}", edge.getRoutingKey(), msg.getDownlinkMsgId(), msg.getErrorMsg()); } if (sessionState.getPendingMsgsMap().isEmpty()) { log.debug("[{}] Pending msgs map is empty. Stopping current iteration", edge.getRoutingKey()); @@ -392,17 +394,17 @@ public final class EdgeGrpcSession implements Closeable { sessionState.setSendDownlinkMsgsFuture(SettableFuture.create()); sessionState.getPendingMsgsMap().clear(); downlinkMsgsPack.forEach(msg -> sessionState.getPendingMsgsMap().put(msg.getDownlinkMsgId(), msg)); - scheduleDownlinkMsgsPackSend(true); + scheduleDownlinkMsgsPackSend(1); return sessionState.getSendDownlinkMsgsFuture(); } - private void scheduleDownlinkMsgsPackSend(boolean firstRun) { + private void scheduleDownlinkMsgsPackSend(int attempt) { Runnable sendDownlinkMsgsTask = () -> { try { if (isConnected() && sessionState.getPendingMsgsMap().values().size() > 0) { List copy = new ArrayList<>(sessionState.getPendingMsgsMap().values()); - if (!firstRun) { - log.warn("[{}] Failed to deliver the batch: {}", this.sessionId, copy); + if (attempt > 1) { + log.warn("[{}] Failed to deliver the batch: {}, attempt: {}", this.sessionId, copy, attempt); } log.trace("[{}] [{}] downlink msg(s) are going to be send.", this.sessionId, copy.size()); for (DownlinkMsg downlinkMsg : copy) { @@ -410,7 +412,13 @@ public final class EdgeGrpcSession implements Closeable { .setDownlinkMsg(downlinkMsg) .build()); } - scheduleDownlinkMsgsPackSend(false); + if (attempt < MAX_DOWNLINK_ATTEMPTS) { + scheduleDownlinkMsgsPackSend(attempt + 1); + } else { + log.warn("[{}] Failed to deliver the batch after {} attempts. Next messages are going to be discarded {}", + this.sessionId, MAX_DOWNLINK_ATTEMPTS, copy); + sessionState.getSendDownlinkMsgsFuture().set(null); + } } else { sessionState.getSendDownlinkMsgsFuture().set(null); } @@ -419,7 +427,7 @@ public final class EdgeGrpcSession implements Closeable { } }; - if (firstRun) { + if (attempt == 1) { sendDownlinkExecutorService.submit(sendDownlinkMsgsTask); } else { sessionState.setScheduledSendDownlinkTask( diff --git a/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java index b33e354e89..74629f5a9e 100644 --- a/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java @@ -25,6 +25,7 @@ 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.gen.edge.v1.AttributeDeleteMsg; +import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.v1.EntityDataProto; import org.thingsboard.server.gen.transport.TransportProtos; @@ -192,4 +193,38 @@ abstract public class BaseTelemetryEdgeTest extends AbstractEdgeTest { return false; } + @Test + public void testTimeseriesDeliveryFailuresForever_deliverOnlyDeviceUpdateMsgs() throws Exception { + int numberOfMsgsToSend = 100; + + edgeImitator.setRandomFailuresOnTimeseriesDownlink(true); + // imitator will generate failure in 100% of timeseries cases + edgeImitator.setFailureProbability(100); + + edgeImitator.expectMessageAmount(numberOfMsgsToSend); + Device device = findDeviceByName("Edge Device 1"); + for (int idx = 1; idx <= numberOfMsgsToSend; idx++) { + String timeseriesData = "{\"data\":{\"idx\":" + idx + "},\"ts\":" + System.currentTimeMillis() + "}"; + JsonNode timeseriesEntityData = mapper.readTree(timeseriesData); + EdgeEvent failedEdgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.TIMESERIES_UPDATED, + device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData); + edgeEventService.saveAsync(failedEdgeEvent).get(); + + EdgeEvent successEdgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.UPDATED, + device.getId().getId(), EdgeEventType.DEVICE, null); + edgeEventService.saveAsync(successEdgeEvent).get(); + + clusterService.onEdgeEventUpdate(tenantId, edge.getId()); + } + + Assert.assertTrue(edgeImitator.waitForMessages(120)); + + List allTelemetryMsgs = edgeImitator.findAllMessagesByType(EntityDataProto.class); + Assert.assertTrue(allTelemetryMsgs.isEmpty()); + + List deviceUpdateMsgs = edgeImitator.findAllMessagesByType(DeviceUpdateMsg.class); + Assert.assertEquals(numberOfMsgsToSend, deviceUpdateMsgs.size()); + + edgeImitator.setRandomFailuresOnTimeseriesDownlink(false); + } }