|
|
|
@ -401,7 +401,8 @@ public class EdgeGrpcSession implements EdgeSession { |
|
|
|
if (state.isConnected() && !pageData.getData().isEmpty()) { |
|
|
|
if (fetcher instanceof GeneralEdgeEventFetcher) { |
|
|
|
long queueSize = pageData.getTotalElements() - ((long) pageLink.getPageSize() * pageLink.getPage()); |
|
|
|
ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.setDownlinkMsgsLag(edge.getTenantId(), edge.getId(), queueSize)); |
|
|
|
ctx.getStatsCounterService().ifPresent(statsCounterService -> |
|
|
|
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_LAG, tenantId, edge.getId(), queueSize)); |
|
|
|
} |
|
|
|
log.trace("[{}][{}][{}] event(s) are going to be processed.", tenantId, edge.getId(), pageData.getData().size()); |
|
|
|
List<DownlinkMsg> downlinkMsgsPack = downlinkMessageMapper.convertToDownlinkMsgsPack(state, pageData.getData()); |
|
|
|
@ -496,7 +497,8 @@ public class EdgeGrpcSession implements EdgeSession { |
|
|
|
ctx.getRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(getTenantId()).edgeId(getEdgeId()) |
|
|
|
.customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg) |
|
|
|
.error("Failed to deliver messages after " + MAX_DOWNLINK_ATTEMPTS + " attempts").build()); |
|
|
|
ctx.getStatsCounterService().ifPresent(statsCounterService -> statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_PERMANENTLY_FAILED, edge.getTenantId(), getEdgeId(), copy.size())); |
|
|
|
ctx.getStatsCounterService().ifPresent(statsCounterService -> |
|
|
|
statsCounterService.recordEvent(EdgeStatsKey.DOWNLINK_MSGS_PERMANENTLY_FAILED, edge.getTenantId(), getEdgeId(), copy.size())); |
|
|
|
stopCurrentSendDownlinkMsgsTask(false); |
|
|
|
} |
|
|
|
} else { |
|
|
|
|