Browse Source

MQTT Gateway API attributes request fix

pull/5796/head
desoliture 5 years ago
parent
commit
b53746bda6
  1. 12
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java
  2. 6
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java
  3. 6
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java
  4. 15
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/JsonConverter.java

12
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<MqttMessage> 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<MqttMessage> processConvertFromGatewayAttributeResponseMsg(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException {
private Optional<MqttMessage> 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));
}
}

6
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<Integer, JsonObject> getPendingAttributesRequests() {
return parent.getPendingAttributesRequests();
}
private boolean isAckExpected(MqttMessage message) {
return message.fixedHeader().qosLevel().value() > 0;
}

6
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java

@ -89,6 +89,7 @@ public class GatewaySessionHandler {
private final ConcurrentMap<MqttTopicMatcher, Integer> mqttQoSMap;
private final ChannelHandlerContext channel;
private final DeviceSessionCtx deviceSessionCtx;
private final Map<Integer, JsonObject> 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<Integer, JsonObject> 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<String> keys;

15
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<TransportProtos.TsKvProto> kvList) {
if (kvList.size() == 1) {
private static void addValues(JsonObject result, List<TransportProtos.TsKvProto> kvList, boolean multipleAttrKeysRequested) {
if (kvList.size() == 1 && !multipleAttrKeysRequested) {
addValueToJson(result, "value", kvList.get(0).getKv());
} else {
JsonObject values;

Loading…
Cancel
Save