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); } }