Browse Source

Merge pull request #14240 from thingsboard/rc

rc
pull/14252/head
Viacheslav Klimov 11 months ago
committed by GitHub
parent
commit
f8c73550e0
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 11
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java
  3. 1
      application/src/test/java/org/thingsboard/server/transport/coap/security/AbstractCoapSecurityIntegrationTest.java
  4. 4
      common/data/src/main/java/org/thingsboard/server/common/data/security/DeviceCredentials.java
  5. 6
      common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceAttributesEventNotificationMsg.java
  6. 3
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java
  7. 3
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  8. 3
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportServerInitializer.java
  9. 3
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java
  10. 6
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java
  11. 4
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java
  12. 3
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java
  13. 3
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewayDeviceSessionContext.java
  14. 11
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java
  15. 49
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java
  16. 3
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionContext.java
  17. 3
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewaySessionHandler.java
  18. 5
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttDeviceAwareSessionContext.java
  19. 12
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttTopicMatcher.java
  20. 3
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java
  21. 4
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsEvictEvent.java
  22. 27
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java
  23. 42
      dao/src/test/java/org/thingsboard/server/dao/service/DeviceCredentialsServiceTest.java
  24. 3
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java

11
application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java

@ -115,14 +115,9 @@ import java.util.concurrent.TimeUnit;
import java.util.function.Consumer; import java.util.function.Consumer;
import java.util.stream.Collectors; import java.util.stream.Collectors;
/**
* @author Andrew Shvayka
*/
@Slf4j @Slf4j
public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
static final String SESSION_TIMEOUT_MESSAGE = "session timeout!";
final TenantId tenantId; final TenantId tenantId;
final DeviceId deviceId; final DeviceId deviceId;
final LinkedHashMapRemoveEldest<UUID, SessionInfoMetaData> sessions; final LinkedHashMapRemoveEldest<UUID, SessionInfoMetaData> sessions;
@ -178,7 +173,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
private EdgeId findRelatedEdgeId() { private EdgeId findRelatedEdgeId() {
List<EntityRelation> result = List<EntityRelation> result =
systemContext.getRelationService().findByToAndType(tenantId, deviceId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); systemContext.getRelationService().findByToAndType(tenantId, deviceId, EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE);
if (result != null && result.size() > 0) { if (result != null && !result.isEmpty()) {
EntityRelation relationToEdge = result.get(0); EntityRelation relationToEdge = result.get(0);
if (relationToEdge.getFrom() != null && relationToEdge.getFrom().getId() != null) { if (relationToEdge.getFrom() != null && relationToEdge.getFrom().getId() != null) {
log.trace("[{}][{}] found edge [{}] for device", tenantId, deviceId, relationToEdge.getFrom().getId()); log.trace("[{}][{}] found edge [{}] for device", tenantId, deviceId, relationToEdge.getFrom().getId());
@ -501,7 +496,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
UUID sessionId = getSessionId(sessionInfo); UUID sessionId = getSessionId(sessionInfo);
DeviceId deviceId = new DeviceId(new UUID(msg.getDeviceIdMSB(), msg.getDeviceIdLSB())); DeviceId deviceId = new DeviceId(new UUID(msg.getDeviceIdMSB(), msg.getDeviceIdLSB()));
ListenableFuture<Void> registrationFuture = systemContext.getClaimDevicesService() ListenableFuture<Void> registrationFuture = systemContext.getClaimDevicesService()
.registerClaimingInfo(tenantId, deviceId, msg.getSecretKey(), msg.getDurationMs()); .registerClaimingInfo(tenantId, deviceId, msg.getSecretKey(), msg.getDurationMs());
Futures.addCallback(registrationFuture, new FutureCallback<>() { Futures.addCallback(registrationFuture, new FutureCallback<>() {
@Override @Override
public void onSuccess(Void result) { public void onSuccess(Void result) {
@ -723,7 +718,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso
toDeviceRpcPendingMap.remove(requestId); toDeviceRpcPendingMap.remove(requestId);
status = RpcStatus.FAILED; status = RpcStatus.FAILED;
response = JacksonUtil.newObjectNode().put("error", "There was a Timeout and all retry " + response = JacksonUtil.newObjectNode().put("error", "There was a Timeout and all retry " +
"attempts have been exhausted. Retry attempts set: " + maxRpcRetries); "attempts have been exhausted. Retry attempts set: " + maxRpcRetries);
} }
} else { } else {
md.setRetries(md.getRetries() + 1); md.setRetries(md.getRetries() + 1);

2
application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java

@ -193,7 +193,7 @@ public abstract class AbstractBulkImportService<E extends HasId<? extends Entity
}); });
return Map.entry(kvType, kvs); return Map.entry(kvType, kvs);
}) })
.filter(kvsEntry -> kvsEntry.getValue().entrySet().size() > 0) .filter(kvsEntry -> !kvsEntry.getValue().entrySet().isEmpty())
.forEach(kvsEntry -> { .forEach(kvsEntry -> {
BulkImportColumnType kvType = kvsEntry.getKey(); BulkImportColumnType kvType = kvsEntry.getKey();
if (kvType == BulkImportColumnType.SHARED_ATTRIBUTE || kvType == BulkImportColumnType.SERVER_ATTRIBUTE) { if (kvType == BulkImportColumnType.SHARED_ATTRIBUTE || kvType == BulkImportColumnType.SERVER_ATTRIBUTE) {

1
application/src/test/java/org/thingsboard/server/transport/coap/security/AbstractCoapSecurityIntegrationTest.java

@ -64,6 +64,7 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.
"coap.server.enabled=true", "coap.server.enabled=true",
"coap.dtls.enabled=true", "coap.dtls.enabled=true",
"coap.dtls.credentials.pem.cert_file=coap/credentials/server/cert.pem", "coap.dtls.credentials.pem.cert_file=coap/credentials/server/cert.pem",
"coap.dtls.x509.skip_validity_check_for_client_cert=true",
"device.connectivity.coaps.enabled=true", "device.connectivity.coaps.enabled=true",
"service.integrations.supported=ALL", "service.integrations.supported=ALL",
"transport.coap.enabled=true", "transport.coap.enabled=true",

4
common/data/src/main/java/org/thingsboard/server/common/data/security/DeviceCredentials.java

@ -24,11 +24,15 @@ import org.thingsboard.server.common.data.HasVersion;
import org.thingsboard.server.common.data.id.DeviceCredentialsId; import org.thingsboard.server.common.data.id.DeviceCredentialsId;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import java.io.Serial;
@Schema @Schema
@EqualsAndHashCode(callSuper = true) @EqualsAndHashCode(callSuper = true)
public class DeviceCredentials extends BaseData<DeviceCredentialsId> implements DeviceCredentialsFilter, HasVersion { public class DeviceCredentials extends BaseData<DeviceCredentialsId> implements DeviceCredentialsFilter, HasVersion {
@Serial
private static final long serialVersionUID = -7869261127032877765L; private static final long serialVersionUID = -7869261127032877765L;
private DeviceId deviceId; private DeviceId deviceId;
private DeviceCredentialsType credentialsType; private DeviceCredentialsType credentialsType;
private String credentialsId; private String credentialsId;

6
common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceAttributesEventNotificationMsg.java

@ -23,16 +23,15 @@ import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg; import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg;
import java.io.Serial;
import java.util.HashSet; import java.util.HashSet;
import java.util.List; import java.util.List;
import java.util.Set; import java.util.Set;
/**
* @author Andrew Shvayka
*/
@Data @Data
public class DeviceAttributesEventNotificationMsg implements ToDeviceActorNotificationMsg { public class DeviceAttributesEventNotificationMsg implements ToDeviceActorNotificationMsg {
@Serial
private static final long serialVersionUID = 2422071590415277039L; private static final long serialVersionUID = 2422071590415277039L;
private final TenantId tenantId; private final TenantId tenantId;
@ -56,4 +55,5 @@ public class DeviceAttributesEventNotificationMsg implements ToDeviceActorNotifi
public MsgType getMsgType() { public MsgType getMsgType() {
return MsgType.DEVICE_ATTRIBUTES_UPDATE_TO_DEVICE_ACTOR_MSG; return MsgType.DEVICE_ATTRIBUTES_UPDATE_TO_DEVICE_ACTOR_MSG;
} }
} }

3
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java

@ -46,9 +46,6 @@ import java.security.cert.X509Certificate;
import java.util.concurrent.CountDownLatch; import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
/**
* Created by valerii.sosliuk on 11/6/16.
*/
@Slf4j @Slf4j
@Component("MqttSslHandlerProvider") @Component("MqttSslHandlerProvider")
@ConditionalOnProperty(prefix = "transport.mqtt.ssl", value = "enabled", havingValue = "true", matchIfMissing = false) @ConditionalOnProperty(prefix = "transport.mqtt.ssl", value = "enabled", havingValue = "true", matchIfMissing = false)

3
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -133,9 +133,6 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetr
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic.parseTopic; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopic.parseTopic;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.parseTopicPublish; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.parseTopicPublish;
/**
* @author Andrew Shvayka
*/
@Slf4j @Slf4j
public class MqttTransportHandler extends ChannelInboundHandlerAdapter implements GenericFutureListener<Future<? super Void>>, SessionMsgListener { public class MqttTransportHandler extends ChannelInboundHandlerAdapter implements GenericFutureListener<Future<? super Void>>, SessionMsgListener {

3
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportServerInitializer.java

@ -25,9 +25,6 @@ import io.netty.handler.ssl.SslHandler;
import org.thingsboard.server.transport.mqtt.limits.IpFilter; import org.thingsboard.server.transport.mqtt.limits.IpFilter;
import org.thingsboard.server.transport.mqtt.limits.ProxyIpFilter; import org.thingsboard.server.transport.mqtt.limits.ProxyIpFilter;
/**
* @author Andrew Shvayka
*/
public class MqttTransportServerInitializer extends ChannelInitializer<SocketChannel> { public class MqttTransportServerInitializer extends ChannelInitializer<SocketChannel> {
private final MqttTransportContext context; private final MqttTransportContext context;

3
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java

@ -35,9 +35,6 @@ import org.thingsboard.server.common.data.TbTransportService;
import java.net.InetSocketAddress; import java.net.InetSocketAddress;
/**
* @author Andrew Shvayka
*/
@Service("MqttTransportService") @Service("MqttTransportService")
@TbMqttTransportComponent @TbMqttTransportComponent
@Slf4j @Slf4j

6
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java

@ -45,10 +45,6 @@ import java.util.UUID;
import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_SOFTWARE_FIRMWARE_RESPONSES_TOPIC_FORMAT; import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_SOFTWARE_FIRMWARE_RESPONSES_TOPIC_FORMAT;
/**
* @author Andrew Shvayka
*/
@Component @Component
@Slf4j @Slf4j
public class JsonMqttAdaptor implements MqttTransportAdaptor { public class JsonMqttAdaptor implements MqttTransportAdaptor {
@ -122,7 +118,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor {
public Optional<MqttMessage> convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException { public Optional<MqttMessage> convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException {
return processConvertFromGatewayAttributeResponseMsg(ctx, deviceName, responseMsg); return processConvertFromGatewayAttributeResponseMsg(ctx, deviceName, responseMsg);
} }
@Override @Override
public Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic) { public Optional<MqttMessage> convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic) {
return Optional.of(createMqttPublishMsg(ctx, topic, JsonConverter.toJson(notificationMsg))); return Optional.of(createMqttPublishMsg(ctx, topic, JsonConverter.toJson(notificationMsg)));

4
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/MqttTransportAdaptor.java

@ -41,9 +41,6 @@ import org.thingsboard.server.transport.mqtt.session.MqttDeviceAwareSessionConte
import java.util.Optional; import java.util.Optional;
/**
* @author Andrew Shvayka
*/
public interface MqttTransportAdaptor { public interface MqttTransportAdaptor {
ByteBufAllocator ALLOCATOR = new UnpooledByteBufAllocator(false); ByteBufAllocator ALLOCATOR = new UnpooledByteBufAllocator(false);
@ -90,4 +87,5 @@ public interface MqttTransportAdaptor {
payload.writeBytes(payloadInBytes); payload.writeBytes(payloadInBytes);
return new MqttPublishMessage(mqttFixedHeader, header, payload); return new MqttPublishMessage(mqttFixedHeader, header, payload);
} }
} }

3
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java

@ -48,7 +48,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
public TransportProtos.PostTelemetryMsg convertToPostTelemetry(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException { public TransportProtos.PostTelemetryMsg convertToPostTelemetry(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException {
DeviceSessionCtx deviceSessionCtx = (DeviceSessionCtx) ctx; DeviceSessionCtx deviceSessionCtx = (DeviceSessionCtx) ctx;
byte[] bytes = toBytes(inbound.payload()); byte[] bytes = toBytes(inbound.payload());
Descriptors.Descriptor telemetryDynamicMsgDescriptor = ProtoConverter.validateDescriptor(deviceSessionCtx.getTelemetryDynamicMsgDescriptor()); Descriptors.Descriptor telemetryDynamicMsgDescriptor = ProtoConverter.validateDescriptor(deviceSessionCtx.getTelemetryDynamicMessageDescriptor());
try { try {
return JsonConverter.convertToTelemetryProto(JsonParser.parseString(ProtoConverter.dynamicMsgToJson(bytes, telemetryDynamicMsgDescriptor))); return JsonConverter.convertToTelemetryProto(JsonParser.parseString(ProtoConverter.dynamicMsgToJson(bytes, telemetryDynamicMsgDescriptor)));
} catch (Exception e) { } catch (Exception e) {
@ -228,4 +228,5 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
private int getRequestId(String topicName, String topic) { private int getRequestId(String topicName, String topic) {
return Integer.parseInt(topicName.substring(topic.length())); return Integer.parseInt(topicName.substring(topic.length()));
} }
} }

3
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewayDeviceSessionContext.java

@ -33,9 +33,6 @@ import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
/**
* Created by ashvayka on 19.01.17.
*/
@ToString(callSuper = true) @ToString(callSuper = true)
@Slf4j @Slf4j
public abstract class AbstractGatewayDeviceSessionContext<T extends AbstractGatewaySessionHandler> extends MqttDeviceAwareSessionContext implements SessionMsgListener { public abstract class AbstractGatewayDeviceSessionContext<T extends AbstractGatewaySessionHandler> extends MqttDeviceAwareSessionContext implements SessionMsgListener {

11
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java

@ -96,9 +96,6 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugConn
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.STATE; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.STATE;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.messageName; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMessageType.messageName;
/**
* Created by ashvayka on 19.01.17.
*/
@Slf4j @Slf4j
public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDeviceSessionContext> { public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDeviceSessionContext> {
@ -116,6 +113,7 @@ public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDev
@Getter @Getter
protected final UUID sessionId; protected final UUID sessionId;
private final ConcurrentMap<String, Lock> deviceCreationLockMap; private final ConcurrentMap<String, Lock> deviceCreationLockMap;
@Getter
private final ConcurrentMap<String, T> devices; private final ConcurrentMap<String, T> devices;
private final ConcurrentMap<String, ListenableFuture<T>> deviceFutures; private final ConcurrentMap<String, ListenableFuture<T>> deviceFutures;
protected final ConcurrentMap<MqttTopicMatcher, Integer> mqttQoSMap; protected final ConcurrentMap<MqttTopicMatcher, Integer> mqttQoSMap;
@ -821,11 +819,7 @@ public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDev
transportService.process(sessionInfo, postTelemetryMsg, pubAckCallback); transportService.process(sessionInfo, postTelemetryMsg, pubAckCallback);
} }
public ConcurrentMap<String, T> getDevices() { protected <T> TransportServiceCallback<Void> getAggregatePubAckCallback(
return this.devices;
}
protected <T>TransportServiceCallback<Void> getAggregatePubAckCallback(
final ChannelHandlerContext ctx, final ChannelHandlerContext ctx,
final int msgId, final int msgId,
final String deviceName, final String deviceName,
@ -915,4 +909,5 @@ public abstract class AbstractGatewaySessionHandler<T extends AbstractGatewayDev
log.trace("Failed to send device disconnect to gateway session", e); log.trace("Failed to send device disconnect to gateway session", e);
} }
} }
} }

49
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java

@ -49,9 +49,6 @@ import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock; import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Consumer; import java.util.function.Consumer;
/**
* @author Andrew Shvayka
*/
@Slf4j @Slf4j
public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
@ -84,13 +81,18 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
private volatile MqttTopicFilter attributesSubscribeTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter(); private volatile MqttTopicFilter attributesSubscribeTopicFilter = MqttTopicFilterFactory.getDefaultAttributesFilter();
@Getter @Getter
private volatile TransportPayloadType payloadType = TransportPayloadType.JSON; private volatile TransportPayloadType payloadType = TransportPayloadType.JSON;
@Getter
private volatile Descriptors.Descriptor attributesDynamicMessageDescriptor; private volatile Descriptors.Descriptor attributesDynamicMessageDescriptor;
@Getter
private volatile Descriptors.Descriptor telemetryDynamicMessageDescriptor; private volatile Descriptors.Descriptor telemetryDynamicMessageDescriptor;
@Getter
private volatile Descriptors.Descriptor rpcResponseDynamicMessageDescriptor; private volatile Descriptors.Descriptor rpcResponseDynamicMessageDescriptor;
@Getter
private volatile DynamicMessage.Builder rpcRequestDynamicMessageBuilder; private volatile DynamicMessage.Builder rpcRequestDynamicMessageBuilder;
private volatile MqttTransportAdaptor adaptor; private volatile MqttTransportAdaptor adaptor;
private volatile boolean jsonPayloadFormatCompatibilityEnabled; private volatile boolean jsonPayloadFormatCompatibilityEnabled;
private volatile boolean useJsonPayloadFormatForDefaultDownlinkTopics; private volatile boolean useJsonPayloadFormatForDefaultDownlinkTopics;
@Getter
private volatile boolean sendAckOnValidationException; private volatile boolean sendAckOnValidationException;
@Getter @Getter
@ -131,26 +133,6 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
return payloadType.equals(TransportPayloadType.JSON); return payloadType.equals(TransportPayloadType.JSON);
} }
public boolean isSendAckOnValidationException() {
return sendAckOnValidationException;
}
public Descriptors.Descriptor getTelemetryDynamicMsgDescriptor() {
return telemetryDynamicMessageDescriptor;
}
public Descriptors.Descriptor getAttributesDynamicMessageDescriptor() {
return attributesDynamicMessageDescriptor;
}
public Descriptors.Descriptor getRpcResponseDynamicMessageDescriptor() {
return rpcResponseDynamicMessageDescriptor;
}
public DynamicMessage.Builder getRpcRequestDynamicMessageBuilder() {
return rpcRequestDynamicMessageBuilder;
}
@Override @Override
public void setDeviceProfile(DeviceProfile deviceProfile) { public void setDeviceProfile(DeviceProfile deviceProfile) {
super.setDeviceProfile(deviceProfile); super.setDeviceProfile(deviceProfile);
@ -166,8 +148,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
private void updateDeviceSessionConfiguration(DeviceProfile deviceProfile) { private void updateDeviceSessionConfiguration(DeviceProfile deviceProfile) {
DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration(); DeviceProfileTransportConfiguration transportConfiguration = deviceProfile.getProfileData().getTransportConfiguration();
if (transportConfiguration.getType().equals(DeviceTransportType.MQTT) && if (transportConfiguration.getType().equals(DeviceTransportType.MQTT) &&
transportConfiguration instanceof MqttDeviceProfileTransportConfiguration) { transportConfiguration instanceof MqttDeviceProfileTransportConfiguration mqttConfig) {
MqttDeviceProfileTransportConfiguration mqttConfig = (MqttDeviceProfileTransportConfiguration) transportConfiguration;
TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttConfig.getTransportPayloadTypeConfiguration(); TransportPayloadTypeConfiguration transportPayloadTypeConfiguration = mqttConfig.getTransportPayloadTypeConfiguration();
payloadType = transportPayloadTypeConfiguration.getTransportPayloadType(); payloadType = transportPayloadTypeConfiguration.getTransportPayloadType();
deviceProfileMqttTransportType = true; deviceProfileMqttTransportType = true;
@ -199,16 +180,12 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
} }
public MqttTransportAdaptor getAdaptor(TopicType topicType) { public MqttTransportAdaptor getAdaptor(TopicType topicType) {
switch (topicType) { return switch (topicType) {
case V2: case V2 -> getDefaultAdaptor();
return getDefaultAdaptor(); case V2_JSON -> context.getJsonMqttAdaptor();
case V2_JSON: case V2_PROTO -> context.getProtoMqttAdaptor();
return context.getJsonMqttAdaptor(); default -> useJsonPayloadFormatForDefaultDownlinkTopics ? context.getJsonMqttAdaptor() : getDefaultAdaptor();
case V2_PROTO: };
return context.getProtoMqttAdaptor();
default:
return useJsonPayloadFormatForDefaultDownlinkTopics ? context.getJsonMqttAdaptor() : getDefaultAdaptor();
}
} }
private MqttTransportAdaptor getDefaultAdaptor() { private MqttTransportAdaptor getDefaultAdaptor() {
@ -269,7 +246,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
} }
} }
public Collection<MqttMessage> getMsgQueueSnapshot(){ public Collection<MqttMessage> getMsgQueueSnapshot() {
return Collections.unmodifiableCollection(msgQueue); return Collections.unmodifiableCollection(msgQueue);
} }

3
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionContext.java

@ -22,9 +22,6 @@ import org.thingsboard.server.common.transport.auth.TransportDeviceInfo;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
/**
* Created by nickAS21 on 26.12.22
*/
@ToString(callSuper = true) @ToString(callSuper = true)
public class GatewayDeviceSessionContext extends AbstractGatewayDeviceSessionContext<GatewaySessionHandler> { public class GatewayDeviceSessionContext extends AbstractGatewayDeviceSessionContext<GatewaySessionHandler> {

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

@ -28,9 +28,6 @@ import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.Optional; import java.util.Optional;
import java.util.UUID; import java.util.UUID;
/**
* Created by nickAS21 on 26.12.22
*/
@Slf4j @Slf4j
public class GatewaySessionHandler extends AbstractGatewaySessionHandler<GatewayDeviceSessionContext> { public class GatewaySessionHandler extends AbstractGatewaySessionHandler<GatewayDeviceSessionContext> {

5
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttDeviceAwareSessionContext.java

@ -25,9 +25,6 @@ import java.util.UUID;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
import java.util.stream.Collectors; import java.util.stream.Collectors;
/**
* Created by ashvayka on 30.08.18.
*/
@ToString(callSuper = true) @ToString(callSuper = true)
public abstract class MqttDeviceAwareSessionContext extends DeviceAwareSessionContext { public abstract class MqttDeviceAwareSessionContext extends DeviceAwareSessionContext {
@ -47,7 +44,7 @@ public abstract class MqttDeviceAwareSessionContext extends DeviceAwareSessionCo
.stream() .stream()
.filter(entry -> entry.getKey().matches(topic)) .filter(entry -> entry.getKey().matches(topic))
.map(Map.Entry::getValue) .map(Map.Entry::getValue)
.collect(Collectors.toList()); .toList();
if (!qosList.isEmpty()) { if (!qosList.isEmpty()) {
return MqttQoS.valueOf(qosList.get(0)); return MqttQoS.valueOf(qosList.get(0));
} else { } else {

12
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttTopicMatcher.java

@ -15,26 +15,25 @@
*/ */
package org.thingsboard.server.transport.mqtt.session; package org.thingsboard.server.transport.mqtt.session;
import lombok.Getter;
import java.util.regex.Pattern; import java.util.regex.Pattern;
public class MqttTopicMatcher { public class MqttTopicMatcher {
@Getter
private final String topic; private final String topic;
private final Pattern topicRegex; private final Pattern topicRegex;
public MqttTopicMatcher(String topic) { public MqttTopicMatcher(String topic) {
if(topic == null){ if (topic == null) {
throw new NullPointerException("topic"); throw new NullPointerException("topic");
} }
this.topic = topic; this.topic = topic;
this.topicRegex = Pattern.compile(topic.replace("+", "[^/]+").replace("#", ".+") + "$"); this.topicRegex = Pattern.compile(topic.replace("+", "[^/]+").replace("#", ".+") + "$");
} }
public String getTopic() { public boolean matches(String topic) {
return topic;
}
public boolean matches(String topic){
return this.topicRegex.matcher(topic).matches(); return this.topicRegex.matcher(topic).matches();
} }
@ -52,4 +51,5 @@ public class MqttTopicMatcher {
public int hashCode() { public int hashCode() {
return topic.hashCode(); return topic.hashCode();
} }
} }

3
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java

@ -64,9 +64,6 @@ import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugMetr
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_REGEXP; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_SPLIT_REGEXP;
import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_STATE_REGEXP; import static org.thingsboard.server.transport.mqtt.util.sparkplug.SparkplugTopicService.TOPIC_STATE_REGEXP;
/**
* Created by nickAS21 on 12.12.22
*/
@Slf4j @Slf4j
@SpecVersion(spec = "sparkplug", version = "3.0.0") @SpecVersion(spec = "sparkplug", version = "3.0.0")
public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler<SparkplugDeviceSessionContext> { public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler<SparkplugDeviceSessionContext> {

4
dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsEvictEvent.java

@ -18,9 +18,9 @@ package org.thingsboard.server.dao.device;
import lombok.Data; import lombok.Data;
@Data @Data
class DeviceCredentialsEvictEvent { public class DeviceCredentialsEvictEvent {
private final String newCedentialsId; private final String newCredentialsId;
private final String oldCredentialsId; private final String oldCredentialsId;
} }

27
dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java

@ -47,6 +47,8 @@ import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.dao.exception.DeviceCredentialsValidationException; import org.thingsboard.server.dao.exception.DeviceCredentialsValidationException;
import org.thingsboard.server.dao.service.validator.DeviceCredentialsDataValidator; import org.thingsboard.server.dao.service.validator.DeviceCredentialsDataValidator;
import java.util.Objects;
import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validateId;
import static org.thingsboard.server.dao.service.Validator.validateString; import static org.thingsboard.server.dao.service.Validator.validateString;
@ -61,8 +63,8 @@ public class DeviceCredentialsServiceImpl extends AbstractCachedEntityService<St
@TransactionalEventListener(classes = DeviceCredentialsEvictEvent.class) @TransactionalEventListener(classes = DeviceCredentialsEvictEvent.class)
@Override @Override
public void handleEvictEvent(DeviceCredentialsEvictEvent event) { public void handleEvictEvent(DeviceCredentialsEvictEvent event) {
cache.evict(event.getNewCedentialsId()); cache.evict(event.getNewCredentialsId());
if (StringUtils.isNotEmpty(event.getOldCredentialsId()) && !event.getNewCedentialsId().equals(event.getOldCredentialsId())) { if (StringUtils.isNotEmpty(event.getOldCredentialsId()) && !event.getNewCredentialsId().equals(event.getOldCredentialsId())) {
cache.evict(event.getOldCredentialsId()); cache.evict(event.getOldCredentialsId());
} }
} }
@ -107,7 +109,7 @@ public class DeviceCredentialsServiceImpl extends AbstractCachedEntityService<St
try { try {
var value = deviceCredentialsDao.saveAndFlush(tenantId, deviceCredentials); var value = deviceCredentialsDao.saveAndFlush(tenantId, deviceCredentials);
publishEvictEvent(new DeviceCredentialsEvictEvent(value.getCredentialsId(), oldDeviceCredentials != null ? oldDeviceCredentials.getCredentialsId() : null)); publishEvictEvent(new DeviceCredentialsEvictEvent(value.getCredentialsId(), oldDeviceCredentials != null ? oldDeviceCredentials.getCredentialsId() : null));
if (oldDeviceCredentials != null) { if (oldDeviceCredentials != null && isCredentialsChanged(oldDeviceCredentials, value)) {
eventPublisher.publishEvent(ActionEntityEvent.builder().tenantId(tenantId).entity(value).entityId(value.getDeviceId()).actionType(ActionType.CREDENTIALS_UPDATED).build()); eventPublisher.publishEvent(ActionEntityEvent.builder().tenantId(tenantId).entity(value).entityId(value.getDeviceId()).actionType(ActionType.CREDENTIALS_UPDATED).build());
} }
return value; return value;
@ -140,13 +142,11 @@ public class DeviceCredentialsServiceImpl extends AbstractCachedEntityService<St
@Override @Override
public JsonNode toCredentialsInfo(DeviceCredentials deviceCredentials) { public JsonNode toCredentialsInfo(DeviceCredentials deviceCredentials) {
switch (deviceCredentials.getCredentialsType()) { return switch (deviceCredentials.getCredentialsType()) {
case ACCESS_TOKEN: case ACCESS_TOKEN -> JacksonUtil.valueToTree(deviceCredentials.getCredentialsId());
return JacksonUtil.valueToTree(deviceCredentials.getCredentialsId()); case X509_CERTIFICATE -> JacksonUtil.valueToTree(deviceCredentials.getCredentialsValue());
case X509_CERTIFICATE: default -> JacksonUtil.fromString(deviceCredentials.getCredentialsValue(), JsonNode.class);
return JacksonUtil.valueToTree(deviceCredentials.getCredentialsValue()); };
}
return JacksonUtil.fromString(deviceCredentials.getCredentialsValue(), JsonNode.class);
} }
private void formatSimpleMqttCredentials(DeviceCredentials deviceCredentials) { private void formatSimpleMqttCredentials(DeviceCredentials deviceCredentials) {
@ -407,4 +407,11 @@ public class DeviceCredentialsServiceImpl extends AbstractCachedEntityService<St
} }
} }
private boolean isCredentialsChanged(DeviceCredentials oldCredentials, DeviceCredentials newCredentials) {
return !Objects.equals(oldCredentials.getCredentialsId(), newCredentials.getCredentialsId())
|| oldCredentials.getCredentialsType() != newCredentials.getCredentialsType()
|| !Objects.equals(oldCredentials.getCredentialsValue(), newCredentials.getCredentialsValue())
|| !Objects.equals(oldCredentials.getDeviceId(), newCredentials.getDeviceId());
}
} }

42
dao/src/test/java/org/thingsboard/server/dao/service/DeviceCredentialsServiceTest.java

@ -19,7 +19,10 @@ import com.datastax.oss.driver.api.core.uuid.Uuids;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Test; import org.junit.Test;
import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Assertions;
import org.mockito.Mockito;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.test.context.bean.override.mockito.MockitoBean;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.DeviceCredentialsId; import org.thingsboard.server.common.data.id.DeviceCredentialsId;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
@ -27,6 +30,7 @@ import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.data.security.DeviceCredentialsType;
import org.thingsboard.server.dao.device.DeviceCredentialsService; import org.thingsboard.server.dao.device.DeviceCredentialsService;
import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.eventsourcing.ActionEntityEvent;
import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.exception.DataValidationException;
@DaoSqlTest @DaoSqlTest
@ -36,6 +40,8 @@ public class DeviceCredentialsServiceTest extends AbstractServiceTest {
DeviceCredentialsService deviceCredentialsService; DeviceCredentialsService deviceCredentialsService;
@Autowired @Autowired
DeviceService deviceService; DeviceService deviceService;
@MockitoBean
ApplicationEventPublisher eventPublisher;
@Test @Test
public void testCreateDeviceCredentials() { public void testCreateDeviceCredentials() {
@ -185,5 +191,39 @@ public class DeviceCredentialsServiceTest extends AbstractServiceTest {
Assert.assertEquals(deviceCredentials, foundDeviceCredentials); Assert.assertEquals(deviceCredentials, foundDeviceCredentials);
deviceService.deleteDevice(tenantId, savedDevice.getId()); deviceService.deleteDevice(tenantId, savedDevice.getId());
} }
}
@Test
public void testUpdateDeviceCredentialsWithSameValuesDoesNotPublishEvent() {
Device device = new Device();
device.setTenantId(tenantId);
device.setName("My device");
device.setType("default");
Device savedDevice = deviceService.saveDevice(device);
try {
DeviceCredentials deviceCredentials = deviceCredentialsService.findDeviceCredentialsByDeviceId(tenantId, savedDevice.getId());
Assert.assertNotNull(deviceCredentials);
DeviceCredentials updatedCredentials = new DeviceCredentials(deviceCredentials.getId());
updatedCredentials.setDeviceId(deviceCredentials.getDeviceId());
updatedCredentials.setCredentialsType(deviceCredentials.getCredentialsType());
updatedCredentials.setCredentialsId(deviceCredentials.getCredentialsId());
updatedCredentials.setCredentialsValue(deviceCredentials.getCredentialsValue());
Mockito.reset(eventPublisher);
DeviceCredentials result = deviceCredentialsService.updateDeviceCredentials(tenantId, updatedCredentials);
Assert.assertEquals(deviceCredentials.getCredentialsId(), result.getCredentialsId());
Assert.assertEquals(deviceCredentials.getCredentialsType(), result.getCredentialsType());
Assert.assertEquals(deviceCredentials.getCredentialsValue(), result.getCredentialsValue());
Assert.assertEquals(deviceCredentials.getDeviceId(), result.getDeviceId());
Mockito.verify(eventPublisher, Mockito.never()).publishEvent(Mockito.any(ActionEntityEvent.class));
} finally {
deviceService.deleteDevice(tenantId, savedDevice.getId());
}
}
}

3
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java

@ -15,9 +15,6 @@
*/ */
package org.thingsboard.rule.engine.api; package org.thingsboard.rule.engine.api;
/**
* Created by ashvayka on 02.04.18.
*/
public interface RuleEngineTelemetryService { public interface RuleEngineTelemetryService {
void saveTimeseries(TimeseriesSaveRequest request); void saveTimeseries(TimeseriesSaveRequest request);

Loading…
Cancel
Save