From b53746bda6bf246f5ad9c28f31b746de739f068a Mon Sep 17 00:00:00 2001 From: desoliture Date: Tue, 28 Dec 2021 16:38:05 +0200 Subject: [PATCH 1/6] 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; From 2f5648c400122b02838ea11d650dc55056dc860c Mon Sep 17 00:00:00 2001 From: desoliture Date: Wed, 29 Dec 2021 17:22:32 +0200 Subject: [PATCH 2/6] refactoring --- .../server/transport/mqtt/adaptors/JsonMqttAdaptor.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 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 0548fa08fb..7aea8f3f91 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 @@ -43,6 +43,7 @@ import java.util.Arrays; import java.util.HashSet; import java.util.Optional; import java.util.Set; +import java.util.Map; import java.util.UUID; import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_SOFTWARE_FIRMWARE_RESPONSES_TOPIC_FORMAT; @@ -124,11 +125,12 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { @Override public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException { - JsonObject request = ((GatewayDeviceSessionCtx) ctx) - .getPendingAttributesRequests() - .getOrDefault(responseMsg.getRequestId(), new JsonObject()); + Map pendingAttributesRequests = ((GatewayDeviceSessionCtx) ctx).getPendingAttributesRequests(); + int requestId = responseMsg.getRequestId(); + JsonObject request = pendingAttributesRequests.getOrDefault(requestId, new JsonObject()); boolean multipleAttrKeysRequested = request.has("keys") && request.get("keys").getAsJsonArray().size() > 1; + pendingAttributesRequests.remove(requestId); return processConvertFromGatewayAttributeResponseMsg(ctx, deviceName, responseMsg, multipleAttrKeysRequested); } From 443bb2280fc62f0fbb4ba023d869e2855b7c233f Mon Sep 17 00:00:00 2001 From: desoliture Date: Fri, 31 Dec 2021 13:27:52 +0200 Subject: [PATCH 3/6] add corresponding test --- .../server/msa/AbstractContainerTest.java | 2 + .../connectivity/MqttGatewayClientTest.java | 71 +++++++++++++++++++ 2 files changed, 73 insertions(+) diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java index 412cf3094b..e7dc6f2c44 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java @@ -22,6 +22,7 @@ import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.collect.ImmutableMap; import com.google.gson.JsonArray; import com.google.gson.JsonObject; +import com.google.gson.JsonParser; import lombok.extern.slf4j.Slf4j; import org.apache.cassandra.cql3.Json; import org.apache.commons.lang3.RandomStringUtils; @@ -69,6 +70,7 @@ public abstract class AbstractContainerTest { protected static String TB_TOKEN; protected static RestClient restClient; protected ObjectMapper mapper = new ObjectMapper(); + protected JsonParser jsonParser = new JsonParser(); @BeforeClass public static void before() throws Exception { diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java index 4e7b7cd05b..dfdc68d9d5 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java @@ -20,6 +20,7 @@ import com.google.common.collect.Sets; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; import com.google.common.util.concurrent.MoreExecutors; +import com.google.gson.JsonArray; import com.google.gson.JsonObject; import com.google.gson.JsonParser; import io.netty.buffer.ByteBuf; @@ -151,6 +152,76 @@ public class MqttGatewayClientTest extends AbstractContainerTest { Assert.assertTrue(verify(actualLatestTelemetry, "attr4", Long.toString(73))); } + @Test + public void responseDataOnAttributesRequestCheck() throws Exception { + Optional createdDeviceCredentials = restClient.getDeviceCredentialsByDeviceId(createdDevice.getId()); + Assert.assertTrue(createdDeviceCredentials.isPresent()); + WsClient wsClient = subscribeToWebSocket(createdDevice.getId(), "CLIENT_SCOPE", CmdsType.ATTR_SUB_CMDS); + JsonObject sharedAttributes = new JsonObject(); + sharedAttributes.addProperty("attr1", "value1"); + sharedAttributes.addProperty("attr2", true); + sharedAttributes.addProperty("attr3", 42.0); + sharedAttributes.addProperty("attr4", 73); + + mqttClient.on("v1/gateway/attributes/response", listener, MqttQoS.AT_LEAST_ONCE).get(); + ResponseEntity sharedAttributesResponse = restClient.getRestTemplate() + .postForEntity(HTTPS_URL + "/api/plugins/telemetry/DEVICE/{deviceId}/SHARED_SCOPE", + mapper.readTree(sharedAttributes.toString()), ResponseEntity.class, + createdDevice.getId()); + Assert.assertTrue(sharedAttributesResponse.getStatusCode().is2xxSuccessful()); + var event = listener.getEvents().poll(10, TimeUnit.SECONDS); + + JsonObject requestData = new JsonObject(); + requestData.addProperty("id", 1); + requestData.addProperty("device", createdDevice.getName()); + requestData.addProperty("client", false); + requestData.addProperty("key", "attr1"); + + mqttClient.on("v1/gateway/attributes/response", listener, MqttQoS.AT_LEAST_ONCE).get(); + mqttClient.publish("v1/gateway/attributes/request", Unpooled.wrappedBuffer(requestData.toString().getBytes())).get(); + event = listener.getEvents().poll(10, TimeUnit.SECONDS); + + JsonObject responseData = jsonParser.parse(Objects.requireNonNull(event).getMessage()).getAsJsonObject(); + Assert.assertTrue(responseData.has("value")); + Assert.assertEquals(sharedAttributes.get("attr1").getAsString(), responseData.get("value").getAsString()); + + requestData = new JsonObject(); + requestData.addProperty("id", 1); + requestData.addProperty("device", createdDevice.getName()); + requestData.addProperty("client", false); + JsonArray keys = new JsonArray(); + keys.add("attr1"); + keys.add("attr2"); + requestData.add("keys", keys); + + mqttClient.on("v1/gateway/attributes/response", listener, MqttQoS.AT_LEAST_ONCE).get(); + mqttClient.publish("v1/gateway/attributes/request", Unpooled.wrappedBuffer(requestData.toString().getBytes())).get(); + event = listener.getEvents().poll(10, TimeUnit.SECONDS); + responseData = jsonParser.parse(Objects.requireNonNull(event).getMessage()).getAsJsonObject(); + + Assert.assertTrue(responseData.has("values")); + Assert.assertEquals(sharedAttributes.get("attr1").getAsString(), responseData.get("values").getAsJsonObject().get("attr1").getAsString()); + Assert.assertEquals(sharedAttributes.get("attr2").getAsString(), responseData.get("values").getAsJsonObject().get("attr2").getAsString()); + + requestData = new JsonObject(); + requestData.addProperty("id", 1); + requestData.addProperty("device", createdDevice.getName()); + requestData.addProperty("client", false); + keys = new JsonArray(); + keys.add("attr1"); + keys.add("undefined"); + requestData.add("keys", keys); + + mqttClient.on("v1/gateway/attributes/response", listener, MqttQoS.AT_LEAST_ONCE).get(); + mqttClient.publish("v1/gateway/attributes/request", Unpooled.wrappedBuffer(requestData.toString().getBytes())).get(); + event = listener.getEvents().poll(10, TimeUnit.SECONDS); + responseData = jsonParser.parse(Objects.requireNonNull(event).getMessage()).getAsJsonObject(); + + Assert.assertTrue(responseData.has("values")); + Assert.assertEquals(sharedAttributes.get("attr1").getAsString(), responseData.get("values").getAsJsonObject().get("attr1").getAsString()); + Assert.assertEquals(1, responseData.get("values").getAsJsonObject().entrySet().size()); + } + @Test public void requestAttributeValuesFromServer() throws Exception { WsClient wsClient = subscribeToWebSocket(createdDevice.getId(), "CLIENT_SCOPE", CmdsType.ATTR_SUB_CMDS); From 35c30b7678e90a7d53b94ca39a891903fa9f840b Mon Sep 17 00:00:00 2001 From: desoliture Date: Thu, 13 Jan 2022 17:48:59 +0200 Subject: [PATCH 4/6] refactor GatewayDeviceSessionCtx and MqttAdaptors refactor GatewayDeviceSessionCtx to determine the value of multipleAttrKeysRequested before calling the JsonMqttAdaptor, add corresponding convertToGatewayPublish method to adaptors interface with multipleAttrKeysRequested parameter, refactor other adaptors to deal with new method --- .../mqtt/adaptors/BackwardCompatibilityAdaptor.java | 5 +++++ .../transport/mqtt/adaptors/JsonMqttAdaptor.java | 13 ++++++------- .../mqtt/adaptors/MqttTransportAdaptor.java | 3 +++ .../transport/mqtt/adaptors/ProtoMqttAdaptor.java | 5 +++++ .../mqtt/session/GatewayDeviceSessionCtx.java | 11 ++++++----- .../mqtt/session/GatewaySessionHandler.java | 7 +++++-- 6 files changed, 30 insertions(+), 14 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java index a608d6c52d..0cc65da1b0 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java @@ -128,6 +128,11 @@ public class BackwardCompatibilityAdaptor implements MqttTransportAdaptor { return protoAdaptor.convertToGatewayPublish(ctx, deviceName, rpcRequest); } + @Override + public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg, boolean multipleAttributeKeysRequested) throws AdaptorException { + return protoAdaptor.convertToGatewayPublish(ctx, deviceName, responseMsg, multipleAttributeKeysRequested); + } + @Override public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase) throws AdaptorException { log.warn("[{}] invoked not implemented adaptor method! ToServerRpcResponseMsg: {} TopicBase: {}", ctx.getSessionId(), rpcResponse, topicBase); 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 7aea8f3f91..d71b87efdc 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 @@ -125,13 +125,12 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { @Override public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException { - Map pendingAttributesRequests = ((GatewayDeviceSessionCtx) ctx).getPendingAttributesRequests(); - int requestId = responseMsg.getRequestId(); - JsonObject request = pendingAttributesRequests.getOrDefault(requestId, new JsonObject()); - boolean multipleAttrKeysRequested = - request.has("keys") && request.get("keys").getAsJsonArray().size() > 1; - pendingAttributesRequests.remove(requestId); - return processConvertFromGatewayAttributeResponseMsg(ctx, deviceName, responseMsg, multipleAttrKeysRequested); + return convertToGatewayPublish(ctx, deviceName, responseMsg, false); + } + + @Override + public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg, boolean multipleAttributeKeysRequested) throws AdaptorException { + return processConvertFromGatewayAttributeResponseMsg(ctx, deviceName, responseMsg, multipleAttributeKeysRequested); } @Override diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java index 53ef439eda..a980e91b53 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java @@ -25,6 +25,7 @@ import io.netty.handler.codec.mqtt.MqttPublishMessage; import io.netty.handler.codec.mqtt.MqttPublishVariableHeader; import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.transport.adaptor.AdaptorException; +import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg; @@ -72,6 +73,8 @@ public interface MqttTransportAdaptor { Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, ToDeviceRpcRequestMsg rpcRequest) throws AdaptorException; + Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg, boolean multipleAttributeKeysRequested) throws AdaptorException; + Optional convertToPublish(MqttDeviceAwareSessionContext ctx, ToServerRpcResponseMsg rpcResponse, String topicBase) throws AdaptorException; ProvisionDeviceRequestMsg convertToProvisionRequestMsg(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException; diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java index 4fa47b367d..4dc9a95d6b 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java @@ -209,6 +209,11 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { return Optional.of(createMqttPublishMsg(ctx, MqttTopics.GATEWAY_RPC_TOPIC, payloadBytes)); } + @Override + public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg, boolean multipleAttributeKeysRequested) throws AdaptorException { + return convertToGatewayPublish(ctx, deviceName, responseMsg); + } + public static byte[] toBytes(ByteBuf inbound) { byte[] bytes = new byte[inbound.readableBytes()]; int readerIndex = inbound.readerIndex(); 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 593205b745..02324fd3e7 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,7 +15,6 @@ */ 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; @@ -28,7 +27,6 @@ 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; @@ -82,7 +80,10 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple @Override public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg response) { try { - parent.getPayloadAdaptor().convertToGatewayPublish(this, getDeviceInfo().getDeviceName(), response).ifPresent(parent::writeAndFlush); + boolean multipleAttrKeysReq = isMultipleAttributeKeysRequested(response.getRequestId()); + parent.getPayloadAdaptor() + .convertToGatewayPublish(this, getDeviceInfo().getDeviceName(), response, multipleAttrKeysReq) + .ifPresent(parent::writeAndFlush); } catch (Exception e) { log.trace("[{}] Failed to convert device attributes response to MQTT msg", sessionId, e); } @@ -138,8 +139,8 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple // This feature is not supported in the TB IoT Gateway yet. } - public Map getPendingAttributesRequests() { - return parent.getPendingAttributesRequests(); + public boolean isMultipleAttributeKeysRequested(int requestId) { + return parent.isMultipleAttributeKeysRequested(requestId); } private boolean isAckExpected(MqttMessage message) { 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 4b4e986d5f..12e9f786bf 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 @@ -108,8 +108,11 @@ public class GatewaySessionHandler { return new ConcurrentReferenceHashMap<>(16, ReferenceType.WEAK); } - public Map getPendingAttributesRequests () { - return pendingAttributesRequests; + public boolean isMultipleAttributeKeysRequested(int requestId) { + JsonObject request = pendingAttributesRequests.getOrDefault(requestId, new JsonObject()); + boolean multipleAttributeKeysRequested = request.has("keys") && request.get("keys").getAsJsonArray().size() > 1; + pendingAttributesRequests.remove(requestId); + return multipleAttributeKeysRequested; } public void onDeviceConnect(MqttPublishMessage mqttMsg) throws AdaptorException { From 7b14fdaf5e0e6901a7a509299df744bb380fd63d Mon Sep 17 00:00:00 2001 From: desoliture Date: Mon, 17 Jan 2022 18:59:39 +0200 Subject: [PATCH 5/6] add isMultipleAttributesRequest field in queue.proto in GetAttributeResponseMsg message --- .../device/DeviceActorMessageProcessor.java | 3 +++ common/cluster-api/src/main/proto/queue.proto | 1 + .../adaptors/BackwardCompatibilityAdaptor.java | 5 ----- .../transport/mqtt/adaptors/JsonMqttAdaptor.java | 15 ++++----------- .../mqtt/adaptors/MqttTransportAdaptor.java | 3 --- .../transport/mqtt/adaptors/ProtoMqttAdaptor.java | 5 ----- .../mqtt/session/GatewayDeviceSessionCtx.java | 9 +-------- .../mqtt/session/GatewaySessionHandler.java | 9 --------- .../common/transport/adaptor/JsonConverter.java | 7 +++---- .../msa/connectivity/MqttGatewayClientTest.java | 1 - 10 files changed, 12 insertions(+), 46 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index 20109ad074..bfcb95a71c 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -450,6 +450,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { .setRequestId(requestId) .setSharedStateMsg(true) .addAllSharedAttributeList(toTsKvProtos(result)) + .setIsMultipleAttributesRequest(request.getSharedAttributeNamesCount() > 1) .build(); sendToTransport(responseMsg, sessionInfo); } @@ -471,6 +472,8 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { .setRequestId(requestId) .addAllClientAttributeList(toTsKvProtos(result.get(0))) .addAllSharedAttributeList(toTsKvProtos(result.get(1))) + .setIsMultipleAttributesRequest( + request.getSharedAttributeNamesCount() > 1 || request.getClientAttributeNamesCount() > 1) .build(); sendToTransport(responseMsg, sessionInfo); } diff --git a/common/cluster-api/src/main/proto/queue.proto b/common/cluster-api/src/main/proto/queue.proto index dca92c1c7a..6565a45740 100644 --- a/common/cluster-api/src/main/proto/queue.proto +++ b/common/cluster-api/src/main/proto/queue.proto @@ -149,6 +149,7 @@ message GetAttributeResponseMsg { int32 requestId = 1; repeated TsKvProto clientAttributeList = 2; repeated TsKvProto sharedAttributeList = 3; + bool isMultipleAttributesRequest = 4; string error = 5; bool sharedStateMsg = 6; } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java index 0cc65da1b0..a608d6c52d 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/BackwardCompatibilityAdaptor.java @@ -128,11 +128,6 @@ public class BackwardCompatibilityAdaptor implements MqttTransportAdaptor { return protoAdaptor.convertToGatewayPublish(ctx, deviceName, rpcRequest); } - @Override - public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg, boolean multipleAttributeKeysRequested) throws AdaptorException { - return protoAdaptor.convertToGatewayPublish(ctx, deviceName, responseMsg, multipleAttributeKeysRequested); - } - @Override public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToServerRpcResponseMsg rpcResponse, String topicBase) throws AdaptorException { log.warn("[{}] invoked not implemented adaptor method! ToServerRpcResponseMsg: {} TopicBase: {}", ctx.getSessionId(), rpcResponse, topicBase); 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 d71b87efdc..8687b02f56 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,7 +34,6 @@ 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; @@ -43,7 +42,6 @@ import java.util.Arrays; import java.util.HashSet; import java.util.Optional; import java.util.Set; -import java.util.Map; import java.util.UUID; import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_SOFTWARE_FIRMWARE_RESPONSES_TOPIC_FORMAT; @@ -125,14 +123,9 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { @Override public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException { - return convertToGatewayPublish(ctx, deviceName, responseMsg, false); + return processConvertFromGatewayAttributeResponseMsg(ctx, deviceName, responseMsg); } - - @Override - public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg, boolean multipleAttributeKeysRequested) throws AdaptorException { - return processConvertFromGatewayAttributeResponseMsg(ctx, deviceName, responseMsg, multipleAttributeKeysRequested); - } - + @Override public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic) { return Optional.of(createMqttPublishMsg(ctx, topic, JsonConverter.toJson(notificationMsg))); @@ -239,11 +232,11 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { } } - private Optional processConvertFromGatewayAttributeResponseMsg(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg, boolean multipleAttrKeysRequested) throws AdaptorException { + private Optional processConvertFromGatewayAttributeResponseMsg(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException { if (!StringUtils.isEmpty(responseMsg.getError())) { throw new AdaptorException(responseMsg.getError()); } else { - JsonObject result = JsonConverter.getJsonObjectForGateway(deviceName, responseMsg, multipleAttrKeysRequested); + JsonObject result = JsonConverter.getJsonObjectForGateway(deviceName, responseMsg); 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/adaptors/MqttTransportAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java index a980e91b53..53ef439eda 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java @@ -25,7 +25,6 @@ import io.netty.handler.codec.mqtt.MqttPublishMessage; import io.netty.handler.codec.mqtt.MqttPublishVariableHeader; import org.thingsboard.server.common.data.ota.OtaPackageType; import org.thingsboard.server.common.transport.adaptor.AdaptorException; -import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg; @@ -73,8 +72,6 @@ public interface MqttTransportAdaptor { Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, ToDeviceRpcRequestMsg rpcRequest) throws AdaptorException; - Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg, boolean multipleAttributeKeysRequested) throws AdaptorException; - Optional convertToPublish(MqttDeviceAwareSessionContext ctx, ToServerRpcResponseMsg rpcResponse, String topicBase) throws AdaptorException; ProvisionDeviceRequestMsg convertToProvisionRequestMsg(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException; diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java index 4dc9a95d6b..4fa47b367d 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java @@ -209,11 +209,6 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { return Optional.of(createMqttPublishMsg(ctx, MqttTopics.GATEWAY_RPC_TOPIC, payloadBytes)); } - @Override - public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg, boolean multipleAttributeKeysRequested) throws AdaptorException { - return convertToGatewayPublish(ctx, deviceName, responseMsg); - } - public static byte[] toBytes(ByteBuf inbound) { byte[] bytes = new byte[inbound.readableBytes()]; int readerIndex = inbound.readerIndex(); 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 02324fd3e7..f41c5668cd 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 @@ -80,10 +80,7 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple @Override public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg response) { try { - boolean multipleAttrKeysReq = isMultipleAttributeKeysRequested(response.getRequestId()); - parent.getPayloadAdaptor() - .convertToGatewayPublish(this, getDeviceInfo().getDeviceName(), response, multipleAttrKeysReq) - .ifPresent(parent::writeAndFlush); + parent.getPayloadAdaptor().convertToGatewayPublish(this, getDeviceInfo().getDeviceName(), response).ifPresent(parent::writeAndFlush); } catch (Exception e) { log.trace("[{}] Failed to convert device attributes response to MQTT msg", sessionId, e); } @@ -139,10 +136,6 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple // This feature is not supported in the TB IoT Gateway yet. } - public boolean isMultipleAttributeKeysRequested(int requestId) { - return parent.isMultipleAttributeKeysRequested(requestId); - } - 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 12e9f786bf..7b45b5f9d9 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,7 +89,6 @@ 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(); @@ -108,13 +107,6 @@ public class GatewaySessionHandler { return new ConcurrentReferenceHashMap<>(16, ReferenceType.WEAK); } - public boolean isMultipleAttributeKeysRequested(int requestId) { - JsonObject request = pendingAttributesRequests.getOrDefault(requestId, new JsonObject()); - boolean multipleAttributeKeysRequested = request.has("keys") && request.get("keys").getAsJsonArray().size() > 1; - pendingAttributesRequests.remove(requestId); - return multipleAttributeKeysRequested; - } - public void onDeviceConnect(MqttPublishMessage mqttMsg) throws AdaptorException { if (isJsonPayloadType()) { onDeviceConnectJson(mqttMsg); @@ -566,7 +558,6 @@ 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 86f90d786f..1e6267ca33 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 @@ -315,17 +315,16 @@ public class JsonConverter { public static JsonObject getJsonObjectForGateway( String deviceName, - TransportProtos.GetAttributeResponseMsg responseMsg, - boolean multipleAttrKeysRequested + TransportProtos.GetAttributeResponseMsg responseMsg ) { JsonObject result = new JsonObject(); result.addProperty("id", responseMsg.getRequestId()); result.addProperty(DEVICE_PROPERTY, deviceName); if (responseMsg.getClientAttributeListCount() > 0) { - addValues(result, responseMsg.getClientAttributeListList(), multipleAttrKeysRequested); + addValues(result, responseMsg.getClientAttributeListList(), responseMsg.getIsMultipleAttributesRequest()); } if (responseMsg.getSharedAttributeListCount() > 0) { - addValues(result, responseMsg.getSharedAttributeListList(), multipleAttrKeysRequested); + addValues(result, responseMsg.getSharedAttributeListList(), responseMsg.getIsMultipleAttributesRequest()); } return result; } diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java index dfdc68d9d5..74dfc2107c 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttGatewayClientTest.java @@ -156,7 +156,6 @@ public class MqttGatewayClientTest extends AbstractContainerTest { public void responseDataOnAttributesRequestCheck() throws Exception { Optional createdDeviceCredentials = restClient.getDeviceCredentialsByDeviceId(createdDevice.getId()); Assert.assertTrue(createdDeviceCredentials.isPresent()); - WsClient wsClient = subscribeToWebSocket(createdDevice.getId(), "CLIENT_SCOPE", CmdsType.ATTR_SUB_CMDS); JsonObject sharedAttributes = new JsonObject(); sharedAttributes.addProperty("attr1", "value1"); sharedAttributes.addProperty("attr2", true); From fc9078a2be2b78ce8ab6bd6b9ba9c2025f147ac3 Mon Sep 17 00:00:00 2001 From: desoliture Date: Thu, 20 Jan 2022 12:58:55 +0200 Subject: [PATCH 6/6] fix case where both shared and client attributes were requested --- .../server/actors/device/DeviceActorMessageProcessor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index bfcb95a71c..baf2e9b03e 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -473,7 +473,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { .addAllClientAttributeList(toTsKvProtos(result.get(0))) .addAllSharedAttributeList(toTsKvProtos(result.get(1))) .setIsMultipleAttributesRequest( - request.getSharedAttributeNamesCount() > 1 || request.getClientAttributeNamesCount() > 1) + request.getSharedAttributeNamesCount() + request.getClientAttributeNamesCount() > 1) .build(); sendToTransport(responseMsg, sessionInfo); }