@ -17,10 +17,13 @@ package org.thingsboard.server.transport.snmp.service;
import com.google.gson.JsonElement ;
import com.google.gson.JsonObject ;
import lombok.Builder ;
import lombok.Data ;
import lombok.Getter ;
import lombok.RequiredArgsConstructor ;
import lombok.extern.slf4j.Slf4j ;
import org.snmp4j.CommandResponder ;
import org.snmp4j.CommandResponderEvent ;
import org.snmp4j.PDU ;
import org.snmp4j.Snmp ;
import org.snmp4j.TransportMapping ;
@ -29,10 +32,15 @@ import org.snmp4j.mp.MPv3;
import org.snmp4j.security.SecurityModels ;
import org.snmp4j.security.SecurityProtocols ;
import org.snmp4j.security.USM ;
import org.snmp4j.smi.Address ;
import org.snmp4j.smi.OctetString ;
import org.snmp4j.smi.TcpAddress ;
import org.snmp4j.smi.UdpAddress ;
import org.snmp4j.transport.DefaultTcpTransportMapping ;
import org.snmp4j.transport.DefaultUdpTransportMapping ;
import org.springframework.beans.factory.annotation.Autowired ;
import org.springframework.beans.factory.annotation.Value ;
import org.springframework.context.annotation.Lazy ;
import org.springframework.stereotype.Service ;
import org.thingsboard.common.util.ThingsBoardExecutors ;
import org.thingsboard.common.util.ThingsBoardThreadFactory ;
@ -45,15 +53,16 @@ import org.thingsboard.server.common.data.transport.snmp.SnmpMethod;
import org.thingsboard.server.common.data.transport.snmp.config.RepeatingQueryingSnmpCommunicationConfig ;
import org.thingsboard.server.common.data.transport.snmp.config.SnmpCommunicationConfig ;
import org.thingsboard.server.common.transport.TransportService ;
import org.thingsboard.server.common.transport.TransportServiceCallback ;
import org.thingsboard.server.common.transport.adaptor.JsonConverter ;
import org.thingsboard.server.gen.transport.TransportProtos ;
import org.thingsboard.server.queue.util.TbSnmpTransportComponent ;
import org.thingsboard.server.transport.snmp.SnmpTransportContext ;
import org.thingsboard.server.transport.snmp.session.DeviceSessionContext ;
import javax.annotation.PostConstruct ;
import javax.annotation.PreDestroy ;
import java.io.IOException ;
import java.util.ArrayList ;
import java.util.Arrays ;
import java.util.Collections ;
import java.util.EnumMap ;
@ -71,9 +80,11 @@ import java.util.stream.Collectors;
@Service
@Slf4j
@RequiredArgsConstructor
public class SnmpTransportService implements TbTransportService {
public class SnmpTransportService implements TbTransportService , CommandResponder {
private final TransportService transportService ;
private final PduService pduService ;
@Autowired @Lazy
private SnmpTransportContext transportContext ;
@Getter
private Snmp snmp ;
@ -83,6 +94,8 @@ public class SnmpTransportService implements TbTransportService {
private final Map < SnmpCommunicationSpec , ResponseDataMapper > responseDataMappers = new EnumMap < > ( SnmpCommunicationSpec . class ) ;
private final Map < SnmpCommunicationSpec , ResponseProcessor > responseProcessors = new EnumMap < > ( SnmpCommunicationSpec . class ) ;
@Value ( "${transport.snmp.bind_port:1620}" )
private Integer snmpBindPort ;
@Value ( "${transport.snmp.response_processing.parallelism_level}" )
private Integer responseProcessingParallelismLevel ;
@Value ( "${transport.snmp.underlying_protocol}" )
@ -114,15 +127,16 @@ public class SnmpTransportService implements TbTransportService {
TransportMapping < ? > transportMapping ;
switch ( snmpUnderlyingProtocol ) {
case "udp" :
transportMapping = new DefaultUdpTransportMapping ( ) ;
transportMapping = new DefaultUdpTransportMapping ( new UdpAddress ( snmpBindPort ) ) ;
break ;
case "tcp" :
transportMapping = new DefaultTcpTransportMapping ( ) ;
transportMapping = new DefaultTcpTransportMapping ( new TcpAddress ( snmpBindPort ) ) ;
break ;
default :
throw new IllegalArgumentException ( "Underlying protocol " + snmpUnderlyingProtocol + " for SNMP is not supported" ) ;
}
snmp = new Snmp ( transportMapping ) ;
snmp . addNotificationListener ( transportMapping , transportMapping . getListenAddress ( ) , this ) ;
snmp . listen ( ) ;
USM usm = new USM ( SecurityProtocols . getInstance ( ) , new OctetString ( MPv3 . createLocalEngineID ( ) ) , 0 ) ;
@ -143,6 +157,7 @@ public class SnmpTransportService implements TbTransportService {
}
} catch ( Exception e ) {
log . error ( "Failed to send SNMP request for device {}: {}" , sessionContext . getDeviceId ( ) , e . toString ( ) ) ;
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , config . getSpec ( ) . getLabel ( ) , e ) ;
}
} , queryingFrequency , queryingFrequency , TimeUnit . MILLISECONDS ) ;
} )
@ -161,18 +176,24 @@ public class SnmpTransportService implements TbTransportService {
}
private void sendRequest ( DeviceSessionContext sessionContext , SnmpCommunicationConfig communicationConfig , Map < String , String > values ) {
PDU request = pduService . createPdu ( sessionContext , communicationConfig , values ) ;
RequestInfo requestInfo = new RequestInfo ( communicationConfig . getSpec ( ) , communicationConfig . getAllMappings ( ) ) ;
sendRequest ( sessionContext , request , requestInfo ) ;
List < PDU > request = pduService . createPdus ( sessionContext , communicationConfig , values ) ;
RequestContext requestContext = RequestContext . builder ( )
. communicationSpec ( communicationConfig . getSpec ( ) )
. method ( communicationConfig . getMethod ( ) )
. responseMappings ( communicationConfig . getAllMappings ( ) )
. requestSize ( request . size ( ) )
. build ( ) ;
sendRequest ( sessionContext , request , requestContext ) ;
}
private void sendRequest ( DeviceSessionContext sessionContext , PDU request , RequestInfo requestInfo ) {
if ( request . size ( ) > 0 ) {
log . trace ( "Executing SNMP request for device {}. Variables bindings: {}" , sessionContext . getDeviceId ( ) , request . getVariableBindings ( ) ) ;
private void sendRequest ( DeviceSessionContext sessionContext , List < PDU > request , RequestContext requestContext ) {
for ( PDU pdu : request ) {
log . debug ( "Executing SNMP request for device {} with {} variable bindings ", sessionContext . getDeviceId ( ) , pdu . size ( ) ) ;
try {
snmp . send ( request , sessionContext . getTarget ( ) , requestInfo , sessionContext ) ;
snmp . send ( pdu , sessionContext . getTarget ( ) , requestContext , sessionContext ) ;
} catch ( IOException e ) {
log . error ( "Failed to send SNMP request to device {}: {}" , sessionContext . getDeviceId ( ) , e . toString ( ) ) ;
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , requestContext . getCommunicationSpec ( ) . getLabel ( ) , e ) ;
}
}
}
@ -215,51 +236,135 @@ public class SnmpTransportService implements TbTransportService {
DataType dataType = snmpMapping . getDataType ( ) ;
PDU request = pduService . createSingleVariablePdu ( sessionContext , snmpMethod , oid , value , dataType ) ;
RequestInfo requestInfo = new RequestInfo ( toDeviceRpcRequestMsg . getRequestId ( ) , communicationConfig . getSpec ( ) , communicationConfig . getAllMappings ( ) ) ;
sendRequest ( sessionContext , request , requestInfo ) ;
RequestContext requestContext = RequestContext . builder ( )
. requestId ( toDeviceRpcRequestMsg . getRequestId ( ) )
. communicationSpec ( communicationConfig . getSpec ( ) )
. method ( snmpMethod )
. responseMappings ( communicationConfig . getAllMappings ( ) )
. requestSize ( 1 )
. build ( ) ;
sendRequest ( sessionContext , List . of ( request ) , requestContext ) ;
}
public void processResponseEvent ( DeviceSessionContext sessionContext , ResponseEvent event ) {
( ( Snmp ) event . getSource ( ) ) . cancel ( event . getRequest ( ) , sessionContext ) ;
RequestContext requestContext = ( RequestContext ) event . getUserObject ( ) ;
if ( event . getError ( ) ! = null ) {
log . warn ( "SNMP response error: {}" , event . getError ( ) . toString ( ) ) ;
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , requestContext . getCommunicationSpec ( ) . getLabel ( ) , new RuntimeException ( event . getError ( ) ) ) ;
return ;
}
PDU response = event . getResponse ( ) ;
if ( response = = null ) {
log . debug ( "No response from SNMP device {}, requestId: {}" , sessionContext . getDeviceId ( ) , event . getRequest ( ) . getRequestID ( ) ) ;
PDU responsePdu = event . getResponse ( ) ;
if ( log . isTraceEnabled ( ) ) {
log . trace ( "Received PDU for device {}: {}" , sessionContext . getDeviceId ( ) , responsePdu ) ;
}
List < PDU > response ;
if ( requestContext . getRequestSize ( ) = = 1 ) {
if ( responsePdu = = null ) {
log . debug ( "No response from SNMP device {}, requestId: {}" , sessionContext . getDeviceId ( ) , event . getRequest ( ) . getRequestID ( ) ) ;
if ( requestContext . getMethod ( ) = = SnmpMethod . GET ) {
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , requestContext . getCommunicationSpec ( ) . getLabel ( ) , new RuntimeException ( "No response from device" ) ) ;
}
return ;
}
response = List . of ( responsePdu ) ;
} else {
List < PDU > responseParts = requestContext . getResponseParts ( ) ;
responseParts . add ( responsePdu ) ;
if ( responseParts . size ( ) = = requestContext . getRequestSize ( ) ) {
response = new ArrayList < > ( ) ;
for ( PDU responsePart : responseParts ) {
if ( responsePart ! = null ) {
response . add ( responsePart ) ;
}
}
log . debug ( "All response parts are collected for request to device {}" , sessionContext . getDeviceId ( ) ) ;
} else {
log . trace ( "Awaiting other response parts for request to device {}" , sessionContext . getDeviceId ( ) ) ;
return ;
}
}
responseProcessingExecutor . execute ( ( ) - > {
try {
processResponse ( sessionContext , response , requestContext ) ;
} catch ( Exception e ) {
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , requestContext . getCommunicationSpec ( ) . getLabel ( ) , e ) ;
}
} ) ;
}
/ *
* SNMP notifications handler
*
* TODO : add check for host uniqueness when saving device ( for backward compatibility - only for the ones using from - device RPC requests )
*
* NOTE : SNMP TRAPs support won ' t work properly when there is more than one SNMP transport ,
* due to load - balancing of requests from devices : session might not be on this instance
* * /
@Override
public void processPdu ( CommandResponderEvent event ) {
Address sourceAddress = event . getPeerAddress ( ) ;
DeviceSessionContext sessionContext = transportContext . getSessions ( ) . stream ( )
. filter ( session - > session . getTarget ( ) . getAddress ( ) . equals ( sourceAddress ) )
. findFirst ( ) . orElse ( null ) ;
if ( sessionContext = = null ) {
log . warn ( "SNMP TRAP processing failed: couldn't find device session for address {}" , sourceAddress ) ;
return ;
}
RequestInfo requestInfo = ( RequestInfo ) event . getUserObject ( ) ;
try {
processIncomingTrap ( sessionContext , event ) ;
} catch ( Throwable e ) {
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) ,
SnmpCommunicationSpec . TO_SERVER_RPC_REQUEST . getLabel ( ) , e ) ;
}
}
private void processIncomingTrap ( DeviceSessionContext sessionContext , CommandResponderEvent event ) {
PDU pdu = event . getPDU ( ) ;
if ( pdu = = null ) {
log . warn ( "Got empty trap from device {}" , sessionContext . getDeviceId ( ) ) ;
throw new IllegalArgumentException ( "Received TRAP with no data" ) ;
}
log . debug ( "Processing SNMP trap from device {} (PDU: {}}" , sessionContext . getDeviceId ( ) , pdu ) ;
SnmpCommunicationConfig communicationConfig = sessionContext . getProfileTransportConfiguration ( ) . getCommunicationConfigs ( ) . stream ( )
. filter ( config - > config . getSpec ( ) = = SnmpCommunicationSpec . TO_SERVER_RPC_REQUEST ) . findFirst ( )
. orElseThrow ( ( ) - > new IllegalArgumentException ( "No config found for to-server RPC requests" ) ) ;
RequestContext requestContext = RequestContext . builder ( )
. communicationSpec ( communicationConfig . getSpec ( ) )
. responseMappings ( communicationConfig . getAllMappings ( ) )
. method ( SnmpMethod . TRAP )
. build ( ) ;
responseProcessingExecutor . execute ( ( ) - > {
processResponse ( sessionContext , response , requestInfo ) ;
processResponse ( sessionContext , List . of ( pdu ) , requestContext ) ;
} ) ;
}
private void processResponse ( DeviceSessionContext sessionContext , PDU response , RequestInfo requestInfo ) {
ResponseProcessor responseProcessor = responseProcessors . get ( requestInfo . getCommunicationSpec ( ) ) ;
private void processResponse ( DeviceSessionContext sessionContext , List < PDU > response , RequestContext requestContext ) {
ResponseProcessor responseProcessor = responseProcessors . get ( requestContext . getCommunicationSpec ( ) ) ;
if ( responseProcessor = = null ) return ;
JsonObject responseData = responseDataMappers . get ( requestInfo . getCommunicationSpec ( ) ) . map ( response , requestInfo ) ;
if ( responseData . entrySet ( ) . isEmpty ( ) ) {
log . debug ( "No values is the SNMP response for device {}. Request id: {}" , sessionContext . getDeviceId ( ) , response . getRequestID ( ) ) ;
return ;
JsonObject responseData = responseDataMappers . get ( requestContext . getCommunicationSpec ( ) ) . map ( response , requestContext ) ;
if ( responseData . size ( ) = = 0 ) {
log . warn ( "No values in the SNMP response for device {}" , sessionContext . getDeviceId ( ) ) ;
throw new IllegalArgumentException ( "No values in the response" ) ;
}
responseProcessor . process ( responseData , requestInfo , sessionContext ) ;
responseProcessor . process ( responseData , requestContext , sessionContext ) ;
reportActivity ( sessionContext . getSessionInfo ( ) ) ;
}
private void configureResponseDataMappers ( ) {
responseDataMappers . put ( SnmpCommunicationSpec . TO_DEVICE_RPC_REQUEST , ( pdu , requestInfo ) - > {
responseDataMappers . put ( SnmpCommunicationSpec . TO_DEVICE_RPC_REQUEST , ( pdus , requestContext ) - > {
JsonObject responseData = new JsonObject ( ) ;
pduService . processPdu ( pdu ) . forEach ( ( oid , value ) - > {
requestInfo . getResponseMappings ( ) . stream ( )
pduService . processPdus ( pdus ) . forEach ( ( oid , value ) - > {
requestContext . getResponseMappings ( ) . stream ( )
. filter ( snmpMapping - > snmpMapping . getOid ( ) . equals ( oid . toDottedString ( ) ) )
. findFirst ( )
. ifPresent ( snmpMapping - > {
@ -269,8 +374,8 @@ public class SnmpTransportService implements TbTransportService {
return responseData ;
} ) ;
ResponseDataMapper defaultResponseDataMapper = ( pdu , requestInfo ) - > {
return pduService . processPdu ( pdu , requestInfo . getResponseMappings ( ) ) ;
ResponseDataMapper defaultResponseDataMapper = ( pdus , requestContext ) - > {
return pduService . processPdus ( pdus , requestContext . getResponseMappings ( ) ) ;
} ;
Arrays . stream ( SnmpCommunicationSpec . values ( ) )
. forEach ( communicationSpec - > {
@ -279,34 +384,39 @@ public class SnmpTransportService implements TbTransportService {
}
private void configureResponseProcessors ( ) {
responseProcessors . put ( SnmpCommunicationSpec . TELEMETRY_QUERYING , ( responseData , requestInfo , sessionContext ) - > {
responseProcessors . put ( SnmpCommunicationSpec . TELEMETRY_QUERYING , ( responseData , requestContext , sessionContext ) - > {
TransportProtos . PostTelemetryMsg postTelemetryMsg = JsonConverter . convertToTelemetryProto ( responseData ) ;
transportService . process ( sessionContext . getSessionInfo ( ) , postTelemetryMsg , null ) ;
log . debug ( "Posted telemetry for SNMP device {}: {}" , sessionContext . getDeviceId ( ) , responseData ) ;
} ) ;
responseProcessors . put ( SnmpCommunicationSpec . CLIENT_ATTRIBUTES_QUERYING , ( responseData , requestInfo , sessionContext ) - > {
responseProcessors . put ( SnmpCommunicationSpec . CLIENT_ATTRIBUTES_QUERYING , ( responseData , requestContext , sessionContext ) - > {
TransportProtos . PostAttributeMsg postAttributesMsg = JsonConverter . convertToAttributesProto ( responseData ) ;
transportService . process ( sessionContext . getSessionInfo ( ) , postAttributesMsg , null ) ;
log . debug ( "Posted attributes for SNMP device {}: {}" , sessionContext . getDeviceId ( ) , responseData ) ;
} ) ;
responseProcessors . put ( SnmpCommunicationSpec . TO_DEVICE_RPC_REQUEST , ( responseData , requestInfo , sessionContext ) - > {
responseProcessors . put ( SnmpCommunicationSpec . TO_DEVICE_RPC_REQUEST , ( responseData , requestContext , sessionContext ) - > {
TransportProtos . ToDeviceRpcResponseMsg rpcResponseMsg = TransportProtos . ToDeviceRpcResponseMsg . newBuilder ( )
. setRequestId ( requestInfo . getRequestId ( ) )
. setRequestId ( requestContext . getRequestId ( ) )
. setPayload ( JsonConverter . toJson ( responseData ) )
. build ( ) ;
transportService . process ( sessionContext . getSessionInfo ( ) , rpcResponseMsg , null ) ;
log . debug ( "Posted RPC response {} for device {}" , responseData , sessionContext . getDeviceId ( ) ) ;
} ) ;
responseProcessors . put ( SnmpCommunicationSpec . TO_SERVER_RPC_REQUEST , ( responseData , requestContext , sessionContext ) - > {
TransportProtos . ToServerRpcRequestMsg toServerRpcRequestMsg = TransportProtos . ToServerRpcRequestMsg . newBuilder ( )
. setRequestId ( 0 )
. setMethodName ( requestContext . getMethod ( ) . name ( ) )
. setParams ( JsonConverter . toJson ( responseData ) )
. build ( ) ;
transportService . process ( sessionContext . getSessionInfo ( ) , toServerRpcRequestMsg , null ) ;
} ) ;
}
private void reportActivity ( TransportProtos . SessionInfoProto sessionInfo ) {
transportService . process ( sessionInfo , TransportProtos . SubscriptionInfoProto . newBuilder ( )
. setAttributeSubscription ( true )
. setRpcSubscription ( true )
. setLastActivityTime ( System . currentTimeMillis ( ) )
. build ( ) , TransportServiceCallback . EMPTY ) ;
transportService . reportActivity ( sessionInfo ) ;
}
@ -335,29 +445,34 @@ public class SnmpTransportService implements TbTransportService {
}
@Data
private static class RequestInfo {
private Integer requestId ;
private SnmpCommunicationSpec communicationSpec ;
private List < SnmpMapping > responseMappings ;
private static class RequestContext {
private final Integer requestId ;
private final SnmpCommunicationSpec communicationSpec ;
private final SnmpMethod method ;
private final List < SnmpMapping > responseMappings ;
public RequestInfo ( Integer requestId , SnmpCommunicationSpec communicationSpec , List < SnmpMapping > responseMappings ) {
this . requestId = requestId ;
this . communicationSpec = communicationSpec ;
this . responseMappings = responseMappings ;
}
private final int requestSize ;
private List < PDU > responseParts ;
public RequestInfo ( SnmpCommunicationSpec communicationSpec , List < SnmpMapping > responseMappings ) {
@Builder
public RequestContext ( Integer requestId , SnmpCommunicationSpec communicationSpec , SnmpMethod method , List < SnmpMapping > responseMappings , int requestSize ) {
this . requestId = requestId ;
this . communicationSpec = communicationSpec ;
this . method = method ;
this . responseMappings = responseMappings ;
this . requestSize = requestSize ;
if ( requestSize > 1 ) {
this . responseParts = Collections . synchronizedList ( new ArrayList < > ( ) ) ;
}
}
}
private interface ResponseDataMapper {
JsonObject map ( PDU pdu , RequestInfo requestInfo ) ;
JsonObject map ( List < PDU > pdus , RequestContext requestContext ) ;
}
private interface ResponseProcessor {
void process ( JsonObject responseData , RequestInfo requestInfo , DeviceSessionContext sessionContext ) ;
void process ( JsonObject responseData , RequestContext requestContext , DeviceSessionContext sessionContext ) ;
}
}