@ -15,6 +15,11 @@
* /
* /
package org.thingsboard.server.transport.snmp.service ;
package org.thingsboard.server.transport.snmp.service ;
import com.google.common.util.concurrent.Futures ;
import com.google.common.util.concurrent.ListenableFuture ;
import com.google.common.util.concurrent.ListenableScheduledFuture ;
import com.google.common.util.concurrent.ListeningScheduledExecutorService ;
import com.google.common.util.concurrent.MoreExecutors ;
import com.google.gson.JsonElement ;
import com.google.gson.JsonElement ;
import com.google.gson.JsonObject ;
import com.google.gson.JsonObject ;
import lombok.Builder ;
import lombok.Builder ;
@ -32,7 +37,7 @@ import org.snmp4j.mp.MPv3;
import org.snmp4j.security.SecurityModels ;
import org.snmp4j.security.SecurityModels ;
import org.snmp4j.security.SecurityProtocols ;
import org.snmp4j.security.SecurityProtocols ;
import org.snmp4j.security.USM ;
import org.snmp4j.security.USM ;
import org.snmp4j.smi.Address ;
import org.snmp4j.smi.Ip Address ;
import org.snmp4j.smi.OctetString ;
import org.snmp4j.smi.OctetString ;
import org.snmp4j.smi.TcpAddress ;
import org.snmp4j.smi.TcpAddress ;
import org.snmp4j.smi.UdpAddress ;
import org.snmp4j.smi.UdpAddress ;
@ -44,6 +49,7 @@ import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service ;
import org.springframework.stereotype.Service ;
import org.thingsboard.common.util.ThingsBoardExecutors ;
import org.thingsboard.common.util.ThingsBoardExecutors ;
import org.thingsboard.common.util.ThingsBoardThreadFactory ;
import org.thingsboard.common.util.ThingsBoardThreadFactory ;
import org.thingsboard.server.common.adaptor.JsonConverter ;
import org.thingsboard.server.common.data.DataConstants ;
import org.thingsboard.server.common.data.DataConstants ;
import org.thingsboard.server.common.data.TbTransportService ;
import org.thingsboard.server.common.data.TbTransportService ;
import org.thingsboard.server.common.data.kv.DataType ;
import org.thingsboard.server.common.data.kv.DataType ;
@ -53,11 +59,11 @@ 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.RepeatingQueryingSnmpCommunicationConfig ;
import org.thingsboard.server.common.data.transport.snmp.config.SnmpCommunicationConfig ;
import org.thingsboard.server.common.data.transport.snmp.config.SnmpCommunicationConfig ;
import org.thingsboard.server.common.transport.TransportService ;
import org.thingsboard.server.common.transport.TransportService ;
import org.thingsboard.server.common.adaptor.JsonConverter ;
import org.thingsboard.server.gen.transport.TransportProtos ;
import org.thingsboard.server.gen.transport.TransportProtos ;
import org.thingsboard.server.queue.util.TbSnmpTransportComponent ;
import org.thingsboard.server.queue.util.TbSnmpTransportComponent ;
import org.thingsboard.server.transport.snmp.SnmpTransportContext ;
import org.thingsboard.server.transport.snmp.SnmpTransportContext ;
import org.thingsboard.server.transport.snmp.session.DeviceSessionContext ;
import org.thingsboard.server.transport.snmp.session.DeviceSessionContext ;
import org.thingsboard.server.transport.snmp.session.ScheduledTask ;
import javax.annotation.PostConstruct ;
import javax.annotation.PostConstruct ;
import javax.annotation.PreDestroy ;
import javax.annotation.PreDestroy ;
@ -71,8 +77,6 @@ import java.util.Map;
import java.util.Optional ;
import java.util.Optional ;
import java.util.concurrent.ExecutorService ;
import java.util.concurrent.ExecutorService ;
import java.util.concurrent.Executors ;
import java.util.concurrent.Executors ;
import java.util.concurrent.ScheduledExecutorService ;
import java.util.concurrent.ScheduledFuture ;
import java.util.concurrent.TimeUnit ;
import java.util.concurrent.TimeUnit ;
import java.util.stream.Collectors ;
import java.util.stream.Collectors ;
@ -80,6 +84,7 @@ import java.util.stream.Collectors;
@Service
@Service
@Slf4j
@Slf4j
@RequiredArgsConstructor
@RequiredArgsConstructor
@SuppressWarnings ( "UnstableApiUsage" )
public class SnmpTransportService implements TbTransportService , CommandResponder {
public class SnmpTransportService implements TbTransportService , CommandResponder {
private final TransportService transportService ;
private final TransportService transportService ;
private final PduService pduService ;
private final PduService pduService ;
@ -88,23 +93,27 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
@Getter
@Getter
private Snmp snmp ;
private Snmp snmp ;
private ScheduledExecutorService queryingExecuto r;
private ListeningScheduledExecutorService schedule r;
private ExecutorService r esponseProcessingE xecutor;
private ExecutorService executor ;
private final Map < SnmpCommunicationSpec , ResponseDataMapper > responseDataMappers = new EnumMap < > ( SnmpCommunicationSpec . class ) ;
private final Map < SnmpCommunicationSpec , ResponseDataMapper > responseDataMappers = new EnumMap < > ( SnmpCommunicationSpec . class ) ;
private final Map < SnmpCommunicationSpec , ResponseProcessor > responseProcessors = new EnumMap < > ( SnmpCommunicationSpec . class ) ;
private final Map < SnmpCommunicationSpec , ResponseProcessor > responseProcessors = new EnumMap < > ( SnmpCommunicationSpec . class ) ;
@Value ( "${transport.snmp.bind_port:1620}" )
@Value ( "${transport.snmp.bind_port:1620}" )
private Integer snmpBindPort ;
private Integer snmpBindPort ;
@Value ( "${transport.snmp.response_processing.parallelism_level}" )
@Value ( "${transport.snmp.response_processing.parallelism_level:4}" )
private Integer responseProcessingParallelismLevel ;
private int responseProcessingThreadPoolSize ;
@Value ( "${transport.snmp.scheduler_thread_pool_size:4}" )
private int schedulerThreadPoolSize ;
@Value ( "${transport.snmp.underlying_protocol}" )
@Value ( "${transport.snmp.underlying_protocol}" )
private String snmpUnderlyingProtocol ;
private String snmpUnderlyingProtocol ;
@Value ( "${transport.snmp.request_chunk_delay_ms:100}" )
private int requestChunkDelayMs ;
@PostConstruct
@PostConstruct
private void init ( ) throws IOException {
private void init ( ) throws IOException {
queryingExecutor = Executors . newScheduledThreadPool ( Runtime . getRuntime ( ) . availableProcessors ( ) , ThingsBoardThreadFactory . forName ( "snmp-querying" ) ) ;
scheduler = MoreExecutors . listeningDecorator ( Executors . newScheduledThreadPool ( schedulerThreadPoolSize , ThingsBoardThreadFactory . forName ( "snmp-querying" ) ) ) ;
r esponseProcessingE xecutor = ThingsBoardExecutors . newWorkStealingPool ( responseProcessingParallelismLevel , "snmp-response-processing" ) ;
executor = ThingsBoardExecutors . newWorkStealingPool ( responseProcessingThreadPoolSize , "snmp-response-processing" ) ;
initializeSnmp ( ) ;
initializeSnmp ( ) ;
configureResponseDataMappers ( ) ;
configureResponseDataMappers ( ) ;
@ -115,11 +124,11 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
@PreDestroy
@PreDestroy
public void stop ( ) {
public void stop ( ) {
if ( queryingExecuto r ! = null ) {
if ( schedule r ! = null ) {
queryingExecuto r. shutdownNow ( ) ;
schedule r. shutdownNow ( ) ;
}
}
if ( r esponseProcessingE xecutor ! = null ) {
if ( executor ! = null ) {
r esponseProcessingE xecutor. shutdownNow ( ) ;
executor . shutdownNow ( ) ;
}
}
}
}
@ -144,38 +153,39 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
}
}
public void createQueryingTasks ( DeviceSessionContext sessionContext ) {
public void createQueryingTasks ( DeviceSessionContext sessionContext ) {
List < ScheduledFuture < ? > > queryingTasks = sessionContext . getProfileTransportConfiguration ( ) . getCommunicationConfigs ( ) . stream ( )
sessionContext . getProfileTransportConfiguration ( ) . getCommunicationConfigs ( ) . stream ( )
. filter ( communicationConfig - > communicationConfig instanceof RepeatingQueryingSnmpCommunicationConfig )
. filter ( communicationConfig - > communicationConfig instanceof RepeatingQueryingSnmpCommunicationConfig )
. map ( config - > {
. forEach ( config - > {
RepeatingQueryingSnmpCommunicationConfig repeatingCommunicationConfig = ( RepeatingQueryingSnmpCommunicationConfig ) config ;
RepeatingQueryingSnmpCommunicationConfig repeatingCommunicationConfig = ( RepeatingQueryingSnmpCommunicationConfig ) config ;
Long queryingFrequency = repeatingCommunicationConfig . getQueryingFrequencyMs ( ) ;
Long queryingFrequency = repeatingCommunicationConfig . getQueryingFrequencyMs ( ) ;
return queryingExecutor . scheduleWithFixedDelay ( ( ) - > {
ScheduledTask scheduledTask = new ScheduledTask ( ) ;
scheduledTask . init ( ( ) - > {
try {
try {
if ( sessionContext . isActive ( ) ) {
if ( sessionContext . isActive ( ) ) {
sendRequest ( sessionContext , repeatingCommunicationConfig ) ;
return sendRequest ( sessionContext , repeatingCommunicationConfig ) ;
}
}
} catch ( Exception e ) {
} catch ( Exception e ) {
log . error ( "Failed to send SNMP request for device {}: {}" , sessionContext . getDeviceId ( ) , e . toString ( ) ) ;
log . error ( "Failed to send SNMP request for device {}: {}" , sessionContext . getDeviceId ( ) , e . toString ( ) ) ;
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , config . getSpec ( ) . getLabel ( ) , e ) ;
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , config . getSpec ( ) . getLabel ( ) , e ) ;
}
}
} , queryingFrequency , queryingFrequency , TimeUnit . MILLISECONDS ) ;
return Futures . immediateVoidFuture ( ) ;
} )
} , queryingFrequency , scheduler ) ;
. collect ( Collectors . toList ( ) ) ;
sessionContext . getQueryingTasks ( ) . add ( scheduledTask ) ;
sessionContext . getQueryingTasks ( ) . addAll ( queryingTasks ) ;
} ) ;
}
}
public void cancelQueryingTasks ( DeviceSessionContext sessionContext ) {
public void cancelQueryingTasks ( DeviceSessionContext sessionContext ) {
sessionContext . getQueryingTasks ( ) . forEach ( task - > task . cancel ( true ) ) ;
sessionContext . getQueryingTasks ( ) . forEach ( ScheduledTask : : cancel ) ;
sessionContext . getQueryingTasks ( ) . clear ( ) ;
sessionContext . getQueryingTasks ( ) . clear ( ) ;
}
}
private void sendRequest ( DeviceSessionContext sessionContext , SnmpCommunicationConfig communicationConfig ) {
private ListenableFuture < Void > sendRequest ( DeviceSessionContext sessionContext , SnmpCommunicationConfig communicationConfig ) {
sendRequest ( sessionContext , communicationConfig , Collections . emptyMap ( ) ) ;
return sendRequest ( sessionContext , communicationConfig , Collections . emptyMap ( ) ) ;
}
}
private void sendRequest ( DeviceSessionContext sessionContext , SnmpCommunicationConfig communicationConfig , Map < String , String > values ) {
private ListenableFuture < Void > sendRequest ( DeviceSessionContext sessionContext , SnmpCommunicationConfig communicationConfig , Map < String , String > values ) {
List < PDU > request = pduService . createPdus ( sessionContext , communicationConfig , values ) ;
List < PDU > request = pduService . createPdus ( sessionContext , communicationConfig , values ) ;
RequestContext requestContext = RequestContext . builder ( )
RequestContext requestContext = RequestContext . builder ( )
. communicationSpec ( communicationConfig . getSpec ( ) )
. communicationSpec ( communicationConfig . getSpec ( ) )
@ -183,19 +193,40 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
. responseMappings ( communicationConfig . getAllMappings ( ) )
. responseMappings ( communicationConfig . getAllMappings ( ) )
. requestSize ( request . size ( ) )
. requestSize ( request . size ( ) )
. build ( ) ;
. build ( ) ;
sendRequest ( sessionContext , request , requestContext ) ;
return sendRequest ( sessionContext , request , requestContext ) ;
}
}
private void sendRequest ( DeviceSessionContext sessionContext , List < PDU > request , RequestContext requestContext ) {
private ListenableFuture < Void > sendRequest ( DeviceSessionContext sessionContext , List < PDU > request , RequestContext requestContext ) {
for ( PDU pdu : request ) {
if ( request . size ( ) < = 1 | | requestChunkDelayMs = = 0 ) {
log . debug ( "Executing SNMP request for device {} with {} variable bindings" , sessionContext . getDeviceId ( ) , pdu . size ( ) ) ;
for ( PDU pdu : request ) {
try {
sendPdu ( pdu , requestContext , sessionContext ) ;
snmp . send ( pdu , sessionContext . getTarget ( ) , requestContext , sessionContext ) ;
}
} catch ( IOException e ) {
return Futures . immediateVoidFuture ( ) ;
log . error ( "Failed to send SNMP request to device {}: {}" , sessionContext . getDeviceId ( ) , e . toString ( ) ) ;
}
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , requestContext . getCommunicationSpec ( ) . getLabel ( ) , e ) ;
List < ListenableFuture < ? > > futures = new ArrayList < > ( ) ;
for ( int i = 0 , delay = 0 ; i < request . size ( ) ; i + + , delay + = requestChunkDelayMs ) {
PDU pdu = request . get ( i ) ;
if ( delay = = 0 ) {
sendPdu ( pdu , requestContext , sessionContext ) ;
} else {
ListenableScheduledFuture < ? > future = scheduler . schedule ( ( ) - > {
sendPdu ( pdu , requestContext , sessionContext ) ;
} , delay , TimeUnit . MILLISECONDS ) ;
futures . add ( future ) ;
}
}
}
}
return Futures . whenAllComplete ( futures ) . call ( ( ) - > null , MoreExecutors . directExecutor ( ) ) ;
}
private void sendPdu ( PDU pdu , RequestContext requestContext , DeviceSessionContext sessionContext ) {
log . debug ( "[{}] Sending SNMP request with {} variable bindings to {}" , sessionContext . getDeviceId ( ) , pdu . size ( ) , sessionContext . getTarget ( ) . getAddress ( ) ) ;
try {
snmp . send ( pdu , sessionContext . getTarget ( ) , requestContext , sessionContext ) ;
} catch ( Exception e ) {
log . error ( "[{}] Failed to send SNMP request" , sessionContext . getDeviceId ( ) , e ) ;
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , requestContext . getCommunicationSpec ( ) . getLabel ( ) , e ) ;
}
}
}
public void onAttributeUpdate ( DeviceSessionContext sessionContext , TransportProtos . AttributeUpdateNotificationMsg attributeUpdateNotification ) {
public void onAttributeUpdate ( DeviceSessionContext sessionContext , TransportProtos . AttributeUpdateNotificationMsg attributeUpdateNotification ) {
@ -251,21 +282,19 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
( ( Snmp ) event . getSource ( ) ) . cancel ( event . getRequest ( ) , sessionContext ) ;
( ( Snmp ) event . getSource ( ) ) . cancel ( event . getRequest ( ) , sessionContext ) ;
RequestContext requestContext = ( RequestContext ) event . getUserObject ( ) ;
RequestContext requestContext = ( RequestContext ) event . getUserObject ( ) ;
if ( event . getError ( ) ! = null ) {
if ( event . getError ( ) ! = null ) {
log . warn ( "SNMP response error: {}" , event . getError ( ) . toString ( ) ) ;
log . warn ( "[{}] SNMP response error: {}" , sessionContext . getDeviceId ( ) , event . getError ( ) . toString ( ) ) ;
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , requestContext . getCommunicationSpec ( ) . getLabel ( ) , new RuntimeException ( event . getError ( ) ) ) ;
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , requestContext . getCommunicationSpec ( ) . getLabel ( ) , new RuntimeException ( event . getError ( ) ) ) ;
return ;
return ;
}
}
PDU responsePdu = event . getResponse ( ) ;
PDU responsePdu = event . getResponse ( ) ;
if ( log . isTraceEnabled ( ) ) {
log . trace ( "[{}] Received PDU: {}" , sessionContext . getDeviceId ( ) , responsePdu ) ;
log . trace ( "Received PDU for device {}: {}" , sessionContext . getDeviceId ( ) , responsePdu ) ;
}
List < PDU > response ;
List < PDU > response ;
if ( requestContext . getRequestSize ( ) = = 1 ) {
if ( requestContext . getRequestSize ( ) = = 1 ) {
if ( responsePdu = = null ) {
if ( responsePdu = = null ) {
log . debug ( "No response from SNMP device {}, requestId: {}" , sessionContext . getDeviceId ( ) , event . getRequest ( ) . getRequestID ( ) ) ;
if ( requestContext . getMethod ( ) = = SnmpMethod . GET ) {
if ( requestContext . getMethod ( ) = = SnmpMethod . GET ) {
log . debug ( "[{}][{}] Empty response from device" , sessionContext . getDeviceId ( ) , event . getRequest ( ) . getRequestID ( ) ) ;
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , requestContext . getCommunicationSpec ( ) . getLabel ( ) , new RuntimeException ( "No response from device" ) ) ;
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , requestContext . getCommunicationSpec ( ) . getLabel ( ) , new RuntimeException ( "No response from device" ) ) ;
}
}
return ;
return ;
@ -281,14 +310,14 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
response . add ( responsePart ) ;
response . add ( responsePart ) ;
}
}
}
}
log . debug ( "All response parts are collected for request to device {} " , sessionContext . getDeviceId ( ) ) ;
log . debug ( "[{}] All {} response parts are collected for request" , sessionContext . getDeviceId ( ) , responseParts . size ( ) ) ;
} else {
} else {
log . trace ( "Awaiting other response parts for request to device {} " , sessionContext . getDeviceId ( ) ) ;
log . trace ( "[{}] Awaiting other response parts for request" , sessionContext . getDeviceId ( ) ) ;
return ;
return ;
}
}
}
}
r esponseProcessingE xecutor. execute ( ( ) - > {
executor . execute ( ( ) - > {
try {
try {
processResponse ( sessionContext , response , requestContext ) ;
processResponse ( sessionContext , response , requestContext ) ;
} catch ( Exception e ) {
} catch ( Exception e ) {
@ -298,24 +327,31 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
}
}
/ *
/ *
* SNMP notifications handler
* SNMP notifications handler
*
*
* TODO : add check for host uniqueness when saving device ( for backward compatibility - only for the ones using from - device RPC requests )
* 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 ,
* 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
* due to load - balancing of requests from devices : session might not be on this instance
* * /
* * /
@Override
@Override
public void processPdu ( CommandResponderEvent event ) {
public void processPdu ( CommandResponderEvent event ) {
Address sourceAddress = event . getPeerAddress ( ) ;
IpAddress sourceAddress = ( IpAddress ) event . getPeerAddress ( ) ;
DeviceSessionContext sessionContext = transportContext . getSessions ( ) . stream ( )
List < DeviceSessionContext > sessions = transportContext . getSessions ( ) . stream ( )
. filter ( session - > session . getTarget ( ) . getAddress ( ) . equals ( sourceAddress ) )
. filter ( session - > ( ( IpAddress ) session . getTarget ( ) . getAddress ( ) ) . getInetAddress ( ) . equals ( sourceAddress . getInetAddress ( ) ) )
. findFirst ( ) . orElse ( null ) ;
. collect ( Collectors . toList ( ) ) ;
if ( sessionContext = = null ) {
if ( sessions . isEmpty ( ) ) {
log . warn ( "SNMP TRAP processing failed: couldn't find device session for address {}" , sourceAddress ) ;
log . warn ( "Couldn't find device session for SNMP TRAP for address {}" , sourceAddress ) ;
return ;
} else if ( sessions . size ( ) > 1 ) {
for ( DeviceSessionContext sessionContext : sessions ) {
transportService . errorEvent ( sessionContext . getTenantId ( ) , sessionContext . getDeviceId ( ) , SnmpCommunicationSpec . TO_SERVER_RPC_REQUEST . getLabel ( ) ,
new IllegalStateException ( "Found multiple devices for host " + sourceAddress . getInetAddress ( ) . getHostAddress ( ) ) ) ;
}
return ;
return ;
}
}
DeviceSessionContext sessionContext = sessions . get ( 0 ) ;
try {
try {
processIncomingTrap ( sessionContext , event ) ;
processIncomingTrap ( sessionContext , event ) ;
} catch ( Throwable e ) {
} catch ( Throwable e ) {
@ -327,11 +363,11 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
private void processIncomingTrap ( DeviceSessionContext sessionContext , CommandResponderEvent event ) {
private void processIncomingTrap ( DeviceSessionContext sessionContext , CommandResponderEvent event ) {
PDU pdu = event . getPDU ( ) ;
PDU pdu = event . getPDU ( ) ;
if ( pdu = = null ) {
if ( pdu = = null ) {
log . warn ( "Got empty trap from device {} " , sessionContext . getDeviceId ( ) ) ;
log . warn ( "[{}] Received empty SNMP trap " , sessionContext . getDeviceId ( ) ) ;
throw new IllegalArgumentException ( "Received TRAP with no data" ) ;
throw new IllegalArgumentException ( "Received TRAP with no data" ) ;
}
}
log . debug ( "Processing SNMP trap from device {} (PDU : {} }" , sessionContext . getDeviceId ( ) , pdu ) ;
log . debug ( "[{}] Processing SNMP trap: {}" , sessionContext . getDeviceId ( ) , pdu ) ;
SnmpCommunicationConfig communicationConfig = sessionContext . getProfileTransportConfiguration ( ) . getCommunicationConfigs ( ) . stream ( )
SnmpCommunicationConfig communicationConfig = sessionContext . getProfileTransportConfiguration ( ) . getCommunicationConfigs ( ) . stream ( )
. filter ( config - > config . getSpec ( ) = = SnmpCommunicationSpec . TO_SERVER_RPC_REQUEST ) . findFirst ( )
. filter ( config - > config . getSpec ( ) = = SnmpCommunicationSpec . TO_SERVER_RPC_REQUEST ) . findFirst ( )
. orElseThrow ( ( ) - > new IllegalArgumentException ( "No config found for to-server RPC requests" ) ) ;
. orElseThrow ( ( ) - > new IllegalArgumentException ( "No config found for to-server RPC requests" ) ) ;
@ -341,7 +377,7 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
. method ( SnmpMethod . TRAP )
. method ( SnmpMethod . TRAP )
. build ( ) ;
. build ( ) ;
r esponseProcessingE xecutor. execute ( ( ) - > {
executor . execute ( ( ) - > {
processResponse ( sessionContext , List . of ( pdu ) , requestContext ) ;
processResponse ( sessionContext , List . of ( pdu ) , requestContext ) ;
} ) ;
} ) ;
}
}
@ -352,7 +388,7 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
JsonObject responseData = responseDataMappers . get ( requestContext . getCommunicationSpec ( ) ) . map ( response , requestContext ) ;
JsonObject responseData = responseDataMappers . get ( requestContext . getCommunicationSpec ( ) ) . map ( response , requestContext ) ;
if ( responseData . size ( ) = = 0 ) {
if ( responseData . size ( ) = = 0 ) {
log . warn ( "No values in the SNMP response for device {} " , sessionContext . getDeviceId ( ) ) ;
log . warn ( "[{}] No values in the response" , sessionContext . getDeviceId ( ) ) ;
throw new IllegalArgumentException ( "No values in the response" ) ;
throw new IllegalArgumentException ( "No values in the response" ) ;
}
}
@ -428,11 +464,11 @@ public class SnmpTransportService implements TbTransportService, CommandResponde
@PreDestroy
@PreDestroy
public void shutdown ( ) {
public void shutdown ( ) {
log . info ( "Stopping SNMP transport!" ) ;
log . info ( "Stopping SNMP transport!" ) ;
if ( queryingExecuto r ! = null ) {
if ( schedule r ! = null ) {
queryingExecuto r. shutdownNow ( ) ;
schedule r. shutdownNow ( ) ;
}
}
if ( r esponseProcessingE xecutor ! = null ) {
if ( executor ! = null ) {
r esponseProcessingE xecutor. shutdownNow ( ) ;
executor . shutdownNow ( ) ;
}
}
if ( snmp ! = null ) {
if ( snmp ! = null ) {
try {
try {