|
|
|
@ -22,6 +22,8 @@ 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; |
|
|
|
@ -30,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; |
|
|
|
@ -49,6 +56,7 @@ import org.thingsboard.server.common.transport.TransportService; |
|
|
|
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; |
|
|
|
@ -72,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; |
|
|
|
@ -84,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}") |
|
|
|
@ -115,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); |
|
|
|
@ -284,6 +297,47 @@ public class SnmpTransportService implements TbTransportService { |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
/* SNMP notifications handler */ |
|
|
|
@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; |
|
|
|
} |
|
|
|
|
|
|
|
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()) |
|
|
|
.build(); |
|
|
|
|
|
|
|
responseProcessingExecutor.execute(() -> { |
|
|
|
processResponse(sessionContext, List.of(pdu), requestContext); |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
private void processResponse(DeviceSessionContext sessionContext, List<PDU> response, RequestContext requestContext) { |
|
|
|
ResponseProcessor responseProcessor = responseProcessors.get(requestContext.getCommunicationSpec()); |
|
|
|
if (responseProcessor == null) return; |
|
|
|
@ -342,6 +396,15 @@ public class SnmpTransportService implements TbTransportService { |
|
|
|
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("TRAP") |
|
|
|
.setParams(JsonConverter.toJson(responseData)) |
|
|
|
.build(); |
|
|
|
transportService.process(sessionContext.getSessionInfo(), toServerRpcRequestMsg, null); |
|
|
|
}); |
|
|
|
} |
|
|
|
|
|
|
|
private void reportActivity(TransportProtos.SessionInfoProto sessionInfo) { |
|
|
|
|