From 7675a573167f99e625ec0b59103914e1191b56e2 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Mon, 12 Jun 2023 15:16:07 +0300 Subject: [PATCH 01/14] Bulk import for SNMP devices --- .../device/DeviceBulkImportService.java | 49 +++++++++++++------ .../importing/csv/BulkImportColumnType.java | 4 ++ .../import-export/import-export.models.ts | 8 +++ .../table-columns-assignment.component.ts | 4 ++ .../assets/locale/locale.constant-en_US.json | 6 +++ 5 files changed, 57 insertions(+), 14 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java index 509e0fdbef..a37e03de1b 100644 --- a/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java +++ b/application/src/main/java/org/thingsboard/server/service/device/DeviceBulkImportService.java @@ -33,7 +33,10 @@ import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials; import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MClientCredential; import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MSecurityMode; +import org.thingsboard.server.common.data.device.data.DefaultDeviceConfiguration; +import org.thingsboard.server.common.data.device.data.DeviceData; import org.thingsboard.server.common.data.device.data.PowerMode; +import org.thingsboard.server.common.data.device.data.SnmpDeviceTransportConfiguration; import org.thingsboard.server.common.data.device.profile.DefaultDeviceProfileConfiguration; import org.thingsboard.server.common.data.device.profile.DeviceProfileData; import org.thingsboard.server.common.data.device.profile.DisabledDeviceProfileProvisionConfiguration; @@ -45,6 +48,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.data.sync.ie.importing.csv.BulkImportColumnType; +import org.thingsboard.server.common.data.transport.snmp.SnmpProtocolVersion; import org.thingsboard.server.dao.device.DeviceCredentialsService; import org.thingsboard.server.dao.device.DeviceProfileService; import org.thingsboard.server.dao.device.DeviceService; @@ -76,18 +80,18 @@ public class DeviceBulkImportService extends AbstractBulkImportService { private final Lock findOrCreateDeviceProfileLock = new ReentrantLock(); @Override - protected void setEntityFields(Device entity, Map fields) { - ObjectNode additionalInfo = getOrCreateAdditionalInfoObj(entity); + protected void setEntityFields(Device device, Map fields) { + ObjectNode additionalInfo = getOrCreateAdditionalInfoObj(device); fields.forEach((columnType, value) -> { switch (columnType) { case NAME: - entity.setName(value); + device.setName(value); break; case TYPE: - entity.setType(value); + device.setType(value); break; case LABEL: - entity.setLabel(value); + device.setLabel(value); break; case DESCRIPTION: additionalInfo.set("description", new TextNode(value)); @@ -96,16 +100,17 @@ public class DeviceBulkImportService extends AbstractBulkImportService { additionalInfo.set("gateway", BooleanNode.valueOf(Boolean.parseBoolean(value))); break; } - entity.setAdditionalInfo(additionalInfo); + device.setAdditionalInfo(additionalInfo); }); + setUpDeviceConfiguration(device, fields); } @Override @SneakyThrows - protected Device saveEntity(SecurityUser user, Device entity, Map fields) { + protected Device saveEntity(SecurityUser user, Device device, Map fields) { DeviceCredentials deviceCredentials; try { - deviceCredentials = createDeviceCredentials(entity.getTenantId(), entity.getId(), fields); + deviceCredentials = createDeviceCredentials(device.getTenantId(), device.getId(), fields); deviceCredentialsService.formatCredentials(deviceCredentials); } catch (Exception e) { throw new DeviceCredentialsValidationException("Invalid device credentials: " + e.getMessage()); @@ -113,15 +118,15 @@ public class DeviceBulkImportService extends AbstractBulkImportService { DeviceProfile deviceProfile; if (deviceCredentials.getCredentialsType() == DeviceCredentialsType.LWM2M_CREDENTIALS) { - deviceProfile = setUpLwM2mDeviceProfile(entity.getTenantId(), entity); - } else if (StringUtils.isNotEmpty(entity.getType())) { - deviceProfile = deviceProfileService.findOrCreateDeviceProfile(entity.getTenantId(), entity.getType()); + deviceProfile = setUpLwM2mDeviceProfile(device.getTenantId(), device); + } else if (StringUtils.isNotEmpty(device.getType())) { + deviceProfile = deviceProfileService.findOrCreateDeviceProfile(device.getTenantId(), device.getType()); } else { - deviceProfile = deviceProfileService.findDefaultDeviceProfile(entity.getTenantId()); + deviceProfile = deviceProfileService.findDefaultDeviceProfile(device.getTenantId()); } - entity.setDeviceProfileId(deviceProfile.getId()); + device.setDeviceProfileId(deviceProfile.getId()); - return tbDeviceService.saveDeviceWithCredentials(entity, deviceCredentials, user); + return tbDeviceService.saveDeviceWithCredentials(device, deviceCredentials, user); } @Override @@ -136,6 +141,22 @@ public class DeviceBulkImportService extends AbstractBulkImportService { entity.setCustomerId(user.getCustomerId()); } + private void setUpDeviceConfiguration(Device device, Map fields) { + if (fields.containsKey(BulkImportColumnType.SNMP_HOST)) { + SnmpDeviceTransportConfiguration transportConfiguration = new SnmpDeviceTransportConfiguration(); + transportConfiguration.setHost(fields.get(BulkImportColumnType.SNMP_HOST)); + transportConfiguration.setPort(Optional.ofNullable(fields.get(BulkImportColumnType.SNMP_PORT)) + .map(Integer::parseInt).orElse(161)); + transportConfiguration.setProtocolVersion(Optional.ofNullable(fields.get(BulkImportColumnType.SNMP_VERSION)) + .map(version -> SnmpProtocolVersion.valueOf(version.toUpperCase())).orElse(SnmpProtocolVersion.V2C)); + transportConfiguration.setCommunity(fields.getOrDefault(BulkImportColumnType.SNMP_COMMUNITY_STRING, "public")); + + DeviceData deviceData = new DeviceData(); + deviceData.setTransportConfiguration(transportConfiguration); + device.setDeviceData(deviceData); + } + } + @SneakyThrows private DeviceCredentials createDeviceCredentials(TenantId tenantId, DeviceId deviceId, Map fields) { DeviceCredentials credentials = new DeviceCredentials(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/sync/ie/importing/csv/BulkImportColumnType.java b/common/data/src/main/java/org/thingsboard/server/common/data/sync/ie/importing/csv/BulkImportColumnType.java index c004187876..87582b4d6d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/sync/ie/importing/csv/BulkImportColumnType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/sync/ie/importing/csv/BulkImportColumnType.java @@ -43,6 +43,10 @@ public enum BulkImportColumnType { LWM2M_SERVER_SECURITY_MODE("securityMode", LwM2MSecurityMode.NO_SEC.name()), LWM2M_SERVER_CLIENT_PUBLIC_KEY_OR_ID("clientPublicKeyOrId"), LWM2M_SERVER_CLIENT_SECRET_KEY("clientSecretKey"), + SNMP_HOST, + SNMP_PORT, + SNMP_VERSION, + SNMP_COMMUNITY_STRING, IS_GATEWAY, DESCRIPTION, ROUTING_KEY, diff --git a/ui-ngx/src/app/modules/home/components/import-export/import-export.models.ts b/ui-ngx/src/app/modules/home/components/import-export/import-export.models.ts index fddfa80d91..5e70c50394 100644 --- a/ui-ngx/src/app/modules/home/components/import-export/import-export.models.ts +++ b/ui-ngx/src/app/modules/home/components/import-export/import-export.models.ts @@ -64,6 +64,10 @@ export enum ImportEntityColumnType { lwm2mServerSecurityMode = 'LWM2M_SERVER_SECURITY_MODE', lwm2mServerClientPublicKeyOrId = 'LWM2M_SERVER_CLIENT_PUBLIC_KEY_OR_ID', lwm2mServerClientSecretKey = 'LWM2M_SERVER_CLIENT_SECRET_KEY', + snmpHost = 'SNMP_HOST', + snmpPort = 'SNMP_PORT', + snmpVersion = 'SNMP_VERSION', + snmpCommunityString = 'SNMP_COMMUNITY_STRING', isGateway = 'IS_GATEWAY', description = 'DESCRIPTION', routingKey = 'ROUTING_KEY', @@ -95,6 +99,10 @@ export const importEntityColumnTypeTranslations = new Map Date: Mon, 12 Jun 2023 21:09:57 +0300 Subject: [PATCH 02/14] Send large SNMP requests in batches --- .../src/main/resources/thingsboard.yml | 1 + .../transport/snmp/SnmpTransportContext.java | 30 +++++++++------ .../transport/snmp/service/PduService.java | 38 ++++++++++++------- .../snmp/service/SnmpTransportService.java | 29 +++++++------- .../snmp/session/DeviceSessionContext.java | 8 ++++ .../transport/snmp/SnmpDeviceSimulatorV2.java | 23 +---------- .../server/transport/snmp/SnmpTestV2.java | 25 ++++-------- .../src/main/resources/tb-snmp-transport.yml | 1 + 8 files changed, 75 insertions(+), 80 deletions(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 96e409373e..0dc9ce9d25 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -959,6 +959,7 @@ transport: parallelism_level: "${SNMP_RESPONSE_PROCESSING_PARALLELISM_LEVEL:20}" # to configure SNMP to work over UDP or TCP underlying_protocol: "${SNMP_UNDERLYING_PROTOCOL:udp}" + max_request_oids: "${SNMP_MAX_REQUEST_OIDS:100}" stats: enabled: "${TB_TRANSPORT_STATS_ENABLED:true}" print-interval-ms: "${TB_TRANSPORT_STATS_PRINT_INTERVAL_MS:60000}" diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java index 5dcffbf62f..ba410cf960 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java @@ -178,17 +178,10 @@ public class SnmpTransportContext extends TransportContext { @Override public void onSuccess(ValidateDeviceCredentialsResponse msg) { if (msg.hasDeviceInfo()) { - SessionInfoProto sessionInfo = SessionInfoCreator.create( - msg, SnmpTransportContext.this, UUID.randomUUID() - ); - - transportService.registerAsyncSession(sessionInfo, deviceSessionContext); - transportService.process(sessionInfo, TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), TransportServiceCallback.EMPTY); - transportService.process(sessionInfo, TransportProtos.SubscribeToRPCMsg.newBuilder().build(), TransportServiceCallback.EMPTY); - - deviceSessionContext.setSessionInfo(sessionInfo); - deviceSessionContext.setDeviceInfo(msg.getDeviceInfo()); - deviceSessionContext.setConnected(true); + registerTransportSession(deviceSessionContext, msg); + deviceSessionContext.setSessionTimeoutHandler(() -> { + registerTransportSession(deviceSessionContext, msg); + }); } else { log.warn("[{}] Failed to process device auth", deviceSessionContext.getDeviceId()); } @@ -201,6 +194,21 @@ public class SnmpTransportContext extends TransportContext { }); } + private void registerTransportSession(DeviceSessionContext deviceSessionContext, ValidateDeviceCredentialsResponse msg) { + SessionInfoProto sessionInfo = SessionInfoCreator.create( + msg, SnmpTransportContext.this, UUID.randomUUID() + ); + log.debug("Registering transport session: {}", sessionInfo); + + transportService.registerAsyncSession(sessionInfo, deviceSessionContext); + transportService.process(sessionInfo, TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), TransportServiceCallback.EMPTY); + transportService.process(sessionInfo, TransportProtos.SubscribeToRPCMsg.newBuilder().build(), TransportServiceCallback.EMPTY); + + deviceSessionContext.setSessionInfo(sessionInfo); + deviceSessionContext.setDeviceInfo(msg.getDeviceInfo()); + deviceSessionContext.setConnected(true); + } + @EventListener(DeviceUpdatedEvent.class) public void onDeviceUpdatedOrCreated(DeviceUpdatedEvent deviceUpdatedEvent) { Device device = deviceUpdatedEvent.getDevice(); diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java index b395672c00..12fab72a6a 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.transport.snmp.service; +import com.google.common.collect.Lists; import com.google.gson.JsonObject; import lombok.extern.slf4j.Slf4j; import org.snmp4j.PDU; @@ -25,6 +26,7 @@ import org.snmp4j.smi.OID; import org.snmp4j.smi.OctetString; import org.snmp4j.smi.Variable; import org.snmp4j.smi.VariableBinding; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.device.data.SnmpDeviceTransportConfiguration; import org.thingsboard.server.common.data.kv.DataType; @@ -35,6 +37,7 @@ import org.thingsboard.server.common.data.transport.snmp.config.SnmpCommunicatio import org.thingsboard.server.queue.util.TbSnmpTransportComponent; import org.thingsboard.server.transport.snmp.session.DeviceSessionContext; +import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -47,21 +50,30 @@ import java.util.stream.IntStream; @Service @Slf4j public class PduService { - public PDU createPdu(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig, Map values) { - PDU pdu = setUpPdu(sessionContext); - pdu.setType(communicationConfig.getMethod().getCode()); - pdu.addAll(communicationConfig.getAllMappings().stream() - .filter(mapping -> values.isEmpty() || values.containsKey(mapping.getKey())) - .map(mapping -> Optional.ofNullable(values.get(mapping.getKey())) - .map(value -> { - Variable variable = toSnmpVariable(value, mapping.getDataType()); - return new VariableBinding(new OID(mapping.getOid()), variable); - }) - .orElseGet(() -> new VariableBinding(new OID(mapping.getOid())))) - .collect(Collectors.toList())); + @Value("${transport.snmp.max_request_oids:100}") + private int maxRequestOids; + + public List createPdus(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig, Map values) { + List pdus = new ArrayList<>(); + List allMappings = communicationConfig.getAllMappings(); + + for (List mappings : Lists.partition(allMappings, maxRequestOids)) { + PDU pdu = setUpPdu(sessionContext); + pdu.setType(communicationConfig.getMethod().getCode()); + pdu.addAll(mappings.stream() + .filter(mapping -> values.isEmpty() || values.containsKey(mapping.getKey())) + .map(mapping -> Optional.ofNullable(values.get(mapping.getKey())) + .map(value -> { + Variable variable = toSnmpVariable(value, mapping.getDataType()); + return new VariableBinding(new OID(mapping.getOid()), variable); + }) + .orElseGet(() -> new VariableBinding(new OID(mapping.getOid())))) + .collect(Collectors.toList())); + pdus.add(pdu); + } - return pdu; + return pdus; } public PDU createSingleVariablePdu(DeviceSessionContext sessionContext, SnmpMethod snmpMethod, String oid, String value, DataType dataType) { 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 6e68a0bd44..45b0d9ec0a 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 @@ -45,7 +45,6 @@ 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; @@ -161,18 +160,20 @@ public class SnmpTransportService implements TbTransportService { } private void sendRequest(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig, Map values) { - PDU request = pduService.createPdu(sessionContext, communicationConfig, values); + List request = pduService.createPdus(sessionContext, communicationConfig, values); RequestInfo requestInfo = new RequestInfo(communicationConfig.getSpec(), communicationConfig.getAllMappings()); sendRequest(sessionContext, request, requestInfo); } - 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()); - try { - snmp.send(request, sessionContext.getTarget(), requestInfo, sessionContext); - } catch (IOException e) { - log.error("Failed to send SNMP request to device {}: {}", sessionContext.getDeviceId(), e.toString()); + private void sendRequest(DeviceSessionContext sessionContext, List request, RequestInfo requestInfo) { + for (PDU pdu : request) { + if (pdu.size() > 0) { + log.trace("Executing SNMP request for device {}. Variables bindings: {}", sessionContext.getDeviceId(), pdu.getVariableBindings()); + try { + snmp.send(pdu, sessionContext.getTarget(), requestInfo, sessionContext); + } catch (IOException e) { + log.error("Failed to send SNMP request to device {}: {}", sessionContext.getDeviceId(), e.toString()); + } } } } @@ -216,7 +217,7 @@ public class SnmpTransportService implements TbTransportService { PDU request = pduService.createSingleVariablePdu(sessionContext, snmpMethod, oid, value, dataType); RequestInfo requestInfo = new RequestInfo(toDeviceRpcRequestMsg.getRequestId(), communicationConfig.getSpec(), communicationConfig.getAllMappings()); - sendRequest(sessionContext, request, requestInfo); + sendRequest(sessionContext, List.of(request), requestInfo); } @@ -247,7 +248,7 @@ public class SnmpTransportService implements TbTransportService { 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()); + log.debug("No values in the SNMP response for device {}. Request id: {}", sessionContext.getDeviceId(), response.getRequestID()); return; } @@ -302,11 +303,7 @@ public class SnmpTransportService implements TbTransportService { } 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); } diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java index 2b7c452ca7..fede19824f 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java @@ -29,6 +29,7 @@ import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.transport.SessionMsgListener; import org.thingsboard.server.common.transport.TransportServiceCallback; +import org.thingsboard.server.common.transport.service.DefaultTransportService; import org.thingsboard.server.common.transport.session.DeviceAwareSessionContext; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg; @@ -63,6 +64,8 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S private final AtomicInteger msgIdSeq = new AtomicInteger(0); @Getter private boolean isActive = true; + @Setter + private Runnable sessionTimeoutHandler; @Getter private final List> queryingTasks = new LinkedList<>(); @@ -137,6 +140,11 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S @Override public void onRemoteSessionCloseCommand(UUID sessionId, SessionCloseNotificationProto sessionCloseNotification) { log.trace("[{}] Received the remote command to close the session: {}", sessionId, sessionCloseNotification.getMessage()); + if (sessionCloseNotification.getMessage().equals(DefaultTransportService.SESSION_EXPIRED_MESSAGE)) { + if (sessionTimeoutHandler != null) { + sessionTimeoutHandler.run(); + } + } } @Override diff --git a/common/transport/snmp/src/test/java/org/thingsboard/server/transport/snmp/SnmpDeviceSimulatorV2.java b/common/transport/snmp/src/test/java/org/thingsboard/server/transport/snmp/SnmpDeviceSimulatorV2.java index 5fd7f3bed0..fa49b17f13 100644 --- a/common/transport/snmp/src/test/java/org/thingsboard/server/transport/snmp/SnmpDeviceSimulatorV2.java +++ b/common/transport/snmp/src/test/java/org/thingsboard/server/transport/snmp/SnmpDeviceSimulatorV2.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.transport.snmp; -import org.snmp4j.CommandResponderEvent; import org.snmp4j.CommunityTarget; import org.snmp4j.PDU; import org.snmp4j.Snmp; @@ -35,7 +34,6 @@ import org.snmp4j.agent.mo.snmp.SnmpTargetMIB; import org.snmp4j.agent.mo.snmp.StorageType; import org.snmp4j.agent.mo.snmp.VacmMIB; import org.snmp4j.agent.security.MutableVACM; -import org.snmp4j.mp.MPv3; import org.snmp4j.mp.SnmpConstants; import org.snmp4j.security.SecurityLevel; import org.snmp4j.security.SecurityModel; @@ -53,27 +51,11 @@ import org.snmp4j.transport.TransportMappings; import java.io.File; import java.io.IOException; import java.util.Map; -import java.util.function.Consumer; import java.util.stream.Collectors; @SuppressWarnings("deprecation") public class SnmpDeviceSimulatorV2 extends BaseAgent { - public static class RequestProcessor extends CommandProcessor { - private final Consumer processor; - - public RequestProcessor(Consumer processor) { - super(new OctetString(MPv3.createLocalEngineID())); - this.processor = processor; - } - - @Override - public void processPdu(CommandResponderEvent event) { - processor.accept(event); - } - } - - private final Target target; private final Address address; private Snmp snmp; @@ -81,10 +63,7 @@ public class SnmpDeviceSimulatorV2 extends BaseAgent { private final String password; public SnmpDeviceSimulatorV2(int port, String password) throws IOException { - super(new File("conf.agent"), new File("bootCounter.agent"), new RequestProcessor(event -> { - System.out.println("aboba"); - ((Snmp) event.getSource()).cancel(event.getPDU(), event1 -> System.out.println("canceled")); - })); + super(new File("conf.agent"), new File("bootCounter.agent"), new CommandProcessor(new OctetString("12312"))); CommunityTarget target = new CommunityTarget(); target.setCommunity(new OctetString(password)); this.address = GenericAddress.parse("udp:0.0.0.0/" + port); 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 bf371058aa..55ed8fb97f 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 @@ -16,6 +16,7 @@ package org.thingsboard.server.transport.snmp; import java.io.IOException; +import java.util.HashMap; import java.util.Map; import java.util.Scanner; @@ -24,24 +25,12 @@ public class SnmpTestV2 { SnmpDeviceSimulatorV2 device = new SnmpDeviceSimulatorV2(1610, "public"); device.start(); - device.setUpMappings(Map.of( - ".1.3.6.1.2.1.1.1.50", "12", - ".1.3.6.1.2.1.2.1.52", "56", - ".1.3.6.1.2.1.3.1.54", "yes", - ".1.3.6.1.2.1.7.1.58", "" - )); - - -// while (true) { -// new Scanner(System.in).nextLine(); -// device.sendTrap("127.0.0.1", 1062, Map.of(".1.3.6.1.2.87.1.56", "12")); -// System.out.println("sent"); -// } - -// Snmp snmp = new Snmp(device.transportMappings[0]); -// device.snmp.addCommandResponder(event -> { -// System.out.println(event); -// }); + 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); + } + device.setUpMappings(mappings); new Scanner(System.in).nextLine(); } diff --git a/transport/snmp/src/main/resources/tb-snmp-transport.yml b/transport/snmp/src/main/resources/tb-snmp-transport.yml index 9f086bcbc5..8fcd4cc2d2 100644 --- a/transport/snmp/src/main/resources/tb-snmp-transport.yml +++ b/transport/snmp/src/main/resources/tb-snmp-transport.yml @@ -93,6 +93,7 @@ transport: parallelism_level: "${SNMP_RESPONSE_PROCESSING_PARALLELISM_LEVEL:20}" # to configure SNMP to work over UDP or TCP underlying_protocol: "${SNMP_UNDERLYING_PROTOCOL:udp}" + max_request_oids: "${SNMP_MAX_REQUEST_OIDS:100}" sessions: inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}" report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:3000}" From c1f5b39cdb348d45752127264b85b494d9ba97ca Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Wed, 14 Jun 2023 13:21:32 +0300 Subject: [PATCH 03/14] Single telemetry message when splitting large SNMP request --- .../transport/snmp/SnmpTransportContext.java | 3 +- .../transport/snmp/service/PduService.java | 15 ++- .../snmp/service/SnmpTransportService.java | 127 +++++++++++------- 3 files changed, 90 insertions(+), 55 deletions(-) diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java index ba410cf960..0d83d45fb5 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java @@ -212,7 +212,7 @@ public class SnmpTransportContext extends TransportContext { @EventListener(DeviceUpdatedEvent.class) public void onDeviceUpdatedOrCreated(DeviceUpdatedEvent deviceUpdatedEvent) { Device device = deviceUpdatedEvent.getDevice(); - log.trace("Got creating or updating device event for device {}", device); + log.debug("Got creating or updating device event for device {}", device); DeviceTransportType transportType = Optional.ofNullable(device.getDeviceData().getTransportConfiguration()) .map(DeviceTransportConfiguration::getType) .orElse(null); @@ -246,6 +246,7 @@ public class SnmpTransportContext extends TransportContext { } public void onDeviceProfileUpdated(DeviceProfile deviceProfile, DeviceSessionContext sessionContext) { + log.debug("Handling device profile {} update event for device {}", deviceProfile.getId(), sessionContext.getDeviceId()); updateDeviceSession(sessionContext, sessionContext.getDevice(), deviceProfile); } diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java index 12fab72a6a..72de03e928 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java @@ -44,7 +44,6 @@ import java.util.Map; import java.util.Objects; import java.util.Optional; import java.util.stream.Collectors; -import java.util.stream.IntStream; @TbSnmpTransportComponent @Service @@ -70,7 +69,9 @@ public class PduService { }) .orElseGet(() -> new VariableBinding(new OID(mapping.getOid())))) .collect(Collectors.toList())); - pdus.add(pdu); + if (pdu.size() > 0) { + pdus.add(pdu); + } } return pdus; @@ -128,8 +129,8 @@ public class PduService { } - public JsonObject processPdu(PDU pdu, List responseMappings) { - Map values = processPdu(pdu); + public JsonObject processPdus(List pdus, List responseMappings) { + Map values = processPdus(pdus); Map mappings = new HashMap<>(); if (responseMappings != null) { @@ -155,9 +156,9 @@ public class PduService { return data; } - public Map processPdu(PDU pdu) { - return IntStream.range(0, pdu.size()) - .mapToObj(pdu::get) + public Map processPdus(List pdus) { + return pdus.stream() + .flatMap(pdu -> pdu.getVariableBindings().stream()) .filter(Objects::nonNull) .filter(variableBinding -> !(variableBinding.getVariable() instanceof Null)) .collect(Collectors.toMap(VariableBinding::getOid, VariableBinding::toValueString)); 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 45b0d9ec0a..cf4aa315a8 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 @@ -17,6 +17,7 @@ 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; @@ -53,6 +54,7 @@ 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; @@ -161,19 +163,21 @@ public class SnmpTransportService implements TbTransportService { private void sendRequest(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig, Map values) { List request = pduService.createPdus(sessionContext, communicationConfig, values); - RequestInfo requestInfo = new RequestInfo(communicationConfig.getSpec(), communicationConfig.getAllMappings()); - sendRequest(sessionContext, request, requestInfo); + RequestContext requestContext = RequestContext.builder() + .communicationSpec(communicationConfig.getSpec()) + .responseMappings(communicationConfig.getAllMappings()) + .requestSize(request.size()) + .build(); + sendRequest(sessionContext, request, requestContext); } - private void sendRequest(DeviceSessionContext sessionContext, List request, RequestInfo requestInfo) { + private void sendRequest(DeviceSessionContext sessionContext, List request, RequestContext requestContext) { for (PDU pdu : request) { - if (pdu.size() > 0) { - log.trace("Executing SNMP request for device {}. Variables bindings: {}", sessionContext.getDeviceId(), pdu.getVariableBindings()); - try { - snmp.send(pdu, sessionContext.getTarget(), requestInfo, sessionContext); - } catch (IOException e) { - log.error("Failed to send SNMP request to device {}: {}", sessionContext.getDeviceId(), e.toString()); - } + log.trace("Executing SNMP request for device {} with {} variable bindings", sessionContext.getDeviceId(), pdu.size()); + try { + snmp.send(pdu, sessionContext.getTarget(), requestContext, sessionContext); + } catch (IOException e) { + log.error("Failed to send SNMP request to device {}: {}", sessionContext.getDeviceId(), e.toString()); } } } @@ -216,51 +220,77 @@ 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, List.of(request), requestInfo); + RequestContext requestContext = RequestContext.builder() + .requestId(toDeviceRpcRequestMsg.getRequestId()) + .communicationSpec(communicationConfig.getSpec()) + .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); - if (event.getError() != null) { log.warn("SNMP response error: {}", event.getError().toString()); return; } - PDU response = event.getResponse(); - if (response == null) { - log.debug("No response from SNMP device {}, requestId: {}", sessionContext.getDeviceId(), event.getRequest().getRequestID()); - return; + PDU responsePdu = event.getResponse(); + if (log.isTraceEnabled()) { + log.trace("Received PDU for device {}: {}", sessionContext.getDeviceId(), responsePdu); + } + RequestContext requestContext = (RequestContext) event.getUserObject(); + + List response; + if (requestContext.getRequestSize() == 1) { + if (responsePdu == null) { + log.debug("No response from SNMP device {}, requestId: {}", sessionContext.getDeviceId(), event.getRequest().getRequestID()); + return; + } + response = List.of(responsePdu); + } else { + List 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.trace("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; + } } - RequestInfo requestInfo = (RequestInfo) event.getUserObject(); responseProcessingExecutor.execute(() -> { - processResponse(sessionContext, response, requestInfo); + processResponse(sessionContext, response, requestContext); }); } - private void processResponse(DeviceSessionContext sessionContext, PDU response, RequestInfo requestInfo) { - ResponseProcessor responseProcessor = responseProcessors.get(requestInfo.getCommunicationSpec()); + private void processResponse(DeviceSessionContext sessionContext, List 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 in the SNMP response for device {}. Request id: {}", sessionContext.getDeviceId(), response.getRequestID()); + 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()); return; } - 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 -> { @@ -270,8 +300,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 -> { @@ -280,21 +310,21 @@ 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); @@ -332,29 +362,32 @@ public class SnmpTransportService implements TbTransportService { } @Data - private static class RequestInfo { - private Integer requestId; - private SnmpCommunicationSpec communicationSpec; - private List responseMappings; + private static class RequestContext { + private final Integer requestId; + private final SnmpCommunicationSpec communicationSpec; + private final List responseMappings; - public RequestInfo(Integer requestId, SnmpCommunicationSpec communicationSpec, List responseMappings) { - this.requestId = requestId; - this.communicationSpec = communicationSpec; - this.responseMappings = responseMappings; - } + private final int requestSize; + private List responseParts; - public RequestInfo(SnmpCommunicationSpec communicationSpec, List responseMappings) { + @Builder + public RequestContext(Integer requestId, SnmpCommunicationSpec communicationSpec, List responseMappings, int requestSize) { + this.requestId = requestId; this.communicationSpec = communicationSpec; 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 pdus, RequestContext requestContext); } private interface ResponseProcessor { - void process(JsonObject responseData, RequestInfo requestInfo, DeviceSessionContext sessionContext); + void process(JsonObject responseData, RequestContext requestContext, DeviceSessionContext sessionContext); } } From a761d0b204745e8ef10c17ec90c6c549415a0f82 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Wed, 14 Jun 2023 13:50:41 +0300 Subject: [PATCH 04/14] Update logging levels for SnmpTransportService --- .../server/transport/snmp/service/SnmpTransportService.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 cf4aa315a8..9177839c23 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 @@ -173,7 +173,7 @@ public class SnmpTransportService implements TbTransportService { private void sendRequest(DeviceSessionContext sessionContext, List request, RequestContext requestContext) { for (PDU pdu : request) { - log.trace("Executing SNMP request for device {} with {} variable bindings", sessionContext.getDeviceId(), pdu.size()); + log.debug("Executing SNMP request for device {} with {} variable bindings", sessionContext.getDeviceId(), pdu.size()); try { snmp.send(pdu, sessionContext.getTarget(), requestContext, sessionContext); } catch (IOException e) { @@ -260,7 +260,7 @@ public class SnmpTransportService implements TbTransportService { response.add(responsePart); } } - log.trace("All response parts are collected for request to device {}", sessionContext.getDeviceId()); + 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; From 277c87240cd2da1f83c1ed9e99c246480bb2dfaf Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Sat, 17 Jun 2023 14:53:26 +0300 Subject: [PATCH 05/14] Fix attributes and rpc subscription for SNMP devices --- .../server/transport/snmp/SnmpTransportContext.java | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java index 0d83d45fb5..49dc3e8c43 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java @@ -201,8 +201,12 @@ public class SnmpTransportContext extends TransportContext { log.debug("Registering transport session: {}", sessionInfo); transportService.registerAsyncSession(sessionInfo, deviceSessionContext); - transportService.process(sessionInfo, TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().build(), TransportServiceCallback.EMPTY); - transportService.process(sessionInfo, TransportProtos.SubscribeToRPCMsg.newBuilder().build(), TransportServiceCallback.EMPTY); + transportService.process(sessionInfo, TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder() + .setSessionType(TransportProtos.SessionType.ASYNC) + .build(), TransportServiceCallback.EMPTY); + transportService.process(sessionInfo, TransportProtos.SubscribeToRPCMsg.newBuilder() + .setSessionType(TransportProtos.SessionType.ASYNC) + .build(), TransportServiceCallback.EMPTY); deviceSessionContext.setSessionInfo(sessionInfo); deviceSessionContext.setDeviceInfo(msg.getDeviceInfo()); From 09b866c41f6fb096206c3c67889a7fe5791150fa Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 22 Jun 2023 18:10:03 +0300 Subject: [PATCH 06/14] Remove unnecessary validation for OID mappings uniqueness --- .../profile/SnmpDeviceProfileTransportConfiguration.java | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/SnmpDeviceProfileTransportConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/SnmpDeviceProfileTransportConfiguration.java index 14b62fe02e..6c89c91006 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/SnmpDeviceProfileTransportConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/profile/SnmpDeviceProfileTransportConfiguration.java @@ -18,7 +18,6 @@ package org.thingsboard.server.common.data.device.profile; import com.fasterxml.jackson.annotation.JsonIgnore; import lombok.Data; import org.thingsboard.server.common.data.DeviceTransportType; -import org.thingsboard.server.common.data.transport.snmp.SnmpMapping; import org.thingsboard.server.common.data.transport.snmp.config.SnmpCommunicationConfig; import java.util.List; @@ -45,8 +44,7 @@ public class SnmpDeviceProfileTransportConfiguration implements DeviceProfileTra private boolean isValid() { return timeoutMs != null && timeoutMs >= 0 && retries != null && retries >= 0 && communicationConfigs != null - && communicationConfigs.stream().allMatch(config -> config != null && config.isValid()) - && communicationConfigs.stream().flatMap(config -> config.getAllMappings().stream()).map(SnmpMapping::getOid) - .distinct().count() == communicationConfigs.stream().mapToInt(config -> config.getAllMappings().size()).sum(); + && communicationConfigs.stream().allMatch(config -> config != null && config.isValid()); } + } From 854e059435484e8c6030f35b9067c307f430b405 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 22 Jun 2023 18:19:16 +0300 Subject: [PATCH 07/14] Option to ignore SNMP response type cast errors --- .../csv/AbstractBulkImportService.java | 5 +- .../src/main/resources/thingsboard.yml | 2 + .../server/common/data/StringUtils.java | 9 ++++ .../common/data/util}/TypeCastUtil.java | 32 ++++++++++--- .../transport/snmp/service/PduService.java | 46 +++++++++++++------ .../src/main/resources/tb-snmp-transport.yml | 2 + 6 files changed, 73 insertions(+), 23 deletions(-) rename {application/src/main/java/org/thingsboard/server/utils => common/data/src/main/java/org/thingsboard/server/common/data/util}/TypeCastUtil.java (58%) diff --git a/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java index 45805c5939..c03dd8dcb7 100644 --- a/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java +++ b/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java @@ -22,6 +22,7 @@ import com.google.gson.JsonPrimitive; import lombok.Data; import lombok.SneakyThrows; import org.apache.commons.lang3.exception.ExceptionUtils; +import org.apache.commons.lang3.tuple.Pair; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.security.core.context.SecurityContext; import org.springframework.security.core.context.SecurityContextHolder; @@ -57,7 +58,7 @@ import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Resource; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import org.thingsboard.server.utils.CsvUtils; -import org.thingsboard.server.utils.TypeCastUtil; +import org.thingsboard.server.common.data.util.TypeCastUtil; import javax.annotation.Nullable; import javax.annotation.PostConstruct; @@ -269,7 +270,7 @@ public abstract class AbstractBulkImportService castResult = TypeCastUtil.castValue(entry.getValue()); + Pair castResult = TypeCastUtil.castValue(entry.getValue()); entityData.getKvs().put(entry.getKey(), new ParsedValue(castResult.getValue(), castResult.getKey())); } }); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 15637ddb9e..1a5443578e 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -975,6 +975,8 @@ transport: # to configure SNMP to work over UDP or TCP underlying_protocol: "${SNMP_UNDERLYING_PROTOCOL:udp}" max_request_oids: "${SNMP_MAX_REQUEST_OIDS:100}" + response: + ignore_type_cast_errors: "${SNMP_RESPONSE_IGNORE_TYPE_CAST_ERRORS:false}" stats: enabled: "${TB_TRANSPORT_STATS_ENABLED:true}" print-interval-ms: "${TB_TRANSPORT_STATS_PRINT_INTERVAL_MS:60000}" diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java b/common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java index a7671f4327..f9d5aec059 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java @@ -163,6 +163,15 @@ public class StringUtils { return false; } + public static boolean equalsAnyIgnoreCase(String string, String... otherStrings) { + for (String otherString : otherStrings) { + if (equalsIgnoreCase(string, otherString)) { + return true; + } + } + return false; + } + public static String substringAfterLast(String str, String sep) { return org.apache.commons.lang3.StringUtils.substringAfterLast(str, sep); } diff --git a/application/src/main/java/org/thingsboard/server/utils/TypeCastUtil.java b/common/data/src/main/java/org/thingsboard/server/common/data/util/TypeCastUtil.java similarity index 58% rename from application/src/main/java/org/thingsboard/server/utils/TypeCastUtil.java rename to common/data/src/main/java/org/thingsboard/server/common/data/util/TypeCastUtil.java index 3bf1efc8f4..cc40c0d13f 100644 --- a/application/src/main/java/org/thingsboard/server/utils/TypeCastUtil.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/util/TypeCastUtil.java @@ -13,35 +13,53 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.utils; +package org.thingsboard.server.common.data.util; import org.apache.commons.lang3.math.NumberUtils; +import org.apache.commons.lang3.tuple.Pair; import org.thingsboard.server.common.data.kv.DataType; import java.math.BigDecimal; -import java.util.Map; public class TypeCastUtil { private TypeCastUtil() {} - public static Map.Entry castValue(String value) { + public static Pair castValue(String value) { if (isNumber(value)) { String formattedValue = value.replace(',', '.'); try { BigDecimal bd = new BigDecimal(formattedValue); if (bd.stripTrailingZeros().scale() > 0 || isSimpleDouble(formattedValue)) { if (bd.scale() <= 16) { - return Map.entry(DataType.DOUBLE, bd.doubleValue()); + return Pair.of(DataType.DOUBLE, bd.doubleValue()); } } else { - return Map.entry(DataType.LONG, bd.longValueExact()); + return Pair.of(DataType.LONG, bd.longValueExact()); } } catch (RuntimeException ignored) {} } else if (value.equalsIgnoreCase("true") || value.equalsIgnoreCase("false")) { - return Map.entry(DataType.BOOLEAN, Boolean.parseBoolean(value)); + return Pair.of(DataType.BOOLEAN, Boolean.parseBoolean(value)); + } + return Pair.of(DataType.STRING, value); + } + + public static Pair castToNumber(String value) { + if (isNumber(value)) { + String formattedValue = value.replace(',', '.'); + BigDecimal bd = new BigDecimal(formattedValue); + if (bd.stripTrailingZeros().scale() > 0 || isSimpleDouble(formattedValue)) { + if (bd.scale() <= 16) { + return Pair.of(DataType.DOUBLE, bd.doubleValue()); + } else { + return Pair.of(DataType.DOUBLE, bd); + } + } else { + return Pair.of(DataType.LONG, bd.longValueExact()); + } + } else { + throw new IllegalArgumentException("'" + value + "' can't be parsed as number"); } - return Map.entry(DataType.STRING, value); } private static boolean isNumber(String value) { diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java index 72de03e928..1cfcfd2c85 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java @@ -17,6 +17,7 @@ package org.thingsboard.server.transport.snmp.service; import com.google.common.collect.Lists; import com.google.gson.JsonObject; +import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.snmp4j.PDU; import org.snmp4j.ScopedPDU; @@ -28,12 +29,14 @@ import org.snmp4j.smi.Variable; import org.snmp4j.smi.VariableBinding; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.device.data.SnmpDeviceTransportConfiguration; import org.thingsboard.server.common.data.kv.DataType; import org.thingsboard.server.common.data.transport.snmp.SnmpMapping; import org.thingsboard.server.common.data.transport.snmp.SnmpMethod; import org.thingsboard.server.common.data.transport.snmp.SnmpProtocolVersion; import org.thingsboard.server.common.data.transport.snmp.config.SnmpCommunicationConfig; +import org.thingsboard.server.common.data.util.TypeCastUtil; import org.thingsboard.server.queue.util.TbSnmpTransportComponent; import org.thingsboard.server.transport.snmp.session.DeviceSessionContext; @@ -48,11 +51,15 @@ import java.util.stream.Collectors; @TbSnmpTransportComponent @Service @Slf4j +@RequiredArgsConstructor public class PduService { @Value("${transport.snmp.max_request_oids:100}") private int maxRequestOids; + @Value("${transport.snmp.response.ignore_type_cast_errors:false}") + private boolean ignoreTypeCastErrors; + public List createPdus(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig, Map values) { List pdus = new ArrayList<>(); List allMappings = communicationConfig.getAllMappings(); @@ -165,20 +172,31 @@ public class PduService { } public void processValue(String key, DataType dataType, String value, JsonObject result) { - switch (dataType) { - case LONG: - result.addProperty(key, Long.parseLong(value)); - break; - case BOOLEAN: - result.addProperty(key, Boolean.parseBoolean(value)); - break; - case DOUBLE: - result.addProperty(key, Double.parseDouble(value)); - break; - case STRING: - case JSON: - default: - result.addProperty(key, value); + try { + switch (dataType) { + case STRING: + case JSON: + result.addProperty(key, value); + break; + case LONG: + case DOUBLE: + result.addProperty(key, TypeCastUtil.castToNumber(value).getValue()); + break; + case BOOLEAN: + if (StringUtils.equalsAnyIgnoreCase(value, "true", "false")) { + result.addProperty(key, Boolean.parseBoolean(value)); + } else { + throw new IllegalArgumentException("Can't parse '" + value + "' as boolean"); + } + break; + } + } catch (IllegalArgumentException e) { + if (ignoreTypeCastErrors) { + log.debug("Ignoring value '{}' for key '{}' because of data type mismatch ({} required)", value, key, dataType); + } else { + throw e; + } } } + } diff --git a/transport/snmp/src/main/resources/tb-snmp-transport.yml b/transport/snmp/src/main/resources/tb-snmp-transport.yml index 066eaa8f3c..5b5fe34af1 100644 --- a/transport/snmp/src/main/resources/tb-snmp-transport.yml +++ b/transport/snmp/src/main/resources/tb-snmp-transport.yml @@ -104,6 +104,8 @@ transport: # to configure SNMP to work over UDP or TCP underlying_protocol: "${SNMP_UNDERLYING_PROTOCOL:udp}" max_request_oids: "${SNMP_MAX_REQUEST_OIDS:100}" + response: + ignore_type_cast_errors: "${SNMP_RESPONSE_IGNORE_TYPE_CAST_ERRORS:false}" sessions: inactivity_timeout: "${TB_TRANSPORT_SESSIONS_INACTIVITY_TIMEOUT:300000}" report_timeout: "${TB_TRANSPORT_SESSIONS_REPORT_TIMEOUT:3000}" From 926f4842300c92f289b65f46a90aaacf9c51bd1f Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 22 Jun 2023 18:19:37 +0300 Subject: [PATCH 08/14] Device events --- .../queue/DefaultTbCoreConsumerService.java | 73 ++++++++++++++++--- common/cluster-api/src/main/proto/queue.proto | 23 ++++++ .../transport/snmp/SnmpCommunicationSpec.java | 16 +++- .../transport/snmp/SnmpTransportContext.java | 43 +++++++---- .../snmp/service/SnmpTransportService.java | 22 +++++- .../snmp/session/DeviceSessionContext.java | 23 +++++- .../server/transport/snmp/SnmpTestV2.java | 9 ++- .../common/transport/TransportService.java | 7 ++ .../service/DefaultTransportService.java | 50 +++++++++++-- 9 files changed, 216 insertions(+), 50 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index dc24d6df33..db0b7c3809 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -28,14 +28,18 @@ import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.common.data.alarm.AlarmInfo; +import org.thingsboard.server.common.data.event.ErrorEvent; +import org.thingsboard.server.common.data.event.Event; +import org.thingsboard.server.common.data.event.LifecycleEvent; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.NotificationRequestId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; +import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTrigger; import org.thingsboard.server.common.data.rpc.RpcError; import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.TbActorMsg; -import org.thingsboard.server.common.data.notification.rule.trigger.NotificationRuleTrigger; +import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; @@ -44,7 +48,9 @@ import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.DeviceStateServiceMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.EdgeNotificationMsgProto; +import org.thingsboard.server.gen.transport.TransportProtos.ErrorEventProto; import org.thingsboard.server.gen.transport.TransportProtos.FromDeviceRPCResponseProto; +import org.thingsboard.server.gen.transport.TransportProtos.LifecycleEventProto; import org.thingsboard.server.gen.transport.TransportProtos.LocalSubscriptionServiceMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionMgrMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.TbAlarmDeleteProto; @@ -63,7 +69,6 @@ import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; -import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.DataDecodingEncodingService; @@ -274,6 +279,10 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService() { @Override public void onSuccess(ValidateDeviceCredentialsResponse msg) { if (msg.hasDeviceInfo()) { - registerTransportSession(deviceSessionContext, msg); - deviceSessionContext.setSessionTimeoutHandler(() -> { - registerTransportSession(deviceSessionContext, msg); + registerTransportSession(sessionContext, msg); + sessionContext.setSessionTimeoutHandler(() -> { + registerTransportSession(sessionContext, msg); }); + transportService.lifecycleEvent(sessionContext.getTenantId(), sessionContext.getDeviceId(), ComponentLifecycleEvent.STARTED, true, null); } else { - log.warn("[{}] Failed to process device auth", deviceSessionContext.getDeviceId()); + log.warn("[{}] Failed to process device auth", sessionContext.getDeviceId()); } } @Override public void onError(Throwable e) { - log.warn("[{}] Failed to process device auth: {}", deviceSessionContext.getDeviceId(), e); + log.warn("[{}] Failed to process device auth: {}", sessionContext.getDeviceId(), e); + transportService.lifecycleEvent(sessionContext.getTenantId(), sessionContext.getDeviceId(), ComponentLifecycleEvent.STARTED, false, e); } }); } 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 9177839c23..d14cafb2c2 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 @@ -144,6 +144,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); }) @@ -165,6 +166,7 @@ public class SnmpTransportService implements TbTransportService { List request = pduService.createPdus(sessionContext, communicationConfig, values); RequestContext requestContext = RequestContext.builder() .communicationSpec(communicationConfig.getSpec()) + .method(communicationConfig.getMethod()) .responseMappings(communicationConfig.getAllMappings()) .requestSize(request.size()) .build(); @@ -178,6 +180,7 @@ public class SnmpTransportService implements TbTransportService { 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); } } } @@ -223,6 +226,7 @@ public class SnmpTransportService implements TbTransportService { RequestContext requestContext = RequestContext.builder() .requestId(toDeviceRpcRequestMsg.getRequestId()) .communicationSpec(communicationConfig.getSpec()) + .method(snmpMethod) .responseMappings(communicationConfig.getAllMappings()) .requestSize(1) .build(); @@ -232,8 +236,10 @@ public class SnmpTransportService implements TbTransportService { 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; } @@ -241,12 +247,14 @@ public class SnmpTransportService implements TbTransportService { if (log.isTraceEnabled()) { log.trace("Received PDU for device {}: {}", sessionContext.getDeviceId(), responsePdu); } - RequestContext requestContext = (RequestContext) event.getUserObject(); List 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); @@ -268,7 +276,11 @@ public class SnmpTransportService implements TbTransportService { } responseProcessingExecutor.execute(() -> { - processResponse(sessionContext, response, requestContext); + try { + processResponse(sessionContext, response, requestContext); + } catch (Exception e) { + transportService.errorEvent(sessionContext.getTenantId(), sessionContext.getDeviceId(), requestContext.getCommunicationSpec().getLabel(), e); + } }); } @@ -279,7 +291,7 @@ public class SnmpTransportService implements TbTransportService { 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()); - return; + throw new IllegalArgumentException("No values in the response"); } responseProcessor.process(responseData, requestContext, sessionContext); @@ -365,15 +377,17 @@ public class SnmpTransportService implements TbTransportService { private static class RequestContext { private final Integer requestId; private final SnmpCommunicationSpec communicationSpec; + private final SnmpMethod method; private final List responseMappings; private final int requestSize; private List responseParts; @Builder - public RequestContext(Integer requestId, SnmpCommunicationSpec communicationSpec, List responseMappings, int requestSize) { + public RequestContext(Integer requestId, SnmpCommunicationSpec communicationSpec, SnmpMethod method, List responseMappings, int requestSize) { this.requestId = requestId; this.communicationSpec = communicationSpec; + this.method = method; this.responseMappings = responseMappings; this.requestSize = requestSize; if (requestSize > 1) { diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java index fede19824f..b95a11e8c2 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.transport.snmp.session; +import lombok.Builder; import lombok.Getter; import lombok.Setter; import lombok.extern.slf4j.Slf4j; @@ -26,7 +27,9 @@ import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.device.data.SnmpDeviceTransportConfiguration; import org.thingsboard.server.common.data.device.profile.SnmpDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.rpc.RpcStatus; +import org.thingsboard.server.common.data.transport.snmp.SnmpCommunicationSpec; import org.thingsboard.server.common.transport.SessionMsgListener; import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.common.transport.service.DefaultTransportService; @@ -58,6 +61,8 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S private SnmpDeviceTransportConfiguration deviceTransportConfiguration; @Getter private final Device device; + @Getter + private final TenantId tenantId; private final SnmpTransportContext snmpTransportContext; @@ -70,7 +75,8 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S @Getter private final List> queryingTasks = new LinkedList<>(); - public DeviceSessionContext(Device device, DeviceProfile deviceProfile, String token, + @Builder + public DeviceSessionContext(TenantId tenantId, Device device, DeviceProfile deviceProfile, String token, SnmpDeviceProfileTransportConfiguration profileTransportConfiguration, SnmpDeviceTransportConfiguration deviceTransportConfiguration, SnmpTransportContext snmpTransportContext) throws Exception { @@ -78,6 +84,7 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S super.setDeviceId(device.getId()); super.setDeviceProfile(deviceProfile); this.device = device; + this.tenantId = tenantId; this.token = token; this.snmpTransportContext = snmpTransportContext; @@ -134,7 +141,11 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S @Override public void onAttributeUpdate(UUID sessionId, AttributeUpdateNotificationMsg attributeUpdateNotification) { log.trace("[{}] Received attributes update notification to device", sessionId); - snmpTransportContext.getSnmpTransportService().onAttributeUpdate(this, attributeUpdateNotification); + try { + snmpTransportContext.getSnmpTransportService().onAttributeUpdate(this, attributeUpdateNotification); + } catch (Exception e) { + snmpTransportContext.getTransportService().errorEvent(getTenantId(), getDeviceId(), SnmpCommunicationSpec.SHARED_ATTRIBUTES_SETTING.getLabel(), e); + } } @Override @@ -150,8 +161,12 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S @Override public void onToDeviceRpcRequest(UUID sessionId, ToDeviceRpcRequestMsg toDeviceRequest) { log.trace("[{}] Received RPC command to device", sessionId); - snmpTransportContext.getSnmpTransportService().onToDeviceRpcRequest(this, toDeviceRequest); - snmpTransportContext.getTransportService().process(getSessionInfo(), toDeviceRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); + try { + snmpTransportContext.getSnmpTransportService().onToDeviceRpcRequest(this, toDeviceRequest); + snmpTransportContext.getTransportService().process(getSessionInfo(), toDeviceRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); + } catch (Exception e) { + snmpTransportContext.getTransportService().errorEvent(getTenantId(), getDeviceId(), SnmpCommunicationSpec.TO_DEVICE_RPC_REQUEST.getLabel(), e); + } } @Override 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 55ed8fb97f..fb7c374d01 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 @@ -26,10 +26,11 @@ public class SnmpTestV2 { device.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); - } +// 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(); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java index cd61d42016..9d52997c46 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java @@ -17,6 +17,9 @@ package org.thingsboard.server.common.transport; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceTransportType; +import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; @@ -139,6 +142,10 @@ public interface TransportService { void reportActivity(SessionInfoProto sessionInfo); + void lifecycleEvent(TenantId tenantId, DeviceId deviceId, ComponentLifecycleEvent eventType, boolean success, Throwable error); + + void errorEvent(TenantId tenantId, DeviceId deviceId, String method, Throwable error); + void deregisterSession(SessionInfoProto sessionInfo); void log(SessionInfoProto sessionInfo, String msg); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 3b827edbae..d33b3cb42f 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -22,6 +22,7 @@ import com.google.gson.Gson; import com.google.gson.JsonObject; import com.google.protobuf.ByteString; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.exception.ExceptionUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.ApplicationEventPublisher; @@ -49,11 +50,12 @@ import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantProfileId; import org.thingsboard.server.common.data.limit.LimitedApi; +import org.thingsboard.server.common.data.notification.rule.trigger.RateLimitsTrigger; +import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.data.rpc.RpcStatus; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.notification.NotificationRuleProcessor; -import org.thingsboard.server.common.data.notification.rule.trigger.RateLimitsTrigger; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.session.SessionMsgType; @@ -817,6 +819,37 @@ public class DefaultTransportService implements TransportService { sessionsToRemove.forEach(sessionsActivity::remove); } + @Override + public void lifecycleEvent(TenantId tenantId, DeviceId deviceId, ComponentLifecycleEvent eventType, boolean success, Throwable error) { + ToCoreMsg msg = ToCoreMsg.newBuilder() + .setLifecycleEventMsg(TransportProtos.LifecycleEventProto.newBuilder() + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setEntityIdMSB(deviceId.getId().getMostSignificantBits()) + .setEntityIdLSB(deviceId.getId().getLeastSignificantBits()) + .setServiceId(serviceInfoProvider.getServiceId()) + .setLcEventType(eventType.name()) + .setSuccess(success) + .setError(error != null ? ExceptionUtils.getStackTrace(error) : "")) + .build(); + sendToCore(tenantId, deviceId, msg, deviceId.getId(), TransportServiceCallback.EMPTY); + } + + @Override + public void errorEvent(TenantId tenantId, DeviceId deviceId, String method, Throwable error) { + ToCoreMsg msg = ToCoreMsg.newBuilder() + .setErrorEventMsg(TransportProtos.ErrorEventProto.newBuilder() + .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) + .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) + .setEntityIdMSB(deviceId.getId().getMostSignificantBits()) + .setEntityIdLSB(deviceId.getId().getLeastSignificantBits()) + .setServiceId(serviceInfoProvider.getServiceId()) + .setMethod(method) + .setError(ExceptionUtils.getStackTrace(error))) + .build(); + sendToCore(tenantId, deviceId, msg, deviceId.getId(), TransportServiceCallback.EMPTY); + } + @Override public SessionMetaData registerSyncSession(TransportProtos.SessionInfoProto sessionInfo, SessionMsgListener listener, long timeout) { SessionMetaData currentSession = new SessionMetaData(sessionInfo, TransportProtos.SessionType.SYNC, listener); @@ -1108,18 +1141,21 @@ public class DefaultTransportService implements TransportService { } protected void sendToDeviceActor(TransportProtos.SessionInfoProto sessionInfo, TransportToDeviceActorMsg toDeviceActorMsg, TransportServiceCallback callback) { - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, getTenantId(sessionInfo), getDeviceId(sessionInfo)); + ToCoreMsg toCoreMsg = ToCoreMsg.newBuilder().setToDeviceActorMsg(toDeviceActorMsg).build(); + sendToCore(getTenantId(sessionInfo), getDeviceId(sessionInfo), toCoreMsg, getRoutingKey(sessionInfo), callback); + } + + private void sendToCore(TenantId tenantId, EntityId entityId, ToCoreMsg msg, UUID routingKey, TransportServiceCallback callback) { + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); if (log.isTraceEnabled()) { - log.trace("[{}][{}] Pushing to topic {} message {}", getTenantId(sessionInfo), getDeviceId(sessionInfo), tpi.getFullTopicName(), toDeviceActorMsg); + log.trace("[{}][{}] Pushing to topic {} message {}", tenantId, entityId, tpi.getFullTopicName(), msg); } + TransportTbQueueCallback transportTbQueueCallback = callback != null ? new TransportTbQueueCallback(callback) : null; tbCoreProducerStats.incrementTotal(); StatsCallback wrappedCallback = new StatsCallback(transportTbQueueCallback, tbCoreProducerStats); - tbCoreMsgProducer.send(tpi, - new TbProtoQueueMsg<>(getRoutingKey(sessionInfo), - ToCoreMsg.newBuilder().setToDeviceActorMsg(toDeviceActorMsg).build()), - wrappedCallback); + tbCoreMsgProducer.send(tpi, new TbProtoQueueMsg<>(routingKey, msg), wrappedCallback); } private void sendToRuleEngine(TenantId tenantId, TbMsg tbMsg, TbQueueCallback callback) { From 0f078b0da2c28386974af64060e1289d8722b848 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Mon, 26 Jun 2023 15:04:31 +0300 Subject: [PATCH 09/14] 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 { From a045718e941587801dcda1b2bc5faacab2a7dcb7 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Mon, 26 Jun 2023 15:05:11 +0300 Subject: [PATCH 10/14] SNMP_BIND_PORT config --- application/src/main/resources/thingsboard.yml | 1 + transport/snmp/src/main/resources/tb-snmp-transport.yml | 1 + 2 files changed, 2 insertions(+) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 1a5443578e..7b2296aa13 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -969,6 +969,7 @@ transport: # value: "${LWM2M_PROTOCOL_STAGE_THREAD_COUNT:4}" snmp: enabled: "${SNMP_ENABLED:true}" + bind_port: "${SNMP_BIND_PORT:1620}" response_processing: # parallelism level for executor (workStealingPool) that is responsible for handling responses from SNMP devices parallelism_level: "${SNMP_RESPONSE_PROCESSING_PARALLELISM_LEVEL:20}" diff --git a/transport/snmp/src/main/resources/tb-snmp-transport.yml b/transport/snmp/src/main/resources/tb-snmp-transport.yml index 5b5fe34af1..6befe63a2c 100644 --- a/transport/snmp/src/main/resources/tb-snmp-transport.yml +++ b/transport/snmp/src/main/resources/tb-snmp-transport.yml @@ -98,6 +98,7 @@ redis: transport: snmp: enabled: "${SNMP_ENABLED:true}" + bind_port: "${SNMP_BIND_PORT:1620}" response_processing: # parallelism level for executor (workStealingPool) that is responsible for handling responses from SNMP devices parallelism_level: "${SNMP_RESPONSE_PROCESSING_PARALLELISM_LEVEL:20}" From 83641e285f92e6087b593427b6415c1094f498b4 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Tue, 27 Jun 2023 11:50:29 +0300 Subject: [PATCH 11/14] Add TODO for SNMP traps processing --- .../common/data/transport/snmp/SnmpMethod.java | 3 ++- .../transport/snmp/service/SnmpTransportService.java | 12 ++++++++++-- 2 files changed, 12 insertions(+), 3 deletions(-) diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/SnmpMethod.java b/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/SnmpMethod.java index 257dd28217..e366cdb8b4 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/SnmpMethod.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/SnmpMethod.java @@ -17,7 +17,8 @@ package org.thingsboard.server.common.data.transport.snmp; public enum SnmpMethod { GET(-96), - SET(-93); + SET(-93), + TRAP(-89); // codes taken from org.snmp4j.PDU class private final int code; 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 2d507df099..38e3077ebb 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 @@ -297,7 +297,14 @@ 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: this won't work properly in a cluster mode, 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(); @@ -331,6 +338,7 @@ public class SnmpTransportService implements TbTransportService, CommandResponde RequestContext requestContext = RequestContext.builder() .communicationSpec(communicationConfig.getSpec()) .responseMappings(communicationConfig.getAllMappings()) + .method(SnmpMethod.TRAP) .build(); responseProcessingExecutor.execute(() -> { @@ -400,7 +408,7 @@ public class SnmpTransportService implements TbTransportService, CommandResponde responseProcessors.put(SnmpCommunicationSpec.TO_SERVER_RPC_REQUEST, (responseData, requestContext, sessionContext) -> { TransportProtos.ToServerRpcRequestMsg toServerRpcRequestMsg = TransportProtos.ToServerRpcRequestMsg.newBuilder() .setRequestId(0) - .setMethodName("TRAP") + .setMethodName(requestContext.getMethod().name()) .setParams(JsonConverter.toJson(responseData)) .build(); transportService.process(sessionContext.getSessionInfo(), toServerRpcRequestMsg, null); From 1aaa017d92b91401a702c4626f9bc3c0ef07d92a Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 7 Sep 2023 13:28:32 +0300 Subject: [PATCH 12/14] Add ports configs for SNMP transport to Docker scripts --- .../server/transport/snmp/service/SnmpTransportService.java | 4 ++-- docker/docker-compose.yml | 2 ++ docker/tb-snmp-transport.env | 1 + 3 files changed, 5 insertions(+), 2 deletions(-) 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 38e3077ebb..0b60db68b1 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 @@ -302,8 +302,8 @@ public class SnmpTransportService implements TbTransportService, CommandResponde * * TODO: add check for host uniqueness when saving device (for backward compatibility - only for the ones using from-device RPC requests) * - * TODO: this won't work properly in a cluster mode, due to load-balancing of requests from devices: - * session might not be on this instance + * 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) { diff --git a/docker/docker-compose.yml b/docker/docker-compose.yml index b4320577c7..78ca337e0a 100644 --- a/docker/docker-compose.yml +++ b/docker/docker-compose.yml @@ -235,6 +235,8 @@ services: tb-snmp-transport: restart: always image: "${DOCKER_REPO}/${SNMP_TRANSPORT_DOCKER_NAME}:${TB_VERSION}" + ports: + - "1620:1620/udp" environment: TB_SERVICE_ID: tb-snmp-transport JAVA_OPTS: "${JAVA_OPTS}" diff --git a/docker/tb-snmp-transport.env b/docker/tb-snmp-transport.env index e2cc39d658..32160de100 100644 --- a/docker/tb-snmp-transport.env +++ b/docker/tb-snmp-transport.env @@ -1,6 +1,7 @@ ZOOKEEPER_ENABLED=true ZOOKEEPER_URL=zookeeper:2181 +SNMP_BIND_PORT=1620 METRICS_ENABLED=true METRICS_ENDPOINTS_EXPOSE=prometheus WEB_APPLICATION_ENABLE=true From a75093faa786ece4c109545ec645f964e8de0cbd Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 7 Sep 2023 13:39:33 +0300 Subject: [PATCH 13/14] Refactor forwardToEventService in core consumer service --- .../service/queue/DefaultTbCoreConsumerService.java | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index db0b7c3809..f6623b62a3 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -24,6 +24,7 @@ import org.springframework.boot.context.event.ApplicationReadyEvent; import org.springframework.context.ApplicationEventPublisher; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; +import org.thingsboard.common.util.DonAsynchron; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.actors.ActorSystemContext; @@ -664,12 +665,9 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService callback.onSuccess(), + callback::onFailure); } private void throwNotHandled(Object msg, TbCallback callback) { From a570101aeee5d5a999cc45fbe8b68172872af634 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 7 Sep 2023 18:13:14 +0300 Subject: [PATCH 14/14] Use DbCallbackExecutor in forwardToEventService --- .../server/service/queue/DefaultTbCoreConsumerService.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index f6623b62a3..fd5a252cc5 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -667,7 +667,8 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService callback.onSuccess(), - callback::onFailure); + callback::onFailure, + actorContext.getDbCallbackExecutor()); } private void throwNotHandled(Object msg, TbCallback callback) {