|
|
|
@ -292,11 +292,11 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
|
|
|
|
protected void processEdgeEvents(EdgeEventFetcher fetcher, PageLink pageLink, SettableFuture<Pair<Long, Long>> result) { |
|
|
|
try { |
|
|
|
log.trace("[{}] Start processing edge events, fetcher = {}, pageLink = {}", sessionId, fetcher.getClass().getSimpleName(), pageLink); |
|
|
|
log.trace("[{}] Start processing edge events, fetcher = {}, pageLink = {}", edge.getId(), fetcher.getClass().getSimpleName(), pageLink); |
|
|
|
processHighPriorityEvents(); |
|
|
|
PageData<EdgeEvent> pageData = fetcher.fetchEdgeEvents(edge.getTenantId(), edge, pageLink); |
|
|
|
if (isConnected() && !pageData.getData().isEmpty()) { |
|
|
|
log.trace("[{}][{}][{}] event(s) are going to be processed.", tenantId, sessionId, pageData.getData().size()); |
|
|
|
log.trace("[{}][{}][{}] event(s) are going to be processed.", tenantId, edge.getId(), pageData.getData().size()); |
|
|
|
List<DownlinkMsg> downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData()); |
|
|
|
Futures.addCallback(sendDownlinkMsgsPack(downlinkMsgsPack), new FutureCallback<>() { |
|
|
|
@Override |
|
|
|
@ -323,16 +323,16 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onFailure(Throwable t) { |
|
|
|
log.error("[{}] Failed to send downlink msgs pack", sessionId, t); |
|
|
|
log.error("[{}] Failed to send downlink msgs pack", edge.getId(), t); |
|
|
|
result.setException(t); |
|
|
|
} |
|
|
|
}, ctx.getGrpcCallbackExecutorService()); |
|
|
|
} else { |
|
|
|
log.trace("[{}] no event(s) found. Stop processing edge events, fetcher = {}, pageLink = {}", sessionId, fetcher.getClass().getSimpleName(), pageLink); |
|
|
|
log.trace("[{}] no event(s) found. Stop processing edge events, fetcher = {}, pageLink = {}", edge.getId(), fetcher.getClass().getSimpleName(), pageLink); |
|
|
|
result.set(null); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.error("[{}] Failed to fetch edge events", sessionId, e); |
|
|
|
log.error("[{}] Failed to fetch edge events", edge.getId(), e); |
|
|
|
result.setException(e); |
|
|
|
} |
|
|
|
} |
|
|
|
@ -459,9 +459,9 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
ctx.getRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId) |
|
|
|
.edgeId(edge.getId()).customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(error).build()); |
|
|
|
} |
|
|
|
log.warn("[{}][{}] {}, attempt: {}", tenantId, sessionId, failureMsg, attempt); |
|
|
|
log.warn("[{}][{}] {}, attempt: {}", tenantId, edge.getId(), failureMsg, attempt); |
|
|
|
} |
|
|
|
log.trace("[{}][{}][{}] downlink msg(s) are going to be send.", tenantId, sessionId, copy.size()); |
|
|
|
log.trace("[{}][{}][{}] downlink msg(s) are going to be send.", tenantId, edge.getId(), copy.size()); |
|
|
|
for (DownlinkMsg downlinkMsg : copy) { |
|
|
|
if (clientMaxInboundMessageSize != 0 && downlinkMsg.getSerializedSize() > clientMaxInboundMessageSize) { |
|
|
|
String error = String.format("Client max inbound message size %s is exceeded. Please increase value of CLOUD_RPC_MAX_INBOUND_MESSAGE_SIZE " + |
|
|
|
@ -483,7 +483,7 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
} else { |
|
|
|
String failureMsg = String.format("Failed to deliver messages: %s", copy); |
|
|
|
log.warn("[{}][{}] Failed to deliver the batch after {} attempts. Next messages are going to be discarded {}", |
|
|
|
tenantId, sessionId, MAX_DOWNLINK_ATTEMPTS, copy); |
|
|
|
tenantId, edge.getId(), MAX_DOWNLINK_ATTEMPTS, copy); |
|
|
|
ctx.getRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId()) |
|
|
|
.customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg) |
|
|
|
.error("Failed to deliver messages after " + MAX_DOWNLINK_ATTEMPTS + " attempts").build()); |
|
|
|
@ -493,7 +493,7 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
stopCurrentSendDownlinkMsgsTask(false); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.warn("[{}][{}] Failed to send downlink msgs. Error msg {}", tenantId, sessionId, e.getMessage(), e); |
|
|
|
log.warn("[{}][{}] Failed to send downlink msgs. Error msg {}", tenantId, edge.getId(), e.getMessage(), e); |
|
|
|
stopCurrentSendDownlinkMsgsTask(true); |
|
|
|
} |
|
|
|
}; |
|
|
|
@ -540,7 +540,7 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
stopCurrentSendDownlinkMsgsTask(false); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.error("[{}][{}] Can't process downlink response message [{}]", tenantId, sessionId, msg, e); |
|
|
|
log.error("[{}][{}] Can't process downlink response message [{}]", tenantId, edge.getId(), msg, e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -555,12 +555,12 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
while ((event = highPriorityQueue.poll()) != null) { |
|
|
|
highPriorityEvents.add(event); |
|
|
|
} |
|
|
|
log.trace("[{}][{}] Sending high priority events {}", tenantId, sessionId, highPriorityEvents.size()); |
|
|
|
log.trace("[{}][{}] Sending high priority events {}", tenantId, edge.getId(), highPriorityEvents.size()); |
|
|
|
List<DownlinkMsg> downlinkMsgsPack = convertToDownlinkMsgsPack(highPriorityEvents); |
|
|
|
sendDownlinkMsgsPack(downlinkMsgsPack).get(); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.error("[{}] Failed to process high priority events", sessionId, e); |
|
|
|
log.error("[{}] Failed to process high priority events", edge.getId(), e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -577,7 +577,7 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
Integer.toUnsignedLong(ctx.getEdgeEventStorageSettings().getMaxReadRecordsCount()), |
|
|
|
ctx.getEdgeEventService()); |
|
|
|
log.trace("[{}][{}] starting processing edge events, previousStartTs = {}, previousStartSeqId = {}", |
|
|
|
tenantId, sessionId, previousStartTs, previousStartSeqId); |
|
|
|
tenantId, edge.getId(), previousStartTs, previousStartSeqId); |
|
|
|
Futures.addCallback(startProcessingEdgeEvents(fetcher), new FutureCallback<>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(@Nullable Pair<Long, Long> newStartTsAndSeqId) { |
|
|
|
@ -586,7 +586,7 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
Futures.addCallback(updateFuture, new FutureCallback<>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(@Nullable AttributesSaveResult saveResult) { |
|
|
|
log.debug("[{}][{}] queue offset was updated [{}]", tenantId, sessionId, newStartTsAndSeqId); |
|
|
|
log.debug("[{}][{}] queue offset was updated [{}]", tenantId, edge.getId(), newStartTsAndSeqId); |
|
|
|
boolean newEventsAvailable; |
|
|
|
if (fetcher.isSeqIdNewCycleStarted()) { |
|
|
|
newEventsAvailable = isNewEdgeEventsAvailable(); |
|
|
|
@ -601,28 +601,28 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onFailure(Throwable t) { |
|
|
|
log.error("[{}][{}] Failed to update queue offset [{}]", tenantId, sessionId, newStartTsAndSeqId, t); |
|
|
|
log.error("[{}][{}] Failed to update queue offset [{}]", tenantId, edge.getId(), newStartTsAndSeqId, t); |
|
|
|
result.setException(t); |
|
|
|
} |
|
|
|
}, ctx.getGrpcCallbackExecutorService()); |
|
|
|
} else { |
|
|
|
log.trace("[{}][{}] newStartTsAndSeqId is null. Skipping iteration without db update", tenantId, sessionId); |
|
|
|
log.trace("[{}][{}] newStartTsAndSeqId is null. Skipping iteration without db update", tenantId, edge.getId()); |
|
|
|
result.set(Boolean.FALSE); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onFailure(Throwable t) { |
|
|
|
log.error("[{}][{}] Failed to process events", tenantId, sessionId, t); |
|
|
|
log.error("[{}][{}] Failed to process events", tenantId, edge.getId(), t); |
|
|
|
result.setException(t); |
|
|
|
} |
|
|
|
}, ctx.getGrpcCallbackExecutorService()); |
|
|
|
} else { |
|
|
|
if (isSyncInProgress()) { |
|
|
|
log.trace("[{}][{}] edge sync is not completed yet. Skipping iteration", tenantId, sessionId); |
|
|
|
log.trace("[{}][{}] edge sync is not completed yet. Skipping iteration", tenantId, edge.getId()); |
|
|
|
result.set(Boolean.TRUE); |
|
|
|
} else { |
|
|
|
log.trace("[{}][{}] edge is not connected. Skipping iteration", tenantId, sessionId); |
|
|
|
log.trace("[{}][{}] edge is not connected. Skipping iteration", tenantId, edge.getId()); |
|
|
|
result.set(null); |
|
|
|
} |
|
|
|
} |
|
|
|
@ -632,7 +632,7 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
protected List<DownlinkMsg> convertToDownlinkMsgsPack(List<EdgeEvent> edgeEvents) { |
|
|
|
List<DownlinkMsg> result = new ArrayList<>(); |
|
|
|
for (EdgeEvent edgeEvent : edgeEvents) { |
|
|
|
log.trace("[{}][{}] converting edge event to downlink msg [{}]", tenantId, sessionId, edgeEvent); |
|
|
|
log.trace("[{}][{}] converting edge event to downlink msg [{}]", tenantId, edge.getId(), edgeEvent); |
|
|
|
DownlinkMsg downlinkMsg = null; |
|
|
|
try { |
|
|
|
switch (edgeEvent.getAction()) { |
|
|
|
@ -641,16 +641,16 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
ASSIGNED_TO_CUSTOMER, UNASSIGNED_FROM_CUSTOMER, ADDED_COMMENT, UPDATED_COMMENT, DELETED_COMMENT -> { |
|
|
|
downlinkMsg = convertEntityEventToDownlink(edgeEvent); |
|
|
|
if (downlinkMsg != null && downlinkMsg.getWidgetTypeUpdateMsgCount() > 0) { |
|
|
|
log.trace("[{}][{}] widgetTypeUpdateMsg message processed, downlinkMsgId = {}", tenantId, sessionId, downlinkMsg.getDownlinkMsgId()); |
|
|
|
log.trace("[{}][{}] widgetTypeUpdateMsg message processed, downlinkMsgId = {}", tenantId, edge.getId(), downlinkMsg.getDownlinkMsgId()); |
|
|
|
} else { |
|
|
|
log.trace("[{}][{}] entity message processed [{}]", tenantId, sessionId, downlinkMsg); |
|
|
|
log.trace("[{}][{}] entity message processed [{}]", tenantId, edge.getId(), downlinkMsg); |
|
|
|
} |
|
|
|
} |
|
|
|
case ATTRIBUTES_UPDATED, POST_ATTRIBUTES, ATTRIBUTES_DELETED, TIMESERIES_UPDATED -> downlinkMsg = ctx.getTelemetryProcessor().convertTelemetryEventToDownlink(edge, edgeEvent); |
|
|
|
default -> log.warn("[{}][{}] Unsupported action type [{}]", tenantId, sessionId, edgeEvent.getAction()); |
|
|
|
default -> log.warn("[{}][{}] Unsupported action type [{}]", tenantId, edge.getId(), edgeEvent.getAction()); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("[{}][{}] Exception during converting edge event to downlink msg", tenantId, sessionId, e); |
|
|
|
log.trace("[{}][{}] Exception during converting edge event to downlink msg", tenantId, edge.getId(), e); |
|
|
|
} |
|
|
|
if (downlinkMsg != null) { |
|
|
|
result.add(downlinkMsg); |
|
|
|
@ -757,19 +757,19 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
private void sendDownlinkMsg(ResponseMsg responseMsg) { |
|
|
|
if (isConnected()) { |
|
|
|
String responseMsgStr = StringUtils.truncate(responseMsg.toString(), 10000); |
|
|
|
log.trace("[{}][{}] Sending downlink msg [{}]", tenantId, sessionId, responseMsgStr); |
|
|
|
log.trace("[{}][{}] Sending downlink msg [{}]", tenantId, edge.getId(), responseMsgStr); |
|
|
|
downlinkMsgLock.lock(); |
|
|
|
String downlinkMsgStr = responseMsg.hasDownlinkMsg() ? String.valueOf(responseMsg.getDownlinkMsg().getDownlinkMsgId()) : responseMsgStr; |
|
|
|
try { |
|
|
|
outputStream.onNext(responseMsg); |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("[{}][{}] Failed to send downlink message [{}]", tenantId, sessionId, downlinkMsgStr, e); |
|
|
|
log.trace("[{}][{}] Failed to send downlink message [{}]", tenantId, edge.getId(), downlinkMsgStr, e); |
|
|
|
connected = false; |
|
|
|
sessionCloseListener.accept(edge, sessionId); |
|
|
|
} finally { |
|
|
|
downlinkMsgLock.unlock(); |
|
|
|
} |
|
|
|
log.trace("[{}][{}] downlink msg successfully sent [{}]", tenantId, sessionId, downlinkMsgStr); |
|
|
|
log.trace("[{}][{}] downlink msg successfully sent [{}]", tenantId, edge.getId(), downlinkMsgStr); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -909,8 +909,8 @@ public abstract class EdgeGrpcSession implements Closeable { |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
String failureMsg = String.format("Can't process uplink msg [%s] from edge", uplinkMsg); |
|
|
|
log.trace("[{}][{}] Can't process uplink msg [{}]", edge.getTenantId(), sessionId, uplinkMsg, e); |
|
|
|
ctx.getRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(edge.getTenantId()).edgeId(edge.getId()) |
|
|
|
log.trace("[{}][{}] Can't process uplink msg [{}]", tenantId, edge.getId(), uplinkMsg, e); |
|
|
|
ctx.getRuleProcessor().process(EdgeCommunicationFailureTrigger.builder().tenantId(tenantId).edgeId(edge.getId()) |
|
|
|
.customerId(edge.getCustomerId()).edgeName(edge.getName()).failureMsg(failureMsg).error(e.getMessage()).build()); |
|
|
|
return Futures.immediateFailedFuture(e); |
|
|
|
} |
|
|
|
|