diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index a1f9800b23..ed73b54ea1 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java @@ -47,6 +47,7 @@ import org.thingsboard.server.common.data.device.data.CoapDeviceTransportConfigu import org.thingsboard.server.common.data.device.data.Lwm2mDeviceTransportConfiguration; import org.thingsboard.server.common.data.device.data.PowerMode; import org.thingsboard.server.common.data.device.data.PowerSavingConfiguration; +import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.ProvisionDeviceProfileCredentials; import org.thingsboard.server.common.data.id.CustomerId; @@ -283,7 +284,12 @@ public class DefaultTransportApiService implements TransportApiService { deviceCreationLock.lock(); try { DeviceProfile deviceProfile = deviceProfileCache.findOrCreateDeviceProfile(gateway.getTenantId(), requestMsg.getDeviceType()); - boolean isSparkplug = ((MqttDeviceProfileTransportConfiguration) deviceProfile.getProfileData().getTransportConfiguration()).isSparkPlug(); + DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); + boolean isSparkplug = false; + if (transportConfiguration instanceof MqttDeviceProfileTransportConfiguration && + ((MqttDeviceProfileTransportConfiguration) transportConfiguration).isSparkPlug()) { + isSparkplug = true; + } Device device = deviceService.findDeviceByTenantIdAndName(gateway.getTenantId(), requestMsg.getDeviceName()); if (device == null) { TenantId tenantId = gateway.getTenantId(); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index d99e2cb801..bf1a996170 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -47,7 +47,6 @@ import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.TransportPayloadType; -import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.MqttTopics; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.OtaPackageId; @@ -979,28 +978,19 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } } - private void checkSparkPlugSession(SessionMetaData sessionMetaData, MqttConnectMessage connectMessage) { - if (((MqttDeviceProfileTransportConfiguration) deviceSessionCtx - .getDeviceProfile() - .getProfileData() - .getTransportConfiguration()) - .isSparkPlug()) { - TransportDeviceInfo device = deviceSessionCtx.getDeviceInfo(); - try { - JsonNode infoNode = context.getMapper().readTree(device.getAdditionalInfo()); - if (infoNode != null) { - SparkplugTopic sparkplugTopic = parseTopic(connectMessage.payload().willTopic()); - SparkplugBProto.Payload payloadBProto = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes()); - if (sparkPlugSessionHandler == null) { - log.error("SparkPlugConnected [{}] [{}]", sparkplugTopic.isNode() ? "node" : "device: " + sparkplugTopic.getDeviceId(), sparkplugTopic.getType()); - sparkPlugSessionHandler = new SparkplugNodeSessionHandler(deviceSessionCtx, sessionId, sparkplugTopic.toString()); - } else { - log.error("SparkPlugReConnected [{}] [{}]", sparkplugTopic.isNode() ? "node" : "device: " + sparkplugTopic.getDeviceId(), sparkplugTopic.getType()); - } - } - } catch (Exception e) { - log.trace("[{}][{}] Failed to fetch sparkplugDevice additional info or sparkplugTopicName", sessionId, device.getDeviceName(), e); + private void checkSparkPlugSession(MqttConnectMessage connectMessage) { + try { + SparkplugTopic sparkplugTopic = parseTopic(connectMessage.payload().willTopic()); + // Test proto + SparkplugBProto.Payload payloadBProto = SparkplugBProto.Payload.parseFrom(connectMessage.payload().willMessageInBytes()); + // + if (sparkPlugSessionHandler == null) { + sparkPlugSessionHandler = new SparkplugNodeSessionHandler(deviceSessionCtx, sessionId, sparkplugTopic.toString()); + } else { + log.warn("SparkPlugNodeReConnected [{}] [{}]", sparkplugTopic.getDeviceId(), sparkplugTopic.getType()); } + } catch (Exception e) { + log.trace("[{}][{}] Failed to fetch sparkplugDevice additional info or sparkplugTopicName", sessionId, deviceSessionCtx.getDeviceInfo().getDeviceName(), e); } } @@ -1038,8 +1028,11 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement @Override public void onSuccess(Void msg) { SessionMetaData sessionMetaData = transportService.registerAsyncSession(deviceSessionCtx.getSessionInfo(), MqttTransportHandler.this); - checkGatewaySession(sessionMetaData); - checkSparkPlugSession(sessionMetaData, connectMessage); + if (deviceSessionCtx.isSparkplug()) { + checkSparkPlugSession(connectMessage); + } else { + checkGatewaySession(sessionMetaData); + } ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED, connectMessage)); deviceSessionCtx.setConnected(true); log.debug("[{}] Client connected!", sessionId); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java index 469ae1034d..e69248dd30 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java @@ -20,6 +20,8 @@ import lombok.Getter; import lombok.Setter; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; +import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.transport.auth.TransportDeviceInfo; import org.thingsboard.server.gen.transport.TransportProtos; @@ -81,4 +83,15 @@ public abstract class DeviceAwareSessionContext implements SessionContext { public void setDisconnected() { this.connected = false; } + + public boolean isSparkplug() { + DeviceProfileTransportConfiguration transportConfiguration = this.deviceProfile.getProfileData().getTransportConfiguration(); + if (transportConfiguration instanceof MqttDeviceProfileTransportConfiguration) { + if (((MqttDeviceProfileTransportConfiguration) transportConfiguration).isSparkPlug()) { + return true; + } + } + return false; + } + }