From b53746bda6bf246f5ad9c28f31b746de739f068a Mon Sep 17 00:00:00 2001 From: desoliture Date: Tue, 28 Dec 2021 16:38:05 +0200 Subject: [PATCH] MQTT Gateway API attributes request fix --- .../transport/mqtt/adaptors/JsonMqttAdaptor.java | 12 +++++++++--- .../mqtt/session/GatewayDeviceSessionCtx.java | 6 ++++++ .../mqtt/session/GatewaySessionHandler.java | 6 ++++++ .../common/transport/adaptor/JsonConverter.java | 15 +++++++++------ 4 files changed, 30 insertions(+), 9 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java index dc2b74a8ea..0548fa08fb 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java @@ -34,6 +34,7 @@ import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.transport.adaptor.AdaptorException; import org.thingsboard.server.common.transport.adaptor.JsonConverter; import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.transport.mqtt.session.GatewayDeviceSessionCtx; import org.thingsboard.server.transport.mqtt.session.MqttDeviceAwareSessionContext; import java.nio.charset.Charset; @@ -123,7 +124,12 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { @Override public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException { - return processConvertFromGatewayAttributeResponseMsg(ctx, deviceName, responseMsg); + JsonObject request = ((GatewayDeviceSessionCtx) ctx) + .getPendingAttributesRequests() + .getOrDefault(responseMsg.getRequestId(), new JsonObject()); + boolean multipleAttrKeysRequested = + request.has("keys") && request.get("keys").getAsJsonArray().size() > 1; + return processConvertFromGatewayAttributeResponseMsg(ctx, deviceName, responseMsg, multipleAttrKeysRequested); } @Override @@ -232,11 +238,11 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { } } - private Optional processConvertFromGatewayAttributeResponseMsg(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException { + private Optional processConvertFromGatewayAttributeResponseMsg(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg, boolean multipleAttrKeysRequested) throws AdaptorException { if (!StringUtils.isEmpty(responseMsg.getError())) { throw new AdaptorException(responseMsg.getError()); } else { - JsonObject result = JsonConverter.getJsonObjectForGateway(deviceName, responseMsg); + JsonObject result = JsonConverter.getJsonObjectForGateway(deviceName, responseMsg, multipleAttrKeysRequested); return Optional.of(createMqttPublishMsg(ctx, MqttTopics.GATEWAY_ATTRIBUTES_RESPONSE_TOPIC, result)); } } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java index f41c5668cd..593205b745 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.transport.mqtt.session; +import com.google.gson.JsonObject; import io.netty.channel.ChannelFuture; import io.netty.handler.codec.mqtt.MqttMessage; import lombok.extern.slf4j.Slf4j; @@ -27,6 +28,7 @@ import org.thingsboard.server.common.transport.auth.TransportDeviceInfo; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto; +import java.util.Map; import java.util.UUID; import java.util.concurrent.ConcurrentMap; @@ -136,6 +138,10 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple // This feature is not supported in the TB IoT Gateway yet. } + public Map getPendingAttributesRequests() { + return parent.getPendingAttributesRequests(); + } + private boolean isAckExpected(MqttMessage message) { return message.fixedHeader().qosLevel().value() > 0; } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java index 7b45b5f9d9..4b4e986d5f 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java @@ -89,6 +89,7 @@ public class GatewaySessionHandler { private final ConcurrentMap mqttQoSMap; private final ChannelHandlerContext channel; private final DeviceSessionCtx deviceSessionCtx; + private final Map pendingAttributesRequests = new ConcurrentHashMap<>(); public GatewaySessionHandler(DeviceSessionCtx deviceSessionCtx, UUID sessionId) { this.context = deviceSessionCtx.getContext(); @@ -107,6 +108,10 @@ public class GatewaySessionHandler { return new ConcurrentReferenceHashMap<>(16, ReferenceType.WEAK); } + public Map getPendingAttributesRequests () { + return pendingAttributesRequests; + } + public void onDeviceConnect(MqttPublishMessage mqttMsg) throws AdaptorException { if (isJsonPayloadType()) { onDeviceConnectJson(mqttMsg); @@ -558,6 +563,7 @@ public class GatewaySessionHandler { if (json.isJsonObject()) { JsonObject jsonObj = json.getAsJsonObject(); int requestId = jsonObj.get("id").getAsInt(); + pendingAttributesRequests.put(requestId, jsonObj); String deviceName = jsonObj.get(DEVICE_PROPERTY).getAsString(); boolean clientScope = jsonObj.get("client").getAsBoolean(); Set keys; 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 be4143e388..86f90d786f 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 @@ -313,16 +313,19 @@ public class JsonConverter { return result; } - public static JsonObject getJsonObjectForGateway(String deviceName, TransportProtos.GetAttributeResponseMsg - responseMsg) { + public static JsonObject getJsonObjectForGateway( + String deviceName, + TransportProtos.GetAttributeResponseMsg responseMsg, + boolean multipleAttrKeysRequested + ) { JsonObject result = new JsonObject(); result.addProperty("id", responseMsg.getRequestId()); result.addProperty(DEVICE_PROPERTY, deviceName); if (responseMsg.getClientAttributeListCount() > 0) { - addValues(result, responseMsg.getClientAttributeListList()); + addValues(result, responseMsg.getClientAttributeListList(), multipleAttrKeysRequested); } if (responseMsg.getSharedAttributeListCount() > 0) { - addValues(result, responseMsg.getSharedAttributeListList()); + addValues(result, responseMsg.getSharedAttributeListList(), multipleAttrKeysRequested); } return result; } @@ -335,8 +338,8 @@ public class JsonConverter { return result; } - private static void addValues(JsonObject result, List kvList) { - if (kvList.size() == 1) { + private static void addValues(JsonObject result, List kvList, boolean multipleAttrKeysRequested) { + if (kvList.size() == 1 && !multipleAttrKeysRequested) { addValueToJson(result, "value", kvList.get(0).getKv()); } else { JsonObject values;