diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/data/SnmpDeviceTransportConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/data/SnmpDeviceTransportConfiguration.java index 68a9d1f218..a7bc143d81 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/data/SnmpDeviceTransportConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/data/SnmpDeviceTransportConfiguration.java @@ -40,7 +40,7 @@ public class SnmpDeviceTransportConfiguration implements DeviceTransportConfigur private String community; /* - * For SNMP v3 with User Based Security Model + * For SNMP v3 * */ private String username; private String securityName; 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 bbf074eeda..0b8efaf8c7 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 @@ -44,7 +44,7 @@ public class SnmpDeviceProfileTransportConfiguration implements DeviceProfileTra @JsonIgnore private boolean isValid() { return timeoutMs != null && timeoutMs >= 0 && retries != null && retries >= 0 - && communicationConfigs != null && !communicationConfigs.isEmpty() + && 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(); 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 a6643ecf1e..8d87144ae8 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 @@ -21,6 +21,5 @@ public enum SnmpCommunicationSpec { CLIENT_ATTRIBUTES_QUERYING, SHARED_ATTRIBUTES_SETTING, - TO_DEVICE_RPC_COMMAND_SETTING, - TO_DEVICE_RPC_RESPONSE_QUERYING + TO_DEVICE_RPC_REQUEST, } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/impl/ToDeviceRpcCommandSettingSnmpCommunicationConfig.java b/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/impl/ToDeviceRpcCommandSettingSnmpCommunicationConfig.java deleted file mode 100644 index 97f3de047b..0000000000 --- a/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/impl/ToDeviceRpcCommandSettingSnmpCommunicationConfig.java +++ /dev/null @@ -1,59 +0,0 @@ -/** - * Copyright © 2016-2021 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.impl; - -import lombok.Data; -import org.thingsboard.server.common.data.kv.DataType; -import org.thingsboard.server.common.data.transport.snmp.SnmpCommunicationSpec; -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.config.SnmpCommunicationConfig; - -import java.util.Arrays; -import java.util.Collections; -import java.util.List; - -@Data -public class ToDeviceRpcCommandSettingSnmpCommunicationConfig implements SnmpCommunicationConfig { - private SnmpMapping mapping; - - @Override - public SnmpCommunicationSpec getSpec() { - return SnmpCommunicationSpec.TO_DEVICE_RPC_COMMAND_SETTING; - } - - @Override - public SnmpMethod getMethod() { - return SnmpMethod.SET; - } - - public void setMapping(SnmpMapping mapping) { - this.mapping = mapping != null ? new SnmpMapping(mapping.getOid(), RPC_COMMAND_KEY_NAME, DataType.STRING) : null; - } - - @Override - public List getAllMappings() { - return Collections.singletonList(mapping); - } - - @Override - public boolean isValid() { - return mapping != null && mapping.isValid(); - } - - public static final String RPC_COMMAND_KEY_NAME = "rpcCommand"; - -} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/impl/ToDeviceRpcResponseQueryingSnmpCommunicationConfig.java b/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/impl/ToDeviceRpcResponseQueryingSnmpCommunicationConfig.java deleted file mode 100644 index 18bce94980..0000000000 --- a/common/data/src/main/java/org/thingsboard/server/common/data/transport/snmp/config/impl/ToDeviceRpcResponseQueryingSnmpCommunicationConfig.java +++ /dev/null @@ -1,60 +0,0 @@ -/** - * Copyright © 2016-2021 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.impl; - -import lombok.Data; -import lombok.EqualsAndHashCode; -import org.thingsboard.server.common.data.kv.DataType; -import org.thingsboard.server.common.data.transport.snmp.SnmpCommunicationSpec; -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.config.RepeatingQueryingSnmpCommunicationConfig; - -import java.util.Collections; -import java.util.List; - -@EqualsAndHashCode(callSuper = true) -@Data -public class ToDeviceRpcResponseQueryingSnmpCommunicationConfig extends RepeatingQueryingSnmpCommunicationConfig { - private SnmpMapping mapping; - - @Override - public SnmpCommunicationSpec getSpec() { - return SnmpCommunicationSpec.TO_DEVICE_RPC_RESPONSE_QUERYING; - } - - @Override - public SnmpMethod getMethod() { - return SnmpMethod.GET; - } - - public void setMapping(SnmpMapping mapping) { - this.mapping = mapping != null ? new SnmpMapping(mapping.getOid(), RPC_RESPONSE_KEY_NAME, DataType.STRING) : null; - } - - @Override - public List getAllMappings() { - return Collections.singletonList(mapping); - } - - @Override - public boolean isValid() { - return true; - } - - public static final String RPC_RESPONSE_KEY_NAME = "rpcResponse"; - -} 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 d66d09ab0b..68d4dd3933 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 @@ -152,6 +152,7 @@ public class SnmpTransportContext extends TransportContext { } } catch (Exception e) { log.error("Failed to update session for SNMP device {}: {}", sessionContext.getDeviceId(), e.getMessage()); + destroyDeviceSession(sessionContext); } } diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduMapper.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java similarity index 76% rename from common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduMapper.java rename to common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java index a3b5f77dfb..720b1e7ae3 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduMapper.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/service/PduService.java @@ -29,6 +29,7 @@ import org.springframework.stereotype.Service; 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.queue.util.TbSnmpTransportComponent; @@ -45,32 +46,16 @@ import java.util.stream.IntStream; @TbSnmpTransportComponent @Service @Slf4j -public class PduMapper { +public class PduService { public PDU createPdu(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig, Map values) { - PDU pdu; - SnmpDeviceTransportConfiguration deviceTransportConfiguration = sessionContext.getDeviceTransportConfiguration(); - SnmpProtocolVersion snmpVersion = deviceTransportConfiguration.getProtocolVersion(); - switch (snmpVersion) { - case V1: - case V2C: - pdu = new PDU(); - break; - case V3: - ScopedPDU scopedPdu = new ScopedPDU(); - scopedPdu.setContextName(new OctetString(deviceTransportConfiguration.getContextName())); - scopedPdu.setContextEngineID(new OctetString(deviceTransportConfiguration.getEngineId())); - pdu = scopedPdu; - break; - default: - throw new UnsupportedOperationException("SNMP version " + snmpVersion + " is not supported"); - } + 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(mapping, value); + Variable variable = toSnmpVariable(value, mapping.getDataType()); return new VariableBinding(new OID(mapping.getOid()), variable); }) .orElseGet(() -> new VariableBinding(new OID(mapping.getOid())))) @@ -79,9 +64,20 @@ public class PduMapper { return pdu; } - private Variable toSnmpVariable(SnmpMapping mapping, String value) { + public PDU createSingleVariablePdu(DeviceSessionContext sessionContext, SnmpMethod snmpMethod, String oid, String value, DataType dataType) { + PDU pdu = setUpPdu(sessionContext); + pdu.setType(snmpMethod.getCode()); + + Variable variable = value == null ? Null.instance : toSnmpVariable(value, dataType); + pdu.add(new VariableBinding(new OID(oid), variable)); + + return pdu; + } + + private Variable toSnmpVariable(String value, DataType dataType) { + dataType = dataType == null ? DataType.STRING : dataType; Variable variable; - switch (mapping.getDataType()) { + switch (dataType) { case LONG: try { variable = new Integer32(Integer.parseInt(value)); @@ -98,37 +94,62 @@ public class PduMapper { return variable; } + private PDU setUpPdu(DeviceSessionContext sessionContext) { + PDU pdu; + SnmpDeviceTransportConfiguration deviceTransportConfiguration = sessionContext.getDeviceTransportConfiguration(); + SnmpProtocolVersion snmpVersion = deviceTransportConfiguration.getProtocolVersion(); + switch (snmpVersion) { + case V1: + case V2C: + pdu = new PDU(); + break; + case V3: + ScopedPDU scopedPdu = new ScopedPDU(); + scopedPdu.setContextName(new OctetString(deviceTransportConfiguration.getContextName())); + scopedPdu.setContextEngineID(new OctetString(deviceTransportConfiguration.getEngineId())); + pdu = scopedPdu; + break; + default: + throw new UnsupportedOperationException("SNMP version " + snmpVersion + " is not supported"); + } + return pdu; + } - public JsonObject processPdu(PDU pdu, DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig) { - List variablesBindings = IntStream.range(0, pdu.size()) - .mapToObj(pdu::get) - .filter(Objects::nonNull) - .filter(variableBinding -> !(variableBinding.getVariable() instanceof Null)) - .collect(Collectors.toList()); - JsonObject data = new JsonObject(); + + public JsonObject processPdu(PDU pdu, List responseMappings) { + Map values = processPdu(pdu); Map mappings = new HashMap<>(); - for (SnmpMapping mapping : communicationConfig.getAllMappings()) { - OID oid = new OID(mapping.getOid()); - mappings.put(oid, mapping); + if (responseMappings != null) { + for (SnmpMapping mapping : responseMappings) { + OID oid = new OID(mapping.getOid()); + mappings.put(oid, mapping); + } } - variablesBindings.forEach(variableBinding -> { - log.trace("Processing variable binding: {}", variableBinding); + JsonObject data = new JsonObject(); + values.forEach((oid, value) -> { + log.trace("Processing variable binding: {} - {}", oid, value); - OID oid = variableBinding.getOid(); SnmpMapping mapping = mappings.get(oid); if (mapping == null) { log.debug("No SNMP mapping for oid {}", oid); return; } - processValue(mapping.getKey(), mapping.getDataType(), variableBinding.toValueString(), data); + processValue(mapping.getKey(), mapping.getDataType(), value, data); }); return data; } + public Map processPdu(PDU pdu) { + return IntStream.range(0, pdu.size()) + .mapToObj(pdu::get) + .filter(Objects::nonNull) + .filter(variableBinding -> !(variableBinding.getVariable() instanceof Null)) + .collect(Collectors.toMap(VariableBinding::getOid, VariableBinding::toValueString)); + } private void processValue(String key, DataType dataType, String value, JsonObject result) { switch (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 f672e36190..5bfee14c41 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 @@ -15,6 +15,7 @@ */ package org.thingsboard.server.transport.snmp.service; +import com.google.gson.JsonElement; import com.google.gson.JsonObject; import lombok.Data; import lombok.Getter; @@ -35,11 +36,12 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.TbTransportService; -import org.thingsboard.server.common.data.id.DeviceProfileId; +import org.thingsboard.server.common.data.kv.DataType; import org.thingsboard.server.common.data.transport.snmp.SnmpCommunicationSpec; +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.config.RepeatingQueryingSnmpCommunicationConfig; import org.thingsboard.server.common.data.transport.snmp.config.SnmpCommunicationConfig; -import org.thingsboard.server.common.data.transport.snmp.config.impl.ToDeviceRpcResponseQueryingSnmpCommunicationConfig; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.common.transport.adaptor.JsonConverter; @@ -50,10 +52,12 @@ import org.thingsboard.server.transport.snmp.session.DeviceSessionContext; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.io.IOException; +import java.util.Arrays; import java.util.Collections; import java.util.EnumMap; import java.util.List; import java.util.Map; +import java.util.Optional; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; @@ -67,13 +71,14 @@ import java.util.stream.Collectors; @RequiredArgsConstructor public class SnmpTransportService implements TbTransportService { private final TransportService transportService; - private final PduMapper pduMapper; + private final PduService pduService; @Getter private Snmp snmp; private ScheduledExecutorService queryingExecutor; private ExecutorService responseProcessingExecutor; + private final Map responseDataMappers = new EnumMap<>(SnmpCommunicationSpec.class); private final Map responseProcessors = new EnumMap<>(SnmpCommunicationSpec.class); @Value("${transport.snmp.response_processing.parallelism_level}") @@ -87,6 +92,7 @@ public class SnmpTransportService implements TbTransportService { responseProcessingExecutor = Executors.newWorkStealingPool(responseProcessingParallelismLevel); initializeSnmp(); + configureResponseDataMappers(); configureResponseProcessors(); log.info("SNMP transport service initialized"); @@ -138,16 +144,19 @@ public class SnmpTransportService implements TbTransportService { } - public void sendRequest(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig) { + private void sendRequest(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig) { sendRequest(sessionContext, communicationConfig, Collections.emptyMap()); } - public void sendRequest(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig, Map values) { - PDU request = pduMapper.createPdu(sessionContext, communicationConfig, values); + private void sendRequest(DeviceSessionContext sessionContext, SnmpCommunicationConfig communicationConfig, Map values) { + PDU request = pduService.createPdu(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()); - RequestInfo requestInfo = new RequestInfo(sessionContext.getDeviceProfile().getId(), communicationConfig); try { snmp.send(request, sessionContext.getTarget(), requestInfo, sessionContext); } catch (IOException e) { @@ -156,6 +165,39 @@ public class SnmpTransportService implements TbTransportService { } } + public void onAttributeUpdate(DeviceSessionContext sessionContext, TransportProtos.AttributeUpdateNotificationMsg attributeUpdateNotification) { + sessionContext.getProfileTransportConfiguration().getCommunicationConfigs().stream() + .filter(config -> config.getSpec() == SnmpCommunicationSpec.SHARED_ATTRIBUTES_SETTING) + .findFirst() + .ifPresent(communicationConfig -> { + Map sharedAttributes = JsonConverter.toJson(attributeUpdateNotification).entrySet().stream() + .collect(Collectors.toMap( + Map.Entry::getKey, + entry -> entry.getValue().isJsonPrimitive() ? entry.getValue().getAsString() : entry.getValue().toString() + )); + sendRequest(sessionContext, communicationConfig, sharedAttributes); + }); + } + + public void onToDeviceRpcRequest(DeviceSessionContext sessionContext, TransportProtos.ToDeviceRpcRequestMsg toDeviceRpcRequestMsg) { + SnmpMethod snmpMethod = SnmpMethod.valueOf(toDeviceRpcRequestMsg.getMethodName()); + JsonObject params = JsonConverter.parse(toDeviceRpcRequestMsg.getParams()).getAsJsonObject(); + + String oid = Optional.ofNullable(params.get("oid")).map(JsonElement::getAsString).orElse(null); + String value = Optional.ofNullable(params.get("value")).map(JsonElement::getAsString).orElse(null); + DataType dataType = Optional.ofNullable(params.get("dataType")).map(e -> DataType.valueOf(e.getAsString())).orElse(DataType.STRING); + + if (oid == null || oid.isEmpty()) { + throw new IllegalArgumentException("OID in to-device RPC request is not specified"); + } + if (value == null && snmpMethod == SnmpMethod.SET) { + throw new IllegalArgumentException("Value must be specified for SNMP method 'SET'"); + } + + PDU request = pduService.createSingleVariablePdu(sessionContext, snmpMethod, oid, value, dataType); + sendRequest(sessionContext, request, new RequestInfo(toDeviceRpcRequestMsg.getRequestId(), SnmpCommunicationSpec.TO_DEVICE_RPC_REQUEST)); + } + public void processResponseEvent(DeviceSessionContext sessionContext, ResponseEvent event) { ((Snmp) event.getSource()).cancel(event.getRequest(), sessionContext); @@ -178,47 +220,59 @@ public class SnmpTransportService implements TbTransportService { } private void processResponse(DeviceSessionContext sessionContext, PDU response, RequestInfo requestInfo) { - ResponseProcessor responseProcessor = responseProcessors.get(requestInfo.getCommunicationConfig().getSpec()); - if (responseProcessor == null) { - return; - } + ResponseProcessor responseProcessor = responseProcessors.get(requestInfo.getCommunicationSpec()); + if (responseProcessor == null) return; + + JsonObject responseData = responseDataMappers.get(requestInfo.getCommunicationSpec()).map(response, requestInfo); - JsonObject responseData = pduMapper.processPdu(response, sessionContext, requestInfo.getCommunicationConfig()); if (responseData.entrySet().isEmpty()) { log.debug("No values is the SNMP response for device {}. Request id: {}", sessionContext.getDeviceId(), response.getRequestID()); return; } - responseProcessor.process(responseData, sessionContext); + responseProcessor.process(responseData, requestInfo, sessionContext); reportActivity(sessionContext.getSessionInfo()); } + private void configureResponseDataMappers() { + responseDataMappers.put(SnmpCommunicationSpec.TO_DEVICE_RPC_REQUEST, (pdu, requestInfo) -> { + JsonObject responseData = new JsonObject(); + pduService.processPdu(pdu).forEach((oid, value) -> { + responseData.addProperty(oid.toDottedString(), value); + }); + return responseData; + }); + + ResponseDataMapper defaultResponseDataMapper = (pdu, requestInfo) -> { + return pduService.processPdu(pdu, requestInfo.getResponseMappings()); + }; + Arrays.stream(SnmpCommunicationSpec.values()) + .forEach(communicationSpec -> { + responseDataMappers.putIfAbsent(communicationSpec, defaultResponseDataMapper); + }); + } + private void configureResponseProcessors() { - responseProcessors.put(SnmpCommunicationSpec.TELEMETRY_QUERYING, (response, sessionContext) -> { - TransportProtos.PostTelemetryMsg postTelemetryMsg = JsonConverter.convertToTelemetryProto(response); + responseProcessors.put(SnmpCommunicationSpec.TELEMETRY_QUERYING, (responseData, requestInfo, sessionContext) -> { + TransportProtos.PostTelemetryMsg postTelemetryMsg = JsonConverter.convertToTelemetryProto(responseData); transportService.process(sessionContext.getSessionInfo(), postTelemetryMsg, null); - log.debug("Posted telemetry for SNMP device {}: {}", sessionContext.getDeviceId(), response); + log.debug("Posted telemetry for SNMP device {}: {}", sessionContext.getDeviceId(), responseData); }); - responseProcessors.put(SnmpCommunicationSpec.CLIENT_ATTRIBUTES_QUERYING, (response, sessionContext) -> { - TransportProtos.PostAttributeMsg postAttributesMsg = JsonConverter.convertToAttributesProto(response); + responseProcessors.put(SnmpCommunicationSpec.CLIENT_ATTRIBUTES_QUERYING, (responseData, requestInfo, sessionContext) -> { + TransportProtos.PostAttributeMsg postAttributesMsg = JsonConverter.convertToAttributesProto(responseData); transportService.process(sessionContext.getSessionInfo(), postAttributesMsg, null); - log.debug("Posted attributes for SNMP device {}: {}", sessionContext.getDeviceId(), response); + log.debug("Posted attributes for SNMP device {}: {}", sessionContext.getDeviceId(), responseData); }); - responseProcessors.put(SnmpCommunicationSpec.TO_DEVICE_RPC_RESPONSE_QUERYING, (response, sessionContext) -> { - String rpcResponse = response.get(ToDeviceRpcResponseQueryingSnmpCommunicationConfig.RPC_RESPONSE_KEY_NAME).getAsString(); + responseProcessors.put(SnmpCommunicationSpec.TO_DEVICE_RPC_REQUEST, (responseData, requestInfo, sessionContext) -> { TransportProtos.ToDeviceRpcResponseMsg rpcResponseMsg = TransportProtos.ToDeviceRpcResponseMsg.newBuilder() - .setPayload(rpcResponse) + .setRequestId(requestInfo.getRequestId()) + .setPayload(JsonConverter.toJson(responseData)) .build(); transportService.process(sessionContext.getSessionInfo(), rpcResponseMsg, null); - log.debug("Processed RPC response from device {}: {}", sessionContext.getDeviceId(), rpcResponse); + log.debug("Posted RPC response {} for device {}", responseData, sessionContext.getDeviceId()); }); - -// responseProcessors.put(, (response, sessionContext) -> { -// TransportProtos.ClaimDeviceMsg claimDeviceMsg = JsonConverter.convertToClaimDeviceProto(sessionContext.getDeviceId(), response); -// transportService.process(sessionContext.getSessionInfo(), claimDeviceMsg, null); -// }); } private void reportActivity(TransportProtos.SessionInfoProto sessionInfo) { @@ -256,12 +310,31 @@ public class SnmpTransportService implements TbTransportService { @Data private static class RequestInfo { - private final DeviceProfileId deviceProfileId; - private final SnmpCommunicationConfig communicationConfig; + private Integer requestId; + private SnmpCommunicationSpec communicationSpec; + private List responseMappings; + + public RequestInfo(Integer requestId, SnmpCommunicationSpec communicationSpec) { + this.requestId = requestId; + this.communicationSpec = communicationSpec; + } + + public RequestInfo(SnmpCommunicationSpec communicationSpec) { + this.communicationSpec = communicationSpec; + } + + public RequestInfo(SnmpCommunicationSpec communicationSpec, List responseMappings) { + this.communicationSpec = communicationSpec; + this.responseMappings = responseMappings; + } + } + + private interface ResponseDataMapper { + JsonObject map(PDU pdu, RequestInfo requestInfo); } private interface ResponseProcessor { - void process(JsonObject responseData, DeviceSessionContext sessionContext); + void process(JsonObject responseData, RequestInfo requestInfo, DeviceSessionContext sessionContext); } } 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 82fdc037ad..a59ce020ba 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 @@ -26,11 +26,7 @@ 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.transport.snmp.SnmpCommunicationSpec; -import org.thingsboard.server.common.data.transport.snmp.config.SnmpCommunicationConfig; -import org.thingsboard.server.common.data.transport.snmp.config.impl.ToDeviceRpcCommandSettingSnmpCommunicationConfig; import org.thingsboard.server.common.transport.SessionMsgListener; -import org.thingsboard.server.common.transport.adaptor.JsonConverter; import org.thingsboard.server.common.transport.session.DeviceAwareSessionContext; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg; @@ -40,15 +36,11 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMs import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcResponseMsg; import org.thingsboard.server.transport.snmp.SnmpTransportContext; -import java.io.IOException; import java.util.LinkedList; import java.util.List; -import java.util.Map; -import java.util.Optional; import java.util.UUID; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.atomic.AtomicInteger; -import java.util.stream.Collectors; @Slf4j public class DeviceSessionContext extends DeviceAwareSessionContext implements SessionMsgListener, ResponseListener { @@ -136,15 +128,7 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S @Override public void onAttributeUpdate(AttributeUpdateNotificationMsg attributeUpdateNotification) { - getCommunicationConfigForSpec(SnmpCommunicationSpec.SHARED_ATTRIBUTES_SETTING) - .ifPresent(communicationConfig -> { - Map sharedAttributes = JsonConverter.toJson(attributeUpdateNotification).entrySet().stream() - .collect(Collectors.toMap( - Map.Entry::getKey, - entry -> entry.getValue().isJsonPrimitive() ? entry.getValue().getAsString() : entry.getValue().toString() - )); - snmpTransportContext.getSnmpTransportService().sendRequest(this, communicationConfig, sharedAttributes); - }); + snmpTransportContext.getSnmpTransportService().onAttributeUpdate(this, attributeUpdateNotification); } @Override @@ -153,23 +137,10 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S @Override public void onToDeviceRpcRequest(ToDeviceRpcRequestMsg toDeviceRequest) { - getCommunicationConfigForSpec(SnmpCommunicationSpec.TO_DEVICE_RPC_COMMAND_SETTING) - .ifPresent(communicationConfig -> { - String value = JsonConverter.toJson(toDeviceRequest, true).toString(); - snmpTransportContext.getSnmpTransportService().sendRequest( - this, communicationConfig, - Map.of(ToDeviceRpcCommandSettingSnmpCommunicationConfig.RPC_COMMAND_KEY_NAME, value) - ); - }); + snmpTransportContext.getSnmpTransportService().onToDeviceRpcRequest(this, toDeviceRequest); } @Override public void onToServerRpcResponse(ToServerRpcResponseMsg toServerResponse) { } - - private Optional getCommunicationConfigForSpec(SnmpCommunicationSpec spec) { - return profileTransportConfiguration.getCommunicationConfigs().stream() - .filter(config -> config.getSpec() == spec) - .findFirst(); - } } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java index fe2f3500a6..fce1c2c9b5 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java @@ -558,6 +558,14 @@ public class JsonConverter { } } + public static JsonElement parse(String json) { + return JSON_PARSER.parse(json); + } + + public static String toJson(JsonElement element) { + return GSON.toJson(element); + } + public static void setTypeCastEnabled(boolean enabled) { isTypeCastEnabled = enabled; } @@ -599,8 +607,7 @@ public class JsonConverter { .build(); } - private static TransportProtos.ProvisionDeviceCredentialsMsg buildProvisionDeviceCredentialsMsg(String - provisionKey, String provisionSecret) { + private static TransportProtos.ProvisionDeviceCredentialsMsg buildProvisionDeviceCredentialsMsg(String provisionKey, String provisionSecret) { return TransportProtos.ProvisionDeviceCredentialsMsg.newBuilder() .setProvisionDeviceKey(provisionKey) .setProvisionDeviceSecret(provisionSecret)