Browse Source

Single telemetry message when splitting large SNMP request

pull/8757/head
ViacheslavKlimov 3 years ago
parent
commit
c1f5b39cdb
  1. 3
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/SnmpTransportContext.java
  2. 15
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java
  3. 127
      common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/SnmpTransportService.java

3
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) @EventListener(DeviceUpdatedEvent.class)
public void onDeviceUpdatedOrCreated(DeviceUpdatedEvent deviceUpdatedEvent) { public void onDeviceUpdatedOrCreated(DeviceUpdatedEvent deviceUpdatedEvent) {
Device device = deviceUpdatedEvent.getDevice(); 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()) DeviceTransportType transportType = Optional.ofNullable(device.getDeviceData().getTransportConfiguration())
.map(DeviceTransportConfiguration::getType) .map(DeviceTransportConfiguration::getType)
.orElse(null); .orElse(null);
@ -246,6 +246,7 @@ public class SnmpTransportContext extends TransportContext {
} }
public void onDeviceProfileUpdated(DeviceProfile deviceProfile, DeviceSessionContext sessionContext) { 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); updateDeviceSession(sessionContext, sessionContext.getDevice(), deviceProfile);
} }

15
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.Objects;
import java.util.Optional; import java.util.Optional;
import java.util.stream.Collectors; import java.util.stream.Collectors;
import java.util.stream.IntStream;
@TbSnmpTransportComponent @TbSnmpTransportComponent
@Service @Service
@ -70,7 +69,9 @@ public class PduService {
}) })
.orElseGet(() -> new VariableBinding(new OID(mapping.getOid())))) .orElseGet(() -> new VariableBinding(new OID(mapping.getOid()))))
.collect(Collectors.toList())); .collect(Collectors.toList()));
pdus.add(pdu); if (pdu.size() > 0) {
pdus.add(pdu);
}
} }
return pdus; return pdus;
@ -128,8 +129,8 @@ public class PduService {
} }
public JsonObject processPdu(PDU pdu, List<SnmpMapping> responseMappings) { public JsonObject processPdus(List<PDU> pdus, List<SnmpMapping> responseMappings) {
Map<OID, String> values = processPdu(pdu); Map<OID, String> values = processPdus(pdus);
Map<OID, SnmpMapping> mappings = new HashMap<>(); Map<OID, SnmpMapping> mappings = new HashMap<>();
if (responseMappings != null) { if (responseMappings != null) {
@ -155,9 +156,9 @@ public class PduService {
return data; return data;
} }
public Map<OID, String> processPdu(PDU pdu) { public Map<OID, String> processPdus(List<PDU> pdus) {
return IntStream.range(0, pdu.size()) return pdus.stream()
.mapToObj(pdu::get) .flatMap(pdu -> pdu.getVariableBindings().stream())
.filter(Objects::nonNull) .filter(Objects::nonNull)
.filter(variableBinding -> !(variableBinding.getVariable() instanceof Null)) .filter(variableBinding -> !(variableBinding.getVariable() instanceof Null))
.collect(Collectors.toMap(VariableBinding::getOid, VariableBinding::toValueString)); .collect(Collectors.toMap(VariableBinding::getOid, VariableBinding::toValueString));

127
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.JsonElement;
import com.google.gson.JsonObject; import com.google.gson.JsonObject;
import lombok.Builder;
import lombok.Data; import lombok.Data;
import lombok.Getter; import lombok.Getter;
import lombok.RequiredArgsConstructor; import lombok.RequiredArgsConstructor;
@ -53,6 +54,7 @@ import org.thingsboard.server.transport.snmp.session.DeviceSessionContext;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import java.io.IOException; import java.io.IOException;
import java.util.ArrayList;
import java.util.Arrays; import java.util.Arrays;
import java.util.Collections; import java.util.Collections;
import java.util.EnumMap; import java.util.EnumMap;
@ -161,19 +163,21 @@ public class SnmpTransportService implements TbTransportService {
private void sendRequest(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig, Map<String, String> values) { private void sendRequest(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig, Map<String, String> values) {
List<PDU> request = pduService.createPdus(sessionContext, communicationConfig, values); List<PDU> request = pduService.createPdus(sessionContext, communicationConfig, values);
RequestInfo requestInfo = new RequestInfo(communicationConfig.getSpec(), communicationConfig.getAllMappings()); RequestContext requestContext = RequestContext.builder()
sendRequest(sessionContext, request, requestInfo); .communicationSpec(communicationConfig.getSpec())
.responseMappings(communicationConfig.getAllMappings())
.requestSize(request.size())
.build();
sendRequest(sessionContext, request, requestContext);
} }
private void sendRequest(DeviceSessionContext sessionContext, List<PDU> request, RequestInfo requestInfo) { private void sendRequest(DeviceSessionContext sessionContext, List<PDU> request, RequestContext requestContext) {
for (PDU pdu : request) { for (PDU pdu : request) {
if (pdu.size() > 0) { log.trace("Executing SNMP request for device {} with {} variable bindings", sessionContext.getDeviceId(), pdu.size());
log.trace("Executing SNMP request for device {}. Variables bindings: {}", sessionContext.getDeviceId(), pdu.getVariableBindings()); try {
try { snmp.send(pdu, sessionContext.getTarget(), requestContext, sessionContext);
snmp.send(pdu, sessionContext.getTarget(), requestInfo, sessionContext); } catch (IOException e) {
} catch (IOException e) { log.error("Failed to send SNMP request to device {}: {}", sessionContext.getDeviceId(), e.toString());
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(); DataType dataType = snmpMapping.getDataType();
PDU request = pduService.createSingleVariablePdu(sessionContext, snmpMethod, oid, value, dataType); PDU request = pduService.createSingleVariablePdu(sessionContext, snmpMethod, oid, value, dataType);
RequestInfo requestInfo = new RequestInfo(toDeviceRpcRequestMsg.getRequestId(), communicationConfig.getSpec(), communicationConfig.getAllMappings()); RequestContext requestContext = RequestContext.builder()
sendRequest(sessionContext, List.of(request), requestInfo); .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) { public void processResponseEvent(DeviceSessionContext sessionContext, ResponseEvent event) {
((Snmp) event.getSource()).cancel(event.getRequest(), sessionContext); ((Snmp) event.getSource()).cancel(event.getRequest(), sessionContext);
if (event.getError() != null) { if (event.getError() != null) {
log.warn("SNMP response error: {}", event.getError().toString()); log.warn("SNMP response error: {}", event.getError().toString());
return; return;
} }
PDU response = event.getResponse(); PDU responsePdu = event.getResponse();
if (response == null) { if (log.isTraceEnabled()) {
log.debug("No response from SNMP device {}, requestId: {}", sessionContext.getDeviceId(), event.getRequest().getRequestID()); log.trace("Received PDU for device {}: {}", sessionContext.getDeviceId(), responsePdu);
return; }
RequestContext requestContext = (RequestContext) event.getUserObject();
List<PDU> 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<PDU> 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(() -> { responseProcessingExecutor.execute(() -> {
processResponse(sessionContext, response, requestInfo); processResponse(sessionContext, response, requestContext);
}); });
} }
private void processResponse(DeviceSessionContext sessionContext, PDU response, RequestInfo requestInfo) { private void processResponse(DeviceSessionContext sessionContext, List<PDU> response, RequestContext requestContext) {
ResponseProcessor responseProcessor = responseProcessors.get(requestInfo.getCommunicationSpec()); ResponseProcessor responseProcessor = responseProcessors.get(requestContext.getCommunicationSpec());
if (responseProcessor == null) return; if (responseProcessor == null) return;
JsonObject responseData = responseDataMappers.get(requestInfo.getCommunicationSpec()).map(response, requestInfo); JsonObject responseData = responseDataMappers.get(requestContext.getCommunicationSpec()).map(response, requestContext);
if (responseData.size() == 0) {
if (responseData.entrySet().isEmpty()) { log.warn("No values in the SNMP response for device {}", sessionContext.getDeviceId());
log.debug("No values in the SNMP response for device {}. Request id: {}", sessionContext.getDeviceId(), response.getRequestID());
return; return;
} }
responseProcessor.process(responseData, requestInfo, sessionContext); responseProcessor.process(responseData, requestContext, sessionContext);
reportActivity(sessionContext.getSessionInfo()); reportActivity(sessionContext.getSessionInfo());
} }
private void configureResponseDataMappers() { 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(); JsonObject responseData = new JsonObject();
pduService.processPdu(pdu).forEach((oid, value) -> { pduService.processPdus(pdus).forEach((oid, value) -> {
requestInfo.getResponseMappings().stream() requestContext.getResponseMappings().stream()
.filter(snmpMapping -> snmpMapping.getOid().equals(oid.toDottedString())) .filter(snmpMapping -> snmpMapping.getOid().equals(oid.toDottedString()))
.findFirst() .findFirst()
.ifPresent(snmpMapping -> { .ifPresent(snmpMapping -> {
@ -270,8 +300,8 @@ public class SnmpTransportService implements TbTransportService {
return responseData; return responseData;
}); });
ResponseDataMapper defaultResponseDataMapper = (pdu, requestInfo) -> { ResponseDataMapper defaultResponseDataMapper = (pdus, requestContext) -> {
return pduService.processPdu(pdu, requestInfo.getResponseMappings()); return pduService.processPdus(pdus, requestContext.getResponseMappings());
}; };
Arrays.stream(SnmpCommunicationSpec.values()) Arrays.stream(SnmpCommunicationSpec.values())
.forEach(communicationSpec -> { .forEach(communicationSpec -> {
@ -280,21 +310,21 @@ public class SnmpTransportService implements TbTransportService {
} }
private void configureResponseProcessors() { 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); TransportProtos.PostTelemetryMsg postTelemetryMsg = JsonConverter.convertToTelemetryProto(responseData);
transportService.process(sessionContext.getSessionInfo(), postTelemetryMsg, null); transportService.process(sessionContext.getSessionInfo(), postTelemetryMsg, null);
log.debug("Posted telemetry for SNMP device {}: {}", sessionContext.getDeviceId(), responseData); 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); TransportProtos.PostAttributeMsg postAttributesMsg = JsonConverter.convertToAttributesProto(responseData);
transportService.process(sessionContext.getSessionInfo(), postAttributesMsg, null); transportService.process(sessionContext.getSessionInfo(), postAttributesMsg, null);
log.debug("Posted attributes for SNMP device {}: {}", sessionContext.getDeviceId(), responseData); 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() TransportProtos.ToDeviceRpcResponseMsg rpcResponseMsg = TransportProtos.ToDeviceRpcResponseMsg.newBuilder()
.setRequestId(requestInfo.getRequestId()) .setRequestId(requestContext.getRequestId())
.setPayload(JsonConverter.toJson(responseData)) .setPayload(JsonConverter.toJson(responseData))
.build(); .build();
transportService.process(sessionContext.getSessionInfo(), rpcResponseMsg, null); transportService.process(sessionContext.getSessionInfo(), rpcResponseMsg, null);
@ -332,29 +362,32 @@ public class SnmpTransportService implements TbTransportService {
} }
@Data @Data
private static class RequestInfo { private static class RequestContext {
private Integer requestId; private final Integer requestId;
private SnmpCommunicationSpec communicationSpec; private final SnmpCommunicationSpec communicationSpec;
private List<SnmpMapping> responseMappings; private final List<SnmpMapping> responseMappings;
public RequestInfo(Integer requestId, SnmpCommunicationSpec communicationSpec, List<SnmpMapping> responseMappings) { private final int requestSize;
this.requestId = requestId; private List<PDU> responseParts;
this.communicationSpec = communicationSpec;
this.responseMappings = responseMappings;
}
public RequestInfo(SnmpCommunicationSpec communicationSpec, List<SnmpMapping> responseMappings) { @Builder
public RequestContext(Integer requestId, SnmpCommunicationSpec communicationSpec, List<SnmpMapping> responseMappings, int requestSize) {
this.requestId = requestId;
this.communicationSpec = communicationSpec; this.communicationSpec = communicationSpec;
this.responseMappings = responseMappings; this.responseMappings = responseMappings;
this.requestSize = requestSize;
if (requestSize > 1) {
this.responseParts = Collections.synchronizedList(new ArrayList<>());
}
} }
} }
private interface ResponseDataMapper { private interface ResponseDataMapper {
JsonObject map(PDU pdu, RequestInfo requestInfo); JsonObject map(List<PDU> pdus, RequestContext requestContext);
} }
private interface ResponseProcessor { private interface ResponseProcessor {
void process(JsonObject responseData, RequestInfo requestInfo, DeviceSessionContext sessionContext); void process(JsonObject responseData, RequestContext requestContext, DeviceSessionContext sessionContext);
} }
} }

Loading…
Cancel
Save