@ -99,9 +99,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private final SslHandler sslHandler ;
private final SslHandler sslHandler ;
private final ConcurrentMap < MqttTopicMatcher , Integer > mqttQoSMap ;
private final ConcurrentMap < MqttTopicMatcher , Integer > mqttQoSMap ;
private volatile SessionInfoProto sessionInfo ;
private final DeviceSessionCtx deviceSessionCtx ;
private volatile InetSocketAddress address ;
private volatile InetSocketAddress address ;
private volatile DeviceSessionCtx deviceSessionCtx ;
private volatile GatewaySessionHandler gatewaySessionHandler ;
private volatile GatewaySessionHandler gatewaySessionHandler ;
MqttTransportHandler ( MqttTransportContext context , SslHandler sslHandler ) {
MqttTransportHandler ( MqttTransportContext context , SslHandler sslHandler ) {
@ -152,7 +151,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
case PINGREQ :
case PINGREQ :
if ( checkConnected ( ctx , msg ) ) {
if ( checkConnected ( ctx , msg ) ) {
ctx . writeAndFlush ( new MqttMessage ( new MqttFixedHeader ( PINGRESP , false , AT_MOST_ONCE , false , 0 ) ) ) ;
ctx . writeAndFlush ( new MqttMessage ( new MqttFixedHeader ( PINGRESP , false , AT_MOST_ONCE , false , 0 ) ) ) ;
transportService . reportActivity ( sessionInfo ) ;
transportService . reportActivity ( deviceSe ssionCtx . getS essionInfo( ) ) ;
}
}
break ;
break ;
case DISCONNECT :
case DISCONNECT :
@ -176,7 +175,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
if ( topicName . startsWith ( MqttTopics . BASE_GATEWAY_API_TOPIC ) ) {
if ( topicName . startsWith ( MqttTopics . BASE_GATEWAY_API_TOPIC ) ) {
if ( gatewaySessionHandler ! = null ) {
if ( gatewaySessionHandler ! = null ) {
handleGatewayPublishMsg ( topicName , msgId , mqttMsg ) ;
handleGatewayPublishMsg ( topicName , msgId , mqttMsg ) ;
transportService . reportActivity ( sessionInfo ) ;
transportService . reportActivity ( deviceSe ssionCtx . getS essionInfo( ) ) ;
}
}
} else {
} else {
processDevicePublish ( ctx , mqttMsg , topicName , msgId ) ;
processDevicePublish ( ctx , mqttMsg , topicName , msgId ) ;
@ -215,26 +214,26 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private void processDevicePublish ( ChannelHandlerContext ctx , MqttPublishMessage mqttMsg , String topicName , int msgId ) {
private void processDevicePublish ( ChannelHandlerContext ctx , MqttPublishMessage mqttMsg , String topicName , int msgId ) {
try {
try {
if ( topicName . equals ( MqttTopics . DEVICE_TELEMETRY_TOPIC ) ) {
if ( deviceSessionCtx . isDeviceTelemetryTopic ( topicName ) ) {
TransportProtos . PostTelemetryMsg postTelemetryMsg = adaptor . convertToPostTelemetry ( deviceSessionCtx , mqttMsg ) ;
TransportProtos . PostTelemetryMsg postTelemetryMsg = adaptor . convertToPostTelemetry ( deviceSessionCtx , mqttMsg ) ;
transportService . process ( sessionInfo , postTelemetryMsg , getPubAckCallback ( ctx , msgId , postTelemetryMsg ) ) ;
transportService . process ( deviceSe ssionCtx . getS essionInfo( ) , postTelemetryMsg , getPubAckCallback ( ctx , msgId , postTelemetryMsg ) ) ;
} else if ( topicName . equals ( MqttTopics . DEVICE_ATTRIBUTES_TOPIC ) ) {
} else if ( deviceSessionCtx . isDeviceAttributesTopic ( topicName ) ) {
TransportProtos . PostAttributeMsg postAttributeMsg = adaptor . convertToPostAttributes ( deviceSessionCtx , mqttMsg ) ;
TransportProtos . PostAttributeMsg postAttributeMsg = adaptor . convertToPostAttributes ( deviceSessionCtx , mqttMsg ) ;
transportService . process ( sessionInfo , postAttributeMsg , getPubAckCallback ( ctx , msgId , postAttributeMsg ) ) ;
transportService . process ( deviceSe ssionCtx . getS essionInfo( ) , postAttributeMsg , getPubAckCallback ( ctx , msgId , postAttributeMsg ) ) ;
} else if ( topicName . startsWith ( MqttTopics . DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX ) ) {
} else if ( topicName . startsWith ( MqttTopics . DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX ) ) {
TransportProtos . GetAttributeRequestMsg getAttributeMsg = adaptor . convertToGetAttributes ( deviceSessionCtx , mqttMsg ) ;
TransportProtos . GetAttributeRequestMsg getAttributeMsg = adaptor . convertToGetAttributes ( deviceSessionCtx , mqttMsg ) ;
transportService . process ( sessionInfo , getAttributeMsg , getPubAckCallback ( ctx , msgId , getAttributeMsg ) ) ;
transportService . process ( deviceSe ssionCtx . getS essionInfo( ) , getAttributeMsg , getPubAckCallback ( ctx , msgId , getAttributeMsg ) ) ;
} else if ( topicName . startsWith ( MqttTopics . DEVICE_RPC_RESPONSE_TOPIC ) ) {
} else if ( topicName . startsWith ( MqttTopics . DEVICE_RPC_RESPONSE_TOPIC ) ) {
TransportProtos . ToDeviceRpcResponseMsg rpcResponseMsg = adaptor . convertToDeviceRpcResponse ( deviceSessionCtx , mqttMsg ) ;
TransportProtos . ToDeviceRpcResponseMsg rpcResponseMsg = adaptor . convertToDeviceRpcResponse ( deviceSessionCtx , mqttMsg ) ;
transportService . process ( sessionInfo , rpcResponseMsg , getPubAckCallback ( ctx , msgId , rpcResponseMsg ) ) ;
transportService . process ( deviceSe ssionCtx . getS essionInfo( ) , rpcResponseMsg , getPubAckCallback ( ctx , msgId , rpcResponseMsg ) ) ;
} else if ( topicName . startsWith ( MqttTopics . DEVICE_RPC_REQUESTS_TOPIC ) ) {
} else if ( topicName . startsWith ( MqttTopics . DEVICE_RPC_REQUESTS_TOPIC ) ) {
TransportProtos . ToServerRpcRequestMsg rpcRequestMsg = adaptor . convertToServerRpcRequest ( deviceSessionCtx , mqttMsg ) ;
TransportProtos . ToServerRpcRequestMsg rpcRequestMsg = adaptor . convertToServerRpcRequest ( deviceSessionCtx , mqttMsg ) ;
transportService . process ( sessionInfo , rpcRequestMsg , getPubAckCallback ( ctx , msgId , rpcRequestMsg ) ) ;
transportService . process ( deviceSe ssionCtx . getS essionInfo( ) , rpcRequestMsg , getPubAckCallback ( ctx , msgId , rpcRequestMsg ) ) ;
} else if ( topicName . equals ( MqttTopics . DEVICE_CLAIM_TOPIC ) ) {
} else if ( topicName . equals ( MqttTopics . DEVICE_CLAIM_TOPIC ) ) {
TransportProtos . ClaimDeviceMsg claimDeviceMsg = adaptor . convertToClaimDevice ( deviceSessionCtx , mqttMsg ) ;
TransportProtos . ClaimDeviceMsg claimDeviceMsg = adaptor . convertToClaimDevice ( deviceSessionCtx , mqttMsg ) ;
transportService . process ( sessionInfo , claimDeviceMsg , getPubAckCallback ( ctx , msgId , claimDeviceMsg ) ) ;
transportService . process ( deviceSe ssionCtx . getS essionInfo( ) , claimDeviceMsg , getPubAckCallback ( ctx , msgId , claimDeviceMsg ) ) ;
} else {
} else {
transportService . reportActivity ( sessionInfo ) ;
transportService . reportActivity ( deviceSe ssionCtx . getS essionInfo( ) ) ;
}
}
} catch ( AdaptorException e ) {
} catch ( AdaptorException e ) {
log . warn ( "[{}] Failed to process publish msg [{}][{}]" , sessionId , topicName , msgId , e ) ;
log . warn ( "[{}] Failed to process publish msg [{}][{}]" , sessionId , topicName , msgId , e ) ;
@ -274,13 +273,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
try {
try {
switch ( topic ) {
switch ( topic ) {
case MqttTopics . DEVICE_ATTRIBUTES_TOPIC : {
case MqttTopics . DEVICE_ATTRIBUTES_TOPIC : {
transportService . process ( sessionInfo , TransportProtos . SubscribeToAttributeUpdatesMsg . newBuilder ( ) . build ( ) , null ) ;
transportService . process ( deviceSe ssionCtx . getS essionInfo( ) , TransportProtos . SubscribeToAttributeUpdatesMsg . newBuilder ( ) . build ( ) , null ) ;
registerSubQoS ( topic , grantedQoSList , reqQoS ) ;
registerSubQoS ( topic , grantedQoSList , reqQoS ) ;
activityReported = true ;
activityReported = true ;
break ;
break ;
}
}
case MqttTopics . DEVICE_RPC_REQUESTS_SUB_TOPIC : {
case MqttTopics . DEVICE_RPC_REQUESTS_SUB_TOPIC : {
transportService . process ( sessionInfo , TransportProtos . SubscribeToRPCMsg . newBuilder ( ) . build ( ) , null ) ;
transportService . process ( deviceSe ssionCtx . getS essionInfo( ) , TransportProtos . SubscribeToRPCMsg . newBuilder ( ) . build ( ) , null ) ;
registerSubQoS ( topic , grantedQoSList , reqQoS ) ;
registerSubQoS ( topic , grantedQoSList , reqQoS ) ;
activityReported = true ;
activityReported = true ;
break ;
break ;
@ -303,7 +302,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
}
}
}
}
if ( ! activityReported ) {
if ( ! activityReported ) {
transportService . reportActivity ( sessionInfo ) ;
transportService . reportActivity ( deviceSe ssionCtx . getS essionInfo( ) ) ;
}
}
ctx . writeAndFlush ( createSubAckMessage ( mqttMsg . variableHeader ( ) . messageId ( ) , grantedQoSList ) ) ;
ctx . writeAndFlush ( createSubAckMessage ( mqttMsg . variableHeader ( ) . messageId ( ) , grantedQoSList ) ) ;
}
}
@ -324,12 +323,14 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
try {
try {
switch ( topicName ) {
switch ( topicName ) {
case MqttTopics . DEVICE_ATTRIBUTES_TOPIC : {
case MqttTopics . DEVICE_ATTRIBUTES_TOPIC : {
transportService . process ( sessionInfo , TransportProtos . SubscribeToAttributeUpdatesMsg . newBuilder ( ) . setUnsubscribe ( true ) . build ( ) , null ) ;
transportService . process ( deviceSessionCtx . getSessionInfo ( ) ,
TransportProtos . SubscribeToAttributeUpdatesMsg . newBuilder ( ) . setUnsubscribe ( true ) . build ( ) , null ) ;
activityReported = true ;
activityReported = true ;
break ;
break ;
}
}
case MqttTopics . DEVICE_RPC_REQUESTS_SUB_TOPIC : {
case MqttTopics . DEVICE_RPC_REQUESTS_SUB_TOPIC : {
transportService . process ( sessionInfo , TransportProtos . SubscribeToRPCMsg . newBuilder ( ) . setUnsubscribe ( true ) . build ( ) , null ) ;
transportService . process ( deviceSessionCtx . getSessionInfo ( ) ,
TransportProtos . SubscribeToRPCMsg . newBuilder ( ) . setUnsubscribe ( true ) . build ( ) , null ) ;
activityReported = true ;
activityReported = true ;
break ;
break ;
}
}
@ -339,7 +340,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
}
}
}
}
if ( ! activityReported ) {
if ( ! activityReported ) {
transportService . reportActivity ( sessionInfo ) ;
transportService . reportActivity ( deviceSe ssionCtx . getS essionInfo( ) ) ;
}
}
ctx . writeAndFlush ( createUnSubAckMessage ( mqttMsg . variableHeader ( ) . messageId ( ) ) ) ;
ctx . writeAndFlush ( createUnSubAckMessage ( mqttMsg . variableHeader ( ) . messageId ( ) ) ) ;
}
}
@ -499,8 +500,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
private void doDisconnect ( ) {
private void doDisconnect ( ) {
if ( deviceSessionCtx . isConnected ( ) ) {
if ( deviceSessionCtx . isConnected ( ) ) {
transportService . process ( sessionInfo , DefaultTransportService . getSessionEventMsg ( SessionEvent . CLOSED ) , null ) ;
transportService . process ( deviceSe ssionCtx . getS essionInfo( ) , DefaultTransportService . getSessionEventMsg ( SessionEvent . CLOSED ) , null ) ;
transportService . deregisterSession ( sessionInfo ) ;
transportService . deregisterSession ( deviceSe ssionCtx . getS essionInfo( ) ) ;
if ( gatewaySessionHandler ! = null ) {
if ( gatewaySessionHandler ! = null ) {
gatewaySessionHandler . onGatewayDisconnect ( ) ;
gatewaySessionHandler . onGatewayDisconnect ( ) ;
}
}
@ -515,11 +516,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
} else {
} else {
deviceSessionCtx . setDeviceInfo ( msg . getDeviceInfo ( ) ) ;
deviceSessionCtx . setDeviceInfo ( msg . getDeviceInfo ( ) ) ;
deviceSessionCtx . setDeviceProfile ( msg . getDeviceProfile ( ) ) ;
deviceSessionCtx . setDeviceProfile ( msg . getDeviceProfile ( ) ) ;
sessionInfo = SessionInfoCreator . create ( msg , context , sessionId ) ;
deviceSessionCtx . setSessionInfo ( SessionInfoCreator . create ( msg , context , sessionId ) ) ;
transportService . process ( sessionInfo , DefaultTransportService . getSessionEventMsg ( SessionEvent . OPEN ) , new TransportServiceCallback < Void > ( ) {
transportService . process ( deviceSe ssionCtx . getS essionInfo( ) , DefaultTransportService . getSessionEventMsg ( SessionEvent . OPEN ) , new TransportServiceCallback < Void > ( ) {
@Override
@Override
public void onSuccess ( Void msg ) {
public void onSuccess ( Void msg ) {
transportService . registerAsyncSession ( sessionInfo , MqttTransportHandler . this ) ;
transportService . registerAsyncSession ( deviceSe ssionCtx . getS essionInfo( ) , MqttTransportHandler . this ) ;
checkGatewaySession ( ) ;
checkGatewaySession ( ) ;
ctx . writeAndFlush ( createMqttConnAckMsg ( CONNECTION_ACCEPTED ) ) ;
ctx . writeAndFlush ( createMqttConnAckMsg ( CONNECTION_ACCEPTED ) ) ;
log . info ( "[{}] Client connected!" , sessionId ) ;
log . info ( "[{}] Client connected!" , sessionId ) ;
@ -581,7 +582,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
@Override
@Override
public void onProfileUpdate ( DeviceProfile deviceProfile ) {
public void onProfileUpdate ( DeviceProfile deviceProfile ) {
deviceSessionCtx . getDeviceInfo ( ) . setDeviceType ( deviceProfile . getName ( ) ) ;
deviceSessionCtx . onProfileUpdate ( deviceProfile ) ;
sessionInfo = SessionInfoProto . newBuilder ( ) . mergeFrom ( sessionInfo ) . setDeviceType ( deviceProfile . getName ( ) ) . build ( ) ;
}
}
}
}