@ -42,7 +42,6 @@ import org.thingsboard.server.common.data.limit.LimitedApi;
import org.thingsboard.server.common.data.notification.rule.trigger.EdgeCommunicationFailureTrigger ;
import org.thingsboard.server.common.data.page.PageData ;
import org.thingsboard.server.common.data.page.PageLink ;
import org.thingsboard.server.common.data.page.SortOrder ;
import org.thingsboard.server.common.data.page.TimePageLink ;
import org.thingsboard.server.common.msg.edge.EdgeEventUpdateMsg ;
import org.thingsboard.server.gen.edge.v1.AlarmCommentUpdateMsg ;
@ -292,11 +291,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 +322,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 +458,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 +482,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 +492,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 +539,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 +554,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 +576,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 +585,7 @@ public abstract class EdgeGrpcSession implements Closeable {
Futures . addCallback ( updateFuture , new FutureCallback < > ( ) {
@Override
public void onSuccess ( @Nullable List < Long > list ) {
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 +600,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 +631,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,17 +640,17 @@ 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 +756,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 +908,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 . ge tT enantId( ) ) . 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 ) ;
}