Browse Source

Attemp to send downlink 10 times max in case edge connected - discard all other attempts and ack failed messages

pull/7478/head
Volodymyr Babak 4 years ago
parent
commit
546b477c34
  1. 24
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  2. 35
      application/src/test/java/org/thingsboard/server/edge/BaseTelemetryEdgeTest.java

24
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 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 static final String QUEUE_START_TS_ATTR_KEY = "queueStartTs";
private final UUID sessionId; private final UUID sessionId;
@ -252,9 +254,9 @@ public final class EdgeGrpcSession implements Closeable {
try { try {
if (msg.getSuccess()) { if (msg.getSuccess()) {
sessionState.getPendingMsgsMap().remove(msg.getDownlinkMsgId()); 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 { } 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()) { if (sessionState.getPendingMsgsMap().isEmpty()) {
log.debug("[{}] Pending msgs map is empty. Stopping current iteration", edge.getRoutingKey()); 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.setSendDownlinkMsgsFuture(SettableFuture.create());
sessionState.getPendingMsgsMap().clear(); sessionState.getPendingMsgsMap().clear();
downlinkMsgsPack.forEach(msg -> sessionState.getPendingMsgsMap().put(msg.getDownlinkMsgId(), msg)); downlinkMsgsPack.forEach(msg -> sessionState.getPendingMsgsMap().put(msg.getDownlinkMsgId(), msg));
scheduleDownlinkMsgsPackSend(true); scheduleDownlinkMsgsPackSend(1);
return sessionState.getSendDownlinkMsgsFuture(); return sessionState.getSendDownlinkMsgsFuture();
} }
private void scheduleDownlinkMsgsPackSend(boolean firstRun) { private void scheduleDownlinkMsgsPackSend(int attempt) {
Runnable sendDownlinkMsgsTask = () -> { Runnable sendDownlinkMsgsTask = () -> {
try { try {
if (isConnected() && sessionState.getPendingMsgsMap().values().size() > 0) { if (isConnected() && sessionState.getPendingMsgsMap().values().size() > 0) {
List<DownlinkMsg> copy = new ArrayList<>(sessionState.getPendingMsgsMap().values()); List<DownlinkMsg> copy = new ArrayList<>(sessionState.getPendingMsgsMap().values());
if (!firstRun) { if (attempt > 1) {
log.warn("[{}] Failed to deliver the batch: {}", this.sessionId, copy); 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()); log.trace("[{}] [{}] downlink msg(s) are going to be send.", this.sessionId, copy.size());
for (DownlinkMsg downlinkMsg : copy) { for (DownlinkMsg downlinkMsg : copy) {
@ -410,7 +412,13 @@ public final class EdgeGrpcSession implements Closeable {
.setDownlinkMsg(downlinkMsg) .setDownlinkMsg(downlinkMsg)
.build()); .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 { } else {
sessionState.getSendDownlinkMsgsFuture().set(null); sessionState.getSendDownlinkMsgsFuture().set(null);
} }
@ -419,7 +427,7 @@ public final class EdgeGrpcSession implements Closeable {
} }
}; };
if (firstRun) { if (attempt == 1) {
sendDownlinkExecutorService.submit(sendDownlinkMsgsTask); sendDownlinkExecutorService.submit(sendDownlinkMsgsTask);
} else { } else {
sessionState.setScheduledSendDownlinkTask( sessionState.setScheduledSendDownlinkTask(

35
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.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.gen.edge.v1.AttributeDeleteMsg; 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.edge.v1.EntityDataProto;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
@ -192,4 +193,38 @@ abstract public class BaseTelemetryEdgeTest extends AbstractEdgeTest {
return false; 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<EntityDataProto> allTelemetryMsgs = edgeImitator.findAllMessagesByType(EntityDataProto.class);
Assert.assertTrue(allTelemetryMsgs.isEmpty());
List<DeviceUpdateMsg> deviceUpdateMsgs = edgeImitator.findAllMessagesByType(DeviceUpdateMsg.class);
Assert.assertEquals(numberOfMsgsToSend, deviceUpdateMsgs.size());
edgeImitator.setRandomFailuresOnTimeseriesDownlink(false);
}
} }

Loading…
Cancel
Save