From 0f078b0da2c28386974af64060e1289d8722b848 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Mon, 26 Jun 2023 15:04:31 +0300 Subject: [PATCH] SNMP traps support --- .../transport/snmp/SnmpCommunicationSpec.java | 4 +- .../snmp/config/SnmpCommunicationConfig.java | 3 +- ...rverRpcRequestSnmpCommunicationConfig.java | 29 ++++++++ .../snmp/service/SnmpTransportService.java | 69 ++++++++++++++++++- .../server/transport/snmp/SnmpTestV2.java | 25 +++++-- ui-ngx/src/app/shared/models/device.models.ts | 12 ++-- 6 files changed, 128 insertions(+), 14 deletions(-) create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/ToServerRpcRequestSnmpCommunicationConfig.java diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/SnmpCommunicationSpec.java b/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/SnmpCommunicationSpec.java index 1e34bb61ab..1e0e5a745b 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/SnmpCommunicationSpec.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/SnmpCommunicationSpec.java @@ -26,7 +26,9 @@ public enum SnmpCommunicationSpec { CLIENT_ATTRIBUTES_QUERYING("clientAttributesQuerying"), SHARED_ATTRIBUTES_SETTING("sharedAttributesSetting"), - TO_DEVICE_RPC_REQUEST("rpcRequest"); + TO_DEVICE_RPC_REQUEST("toDeviceRpcRequest"), + + TO_SERVER_RPC_REQUEST("toServerRpcRequest"); private final String label; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/SnmpCommunicationConfig.java b/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/SnmpCommunicationConfig.java index 36a7f64d3a..d09c9c147c 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/SnmpCommunicationConfig.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/SnmpCommunicationConfig.java @@ -37,7 +37,8 @@ import java.util.List; @Type(value = TelemetryQueryingSnmpCommunicationConfig.class, name = "TELEMETRY_QUERYING"), @Type(value = ClientAttributesQueryingSnmpCommunicationConfig.class, name = "CLIENT_ATTRIBUTES_QUERYING"), @Type(value = SharedAttributesSettingSnmpCommunicationConfig.class, name = "SHARED_ATTRIBUTES_SETTING"), - @Type(value = ToDeviceRpcRequestSnmpCommunicationConfig.class, name = "TO_DEVICE_RPC_REQUEST") + @Type(value = ToDeviceRpcRequestSnmpCommunicationConfig.class, name = "TO_DEVICE_RPC_REQUEST"), + @Type(value = ToServerRpcRequestSnmpCommunicationConfig.class, name = "TO_SERVER_RPC_REQUEST") }) public interface SnmpCommunicationConfig extends Serializable { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/ToServerRpcRequestSnmpCommunicationConfig.java b/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/ToServerRpcRequestSnmpCommunicationConfig.java new file mode 100644 index 0000000000..f8a7968b82 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/ToServerRpcRequestSnmpCommunicationConfig.java @@ -0,0 +1,29 @@ +/** + * Copyright © 2016-2023 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.transport.snmp.config; + +import org.thingsboard.server.common.data.transport.snmp.SnmpCommunicationSpec; + +public class ToServerRpcRequestSnmpCommunicationConfig extends MultipleMappingsSnmpCommunicationConfig { + + private static final long serialVersionUID = 4851028734093214L; + + @Override + public SnmpCommunicationSpec getSpec() { + return SnmpCommunicationSpec.TO_SERVER_RPC_REQUEST; + } + +} diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java index d14cafb2c2..2d507df099 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java @@ -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 responseDataMappers = new EnumMap<>(SnmpCommunicationSpec.class); private final Map 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 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) { diff --git a/common/transport/snmp/src/test/java/org/thingsboard/server/transport/snmp/SnmpTestV2.java b/common/transport/snmp/src/test/java/org/thingsboard/server/transport/snmp/SnmpTestV2.java index fb7c374d01..5367ef93cf 100644 --- a/common/transport/snmp/src/test/java/org/thingsboard/server/transport/snmp/SnmpTestV2.java +++ b/common/transport/snmp/src/test/java/org/thingsboard/server/transport/snmp/SnmpTestV2.java @@ -21,19 +21,36 @@ import java.util.Map; import java.util.Scanner; public class SnmpTestV2 { + + private static final Scanner scanner = new Scanner(System.in); + public static void main(String[] args) throws IOException { - SnmpDeviceSimulatorV2 device = new SnmpDeviceSimulatorV2(1610, "public"); + SnmpDeviceSimulatorV2 client = new SnmpDeviceSimulatorV2(1610, "public"); - device.start(); + client.start(); Map mappings = new HashMap<>(); // for (int i = 1; i <= 500; i++) { // String oid = String.format(".1.3.6.1.2.1.%s.1.52", i); // mappings.put(oid, "value_" + i); // } mappings.put("1.3.6.1.2.1.266.1.52", "****"); - device.setUpMappings(mappings); - new Scanner(System.in).nextLine(); + client.setUpMappings(mappings); + inputTraps(client); + + scanner.nextLine(); + } + + private static void inputTraps(SnmpDeviceSimulatorV2 client) throws IOException { + while (true) { + String data = scanner.nextLine(); + if (!data.isEmpty()) { + client.sendTrap("127.0.0.1", 1620, Map.of( + "1.3.6.1.2.1.266.1.52", data + " (266)", + "1.3.6.1.2.1.267.1.52", data + " (267)" + )); + } + } } } diff --git a/ui-ngx/src/app/shared/models/device.models.ts b/ui-ngx/src/app/shared/models/device.models.ts index f7802b3a8b..7679c3344a 100644 --- a/ui-ngx/src/app/shared/models/device.models.ts +++ b/ui-ngx/src/app/shared/models/device.models.ts @@ -290,14 +290,16 @@ export enum SnmpSpecType { TELEMETRY_QUERYING = 'TELEMETRY_QUERYING', CLIENT_ATTRIBUTES_QUERYING = 'CLIENT_ATTRIBUTES_QUERYING', SHARED_ATTRIBUTES_SETTING = 'SHARED_ATTRIBUTES_SETTING', - TO_DEVICE_RPC_REQUEST = 'TO_DEVICE_RPC_REQUEST' + TO_DEVICE_RPC_REQUEST = 'TO_DEVICE_RPC_REQUEST', + TO_SERVER_RPC_REQUEST = 'TO_SERVER_RPC_REQUEST' } export const SnmpSpecTypeTranslationMap = new Map([ - [SnmpSpecType.TELEMETRY_QUERYING, ' Telemetry'], - [SnmpSpecType.CLIENT_ATTRIBUTES_QUERYING, 'Client attributes'], - [SnmpSpecType.SHARED_ATTRIBUTES_SETTING, 'Shared attributes'], - [SnmpSpecType.TO_DEVICE_RPC_REQUEST, 'RPC request'] + [SnmpSpecType.TELEMETRY_QUERYING, ' Telemetry (SNMP GET)'], + [SnmpSpecType.CLIENT_ATTRIBUTES_QUERYING, 'Client attributes (SNMP GET)'], + [SnmpSpecType.SHARED_ATTRIBUTES_SETTING, 'Shared attributes (SNMP SET)'], + [SnmpSpecType.TO_DEVICE_RPC_REQUEST, 'To-device RPC request (SNMP GET/SET)'], + [SnmpSpecType.TO_SERVER_RPC_REQUEST, 'From-device RPC request (SNMP TRAP)'] ]); export interface SnmpCommunicationConfig {