Browse Source

Refactor

pull/4455/head
Viacheslav Klimov 5 years ago
parent
commit
415bf570ba
  1. 26
      common/data/src/main/java/org/thingsboard/server/common/data/device/data/SnmpDeviceTransportConfiguration.java
  2. 18
      common/data/src/main/java/org/thingsboard/server/common/data/device/profile/SnmpDeviceProfileTransportConfiguration.java
  3. 2
      common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/configs/RepeatingQueryingSnmpCommunicationConfig.java
  4. 5
      common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/configs/SharedAttributesSettingSnmpCommunicationConfig.java
  5. 8
      common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/configs/SnmpCommunicationConfig.java
  6. 3
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpAuthService.java
  7. 54
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java
  8. 27
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java
  9. 14
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java

26
common/data/src/main/java/org/thingsboard/server/common/data/device/data/SnmpDeviceTransportConfiguration.java

@ -17,16 +17,21 @@ package org.thingsboard.server.common.data.device.data;
import com.fasterxml.jackson.annotation.JsonIgnore;
import lombok.Data;
import lombok.ToString;
import org.apache.commons.lang3.ObjectUtils;
import org.apache.commons.lang3.StringUtils;
import org.thingsboard.server.common.data.DeviceTransportType;
import org.thingsboard.server.common.data.transport.snmp.AuthenticationProtocol;
import org.thingsboard.server.common.data.transport.snmp.PrivacyProtocol;
import org.thingsboard.server.common.data.transport.snmp.SnmpProtocolVersion;
import java.util.Objects;
@Data
@ToString(of = {"host", "port", "protocolVersion"})
public class SnmpDeviceTransportConfiguration implements DeviceTransportConfiguration {
private String address;
private int port;
private String host;
private Integer port;
private SnmpProtocolVersion protocolVersion;
/*
@ -60,6 +65,21 @@ public class SnmpDeviceTransportConfiguration implements DeviceTransportConfigur
@JsonIgnore
private boolean isValid() {
return true;
boolean isValid = StringUtils.isNotBlank(host) && port != null && protocolVersion != null;
if (isValid) {
switch (protocolVersion) {
case V1:
case V2C:
isValid = StringUtils.isNotEmpty(community);
break;
case V3:
isValid = StringUtils.isNotBlank(username) && StringUtils.isNotBlank(securityName)
&& contextName != null && authenticationProtocol != null
&& StringUtils.isNotBlank(authenticationPassphrase)
&& privacyProtocol != null && privacyPassphrase != null && engineId != null;
break;
}
}
return isValid;
}
}

18
common/data/src/main/java/org/thingsboard/server/common/data/device/profile/SnmpDeviceProfileTransportConfiguration.java

@ -17,15 +17,21 @@ package org.thingsboard.server.common.data.device.profile;
import com.fasterxml.jackson.annotation.JsonIgnore;
import lombok.Data;
import org.apache.commons.lang3.ArrayUtils;
import org.thingsboard.server.common.data.DeviceTransportType;
import org.thingsboard.server.common.data.transport.snmp.SnmpMapping;
import org.thingsboard.server.common.data.transport.snmp.configs.SnmpCommunicationConfig;
import java.util.Collections;
import java.util.List;
import java.util.function.Function;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@Data
public class SnmpDeviceProfileTransportConfiguration implements DeviceProfileTransportConfiguration {
private int timeoutMs;
private int retries;
private Integer timeoutMs;
private Integer retries;
private List<SnmpCommunicationConfig> communicationConfigs;
@Override
@ -36,12 +42,16 @@ public class SnmpDeviceProfileTransportConfiguration implements DeviceProfileTra
@Override
public void validate() {
if (!isValid()) {
throw new IllegalArgumentException("Transport configuration is not valid");
throw new IllegalArgumentException("SNMP transport configuration is not valid");
}
}
@JsonIgnore
private boolean isValid() {
return true;
return timeoutMs != null && timeoutMs >= 0 && retries != null && retries >= 0
&& communicationConfigs != null && !communicationConfigs.isEmpty()
&& communicationConfigs.stream().allMatch(config -> config != null && config.isValid())
&& communicationConfigs.stream().flatMap(config -> config.getMappings().stream()).map(SnmpMapping::getOid)
.distinct().count() == communicationConfigs.stream().mapToInt(config -> config.getMappings().size()).sum();
}
}

2
common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/configs/RepeatingQueryingSnmpCommunicationConfig.java

@ -31,6 +31,6 @@ public abstract class RepeatingQueryingSnmpCommunicationConfig extends SnmpCommu
@Override
public boolean isValid() {
return true;
return super.isValid() && queryingFrequencyMs != null && queryingFrequencyMs > 0;
}
}

5
common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/configs/SharedAttributesSettingSnmpCommunicationConfig.java

@ -28,9 +28,4 @@ public class SharedAttributesSettingSnmpCommunicationConfig extends SnmpCommunic
public SnmpMethod getMethod() {
return SnmpMethod.SET;
}
@Override
public boolean isValid() {
return true;
}
}

8
common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/configs/SnmpCommunicationConfig.java

@ -49,12 +49,6 @@ public abstract class SnmpCommunicationConfig {
@JsonIgnore
public boolean isValid() {
return true;
}
public void validate() {
if (!isValid()) {
throw new IllegalArgumentException("Communication config is not valid");
}
return mappings != null && !mappings.isEmpty() && mappings.stream().allMatch(mapping -> mapping != null && mapping.isValid());
}
}

3
common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpAuthService.java

@ -24,7 +24,6 @@ import org.snmp4j.security.SecurityLevel;
import org.snmp4j.security.SecurityModel;
import org.snmp4j.security.SecurityProtocols;
import org.snmp4j.security.USM;
import org.snmp4j.security.UsmUser;
import org.snmp4j.smi.GenericAddress;
import org.snmp4j.smi.OID;
import org.snmp4j.smi.OctetString;
@ -98,7 +97,7 @@ public class SnmpAuthService {
throw new UnsupportedOperationException("SNMP protocol version " + protocolVersion + " is not supported");
}
target.setAddress(GenericAddress.parse(snmpUnderlyingProtocol + ":" + deviceTransportConfig.getAddress() + "/" + deviceTransportConfig.getPort()));
target.setAddress(GenericAddress.parse(snmpUnderlyingProtocol + ":" + deviceTransportConfig.getHost() + "/" + deviceTransportConfig.getPort()));
target.setTimeout(profileTransportConfig.getTimeoutMs());
target.setRetries(profileTransportConfig.getRetries());
target.setVersion(protocolVersion.getCode());

54
common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java

@ -89,18 +89,12 @@ public class SnmpTransportContext extends TransportContext {
managedDevicesIds.stream()
.map(protoEntityService::getDeviceById)
.collect(Collectors.toList())
.forEach(device -> {
try {
establishDeviceSession(device);
} catch (Exception e) {
log.error("Failed to establish session for SNMP device {}: {}", device.getId(), e.getMessage());
}
});
.forEach(this::establishDeviceSession);
}
private void establishDeviceSession(Device device) {
if (device == null) return;
log.info("Establishing SNMP device session for device {}", device.getId());
log.info("Establishing SNMP session for device {}", device.getId());
DeviceProfileId deviceProfileId = device.getDeviceProfileId();
DeviceProfile deviceProfile = deviceProfileCache.get(deviceProfileId);
@ -114,18 +108,24 @@ public class SnmpTransportContext extends TransportContext {
SnmpDeviceProfileTransportConfiguration profileTransportConfiguration = (SnmpDeviceProfileTransportConfiguration) deviceProfile.getProfileData().getTransportConfiguration();
SnmpDeviceTransportConfiguration deviceTransportConfiguration = (SnmpDeviceTransportConfiguration) device.getDeviceData().getTransportConfiguration();
DeviceSessionContext deviceSessionContext = new DeviceSessionContext(
device, deviceProfile, credentials.getCredentialsId(),
profileTransportConfiguration, deviceTransportConfiguration, this
);
registerSessionMsgListener(deviceSessionContext);
DeviceSessionContext deviceSessionContext;
try {
deviceSessionContext = new DeviceSessionContext(
device, deviceProfile, credentials.getCredentialsId(),
profileTransportConfiguration, deviceTransportConfiguration, this
);
registerSessionMsgListener(deviceSessionContext);
} catch (Exception e) {
log.error("Failed to establish session for SNMP device {}: {}", device.getId(), e.getMessage());
return;
}
sessions.put(device.getId(), deviceSessionContext);
snmpTransportService.createQueryingTasks(deviceSessionContext);
log.info("Established SNMP device session for device {}", device.getId());
}
private void updateDeviceSession(DeviceSessionContext sessionContext, Device device, DeviceProfile deviceProfile) {
log.info("Updating SNMP device session for device {}", device.getId());
log.info("Updating SNMP session for device {}", device.getId());
DeviceCredentials credentials = protoEntityService.getDeviceCredentialsByDeviceId(device.getId());
if (credentials.getCredentialsType() != DeviceCredentialsType.ACCESS_TOKEN) {
@ -137,16 +137,20 @@ public class SnmpTransportContext extends TransportContext {
SnmpDeviceProfileTransportConfiguration newProfileTransportConfiguration = (SnmpDeviceProfileTransportConfiguration) deviceProfile.getProfileData().getTransportConfiguration();
SnmpDeviceTransportConfiguration newDeviceTransportConfiguration = (SnmpDeviceTransportConfiguration) device.getDeviceData().getTransportConfiguration();
if (!newProfileTransportConfiguration.equals(sessionContext.getProfileTransportConfiguration())) {
sessionContext.setProfileTransportConfiguration(newProfileTransportConfiguration);
sessionContext.initializeTarget(newProfileTransportConfiguration, newDeviceTransportConfiguration);
snmpTransportService.cancelQueryingTasks(sessionContext);
snmpTransportService.createQueryingTasks(sessionContext);
} else if (!newDeviceTransportConfiguration.equals(sessionContext.getDeviceTransportConfiguration())) {
sessionContext.setDeviceTransportConfiguration(newDeviceTransportConfiguration);
sessionContext.initializeTarget(newProfileTransportConfiguration, newDeviceTransportConfiguration);
} else {
log.trace("Configuration of the device {} was not updated", device);
try {
if (!newProfileTransportConfiguration.equals(sessionContext.getProfileTransportConfiguration())) {
sessionContext.setProfileTransportConfiguration(newProfileTransportConfiguration);
sessionContext.initializeTarget(newProfileTransportConfiguration, newDeviceTransportConfiguration);
snmpTransportService.cancelQueryingTasks(sessionContext);
snmpTransportService.createQueryingTasks(sessionContext);
} else if (!newDeviceTransportConfiguration.equals(sessionContext.getDeviceTransportConfiguration())) {
sessionContext.setDeviceTransportConfiguration(newDeviceTransportConfiguration);
sessionContext.initializeTarget(newProfileTransportConfiguration, newDeviceTransportConfiguration);
} else {
log.trace("Configuration of the device {} was not updated", device);
}
} catch (Exception e) {
log.error("Failed to update session for SNMP device {}: {}", sessionContext.getDeviceId(), e.getMessage());
}
}
@ -156,8 +160,8 @@ public class SnmpTransportContext extends TransportContext {
sessionContext.close();
snmpAuthService.cleanUpSnmpAuthInfo(sessionContext);
transportService.deregisterSession(sessionContext.getSessionInfo());
sessions.remove(sessionContext.getDeviceId());
snmpTransportService.cancelQueryingTasks(sessionContext);
sessions.remove(sessionContext.getDeviceId());
log.trace("Unregistered and removed session");
}

27
common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java

@ -28,9 +28,11 @@ import org.snmp4j.mp.MPv3;
import org.snmp4j.security.SecurityModels;
import org.snmp4j.security.SecurityProtocols;
import org.snmp4j.security.USM;
import org.snmp4j.smi.Integer32;
import org.snmp4j.smi.Null;
import org.snmp4j.smi.OID;
import org.snmp4j.smi.OctetString;
import org.snmp4j.smi.Variable;
import org.snmp4j.smi.VariableBinding;
import org.snmp4j.transport.DefaultTcpTransportMapping;
import org.snmp4j.transport.DefaultUdpTransportMapping;
@ -143,7 +145,7 @@ public class SnmpTransportService implements TbTransportService {
sendRequest(sessionContext, communicationConfig);
}
} catch (Exception e) {
log.error("Failed to send SNMP request for device {}: {}", sessionContext.getDeviceId(), e.getMessage());
log.error("Failed to send SNMP request for device {}: {}", sessionContext.getDeviceId(), e.toString());
}
}, queryingFrequency, queryingFrequency, TimeUnit.MILLISECONDS);
}
@ -192,7 +194,24 @@ public class SnmpTransportService implements TbTransportService {
pdu.addAll(communicationConfig.getMappings().stream()
.filter(mapping -> values.isEmpty() || values.containsKey(mapping.getKey()))
.map(mapping -> Optional.ofNullable(values.get(mapping.getKey()))
.map(value -> new VariableBinding(new OID(mapping.getOid()), new OctetString(values.get(mapping.getKey()))))
.map(value -> {
Variable variable;
switch (mapping.getDataType()) {
case LONG:
try {
variable = new Integer32(Integer.parseInt(value));
break;
} catch (NumberFormatException ignored) {
}
case DOUBLE:
case BOOLEAN:
case STRING:
case JSON:
default:
variable = new OctetString(value);
}
return new VariableBinding(new OID(mapping.getOid()), variable);
})
.orElseGet(() -> new VariableBinding(new OID(mapping.getOid()))))
.collect(Collectors.toList()));
@ -267,7 +286,9 @@ public class SnmpTransportService implements TbTransportService {
responses.forEach((spec, response) -> {
Optional.ofNullable(responseProcessors.get(spec))
.ifPresent(responseProcessor -> {
responseProcessor.accept(response, sessionContext);
if (!response.entrySet().isEmpty()) {
responseProcessor.accept(response, sessionContext);
}
});
});

14
common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java

@ -63,8 +63,6 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S
private final Device device;
private final SnmpTransportContext snmpTransportContext;
private final SnmpTransportService snmpTransportService;
private final SnmpAuthService snmpAuthService;
@Getter
@Setter
@ -80,7 +78,7 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S
public DeviceSessionContext(Device device, DeviceProfile deviceProfile, String token,
SnmpDeviceProfileTransportConfiguration profileTransportConfiguration,
SnmpDeviceTransportConfiguration deviceTransportConfiguration,
SnmpTransportContext snmpTransportContext) {
SnmpTransportContext snmpTransportContext) throws Exception {
super(UUID.randomUUID());
super.setDeviceId(device.getId());
super.setDeviceProfile(deviceProfile);
@ -88,8 +86,6 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S
this.token = token;
this.snmpTransportContext = snmpTransportContext;
this.snmpTransportService = snmpTransportContext.getSnmpTransportService();
this.snmpAuthService = snmpTransportContext.getSnmpAuthService();
this.profileTransportConfiguration = profileTransportConfiguration;
this.deviceTransportConfiguration = deviceTransportConfiguration;
@ -113,13 +109,13 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S
@Override
public void onResponse(ResponseEvent event) {
if (isActive) {
snmpTransportService.processResponseEvent(this, event);
snmpTransportContext.getSnmpTransportService().processResponseEvent(this, event);
}
}
public void initializeTarget(SnmpDeviceProfileTransportConfiguration profileTransportConfig, SnmpDeviceTransportConfiguration deviceTransportConfig) {
public void initializeTarget(SnmpDeviceProfileTransportConfiguration profileTransportConfig, SnmpDeviceTransportConfiguration deviceTransportConfig) throws Exception {
log.trace("Initializing target for SNMP session of device {}", device);
this.target = snmpAuthService.setUpSnmpTarget(profileTransportConfig, deviceTransportConfig);
this.target = snmpTransportContext.getSnmpAuthService().setUpSnmpTarget(profileTransportConfig, deviceTransportConfig);
log.info("SNMP target initialized: {}", target);
}
@ -152,7 +148,7 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S
entry -> entry.getValue().isJsonPrimitive() ? entry.getValue().getAsString() : entry.getValue().toString()
));
try {
snmpTransportService.sendRequest(this, communicationConfig, sharedAttributes);
snmpTransportContext.getSnmpTransportService().sendRequest(this, communicationConfig, sharedAttributes);
} catch (Exception e) {
log.error("Failed to send request with shared attributes to SNMP device {}: {}", getDeviceId(), e.getMessage());
}

Loading…
Cancel
Save