From fdfcad8665dc2d5aa45c6e0b7661bc3093e60847 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Mon, 27 Oct 2025 10:23:06 +0200 Subject: [PATCH 1/2] Fix bulk import --- .../device/DeviceActorMessageProcessor.java | 11 ++--- .../csv/AbstractBulkImportService.java | 2 +- .../data/security/DeviceCredentials.java | 4 ++ .../DeviceAttributesEventNotificationMsg.java | 6 +-- .../mqtt/MqttSslHandlerProvider.java | 3 -- .../transport/mqtt/MqttTransportHandler.java | 3 -- .../mqtt/MqttTransportServerInitializer.java | 3 -- .../transport/mqtt/MqttTransportService.java | 3 -- .../mqtt/adaptors/JsonMqttAdaptor.java | 6 +-- .../mqtt/adaptors/MqttTransportAdaptor.java | 4 +- .../mqtt/adaptors/ProtoMqttAdaptor.java | 3 +- .../AbstractGatewayDeviceSessionContext.java | 3 -- .../AbstractGatewaySessionHandler.java | 11 ++--- .../mqtt/session/DeviceSessionCtx.java | 49 +++++-------------- .../session/GatewayDeviceSessionContext.java | 3 -- .../mqtt/session/GatewaySessionHandler.java | 3 -- .../MqttDeviceAwareSessionContext.java | 5 +- .../mqtt/session/MqttTopicMatcher.java | 12 ++--- .../session/SparkplugNodeSessionHandler.java | 3 -- .../device/DeviceCredentialsEvictEvent.java | 4 +- .../device/DeviceCredentialsServiceImpl.java | 27 ++++++---- .../service/DeviceCredentialsServiceTest.java | 42 +++++++++++++++- .../api/RuleEngineTelemetryService.java | 3 -- 23 files changed, 98 insertions(+), 115 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 c143214e2b..b24b803950 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 @@ -115,14 +115,9 @@ import java.util.concurrent.TimeUnit; import java.util.function.Consumer; import java.util.stream.Collectors; - -/** - * @author Andrew Shvayka - */ @Slf4j public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { - static final String SESSION_TIMEOUT_MESSAGE = "session timeout!"; final TenantId tenantId; final DeviceId deviceId; final LinkedHashMapRemoveEldest sessions; @@ -178,7 +173,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso private EdgeId findRelatedEdgeId() { List result = 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); if (relationToEdge.getFrom() != null && relationToEdge.getFrom().getId() != null) { log.trace("[{}][{}] found edge [{}] for device", tenantId, deviceId, relationToEdge.getFrom().getId()); @@ -501,7 +496,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso UUID sessionId = getSessionId(sessionInfo); DeviceId deviceId = new DeviceId(new UUID(msg.getDeviceIdMSB(), msg.getDeviceIdLSB())); ListenableFuture registrationFuture = systemContext.getClaimDevicesService() - .registerClaimingInfo(tenantId, deviceId, msg.getSecretKey(), msg.getDurationMs()); + .registerClaimingInfo(tenantId, deviceId, msg.getSecretKey(), msg.getDurationMs()); Futures.addCallback(registrationFuture, new FutureCallback<>() { @Override public void onSuccess(Void result) { @@ -723,7 +718,7 @@ public class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcesso toDeviceRpcPendingMap.remove(requestId); status = RpcStatus.FAILED; 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 { md.setRetries(md.getRetries() + 1); diff --git a/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java b/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java index 9850e2d1a1..63a4f100e1 100644 --- a/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java +++ b/application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java @@ -186,7 +186,7 @@ public abstract class AbstractBulkImportService kvs.add(dataEntry.getKey().getKey(), dataEntry.getValue().toJsonPrimitive())); return Map.entry(kvType, kvs); }) - .filter(kvsEntry -> kvsEntry.getValue().entrySet().size() > 0) + .filter(kvsEntry -> !kvsEntry.getValue().entrySet().isEmpty()) .forEach(kvsEntry -> { BulkImportColumnType kvType = kvsEntry.getKey(); if (kvType == BulkImportColumnType.SHARED_ATTRIBUTE || kvType == BulkImportColumnType.SERVER_ATTRIBUTE) { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/security/DeviceCredentials.java b/common/data/src/main/java/org/thingsboard/server/common/data/security/DeviceCredentials.java index ad73510d9e..6d1dbdcb4a 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/security/DeviceCredentials.java +++ b/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.DeviceId; +import java.io.Serial; + @Schema @EqualsAndHashCode(callSuper = true) public class DeviceCredentials extends BaseData implements DeviceCredentialsFilter, HasVersion { + @Serial private static final long serialVersionUID = -7869261127032877765L; + private DeviceId deviceId; private DeviceCredentialsType credentialsType; private String credentialsId; diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceAttributesEventNotificationMsg.java b/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceAttributesEventNotificationMsg.java index efee3dcb97..f647714c63 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceAttributesEventNotificationMsg.java +++ b/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.ToDeviceActorNotificationMsg; +import java.io.Serial; import java.util.HashSet; import java.util.List; import java.util.Set; -/** - * @author Andrew Shvayka - */ @Data public class DeviceAttributesEventNotificationMsg implements ToDeviceActorNotificationMsg { + @Serial private static final long serialVersionUID = 2422071590415277039L; private final TenantId tenantId; @@ -56,4 +55,5 @@ public class DeviceAttributesEventNotificationMsg implements ToDeviceActorNotifi public MsgType getMsgType() { return MsgType.DEVICE_ATTRIBUTES_UPDATE_TO_DEVICE_ACTOR_MSG; } + } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java index a756b84fe8..867b39cd23 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java +++ b/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.TimeUnit; -/** - * Created by valerii.sosliuk on 11/6/16. - */ @Slf4j @Component("MqttSslHandlerProvider") @ConditionalOnProperty(prefix = "transport.mqtt.ssl", value = "enabled", havingValue = "true", matchIfMissing = false) 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 7f9da9e0b9..4397802951 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 @@ -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.SparkplugTopicService.parseTopicPublish; -/** - * @author Andrew Shvayka - */ @Slf4j public class MqttTransportHandler extends ChannelInboundHandlerAdapter implements GenericFutureListener>, SessionMsgListener { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportServerInitializer.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportServerInitializer.java index c630def89f..5d127fccb7 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportServerInitializer.java +++ b/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.ProxyIpFilter; -/** - * @author Andrew Shvayka - */ public class MqttTransportServerInitializer extends ChannelInitializer { private final MqttTransportContext context; diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java index c7ff8912aa..d62f35d561 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java +++ b/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; -/** - * @author Andrew Shvayka - */ @Service("MqttTransportService") @TbMqttTransportComponent @Slf4j 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 64f7e5df4a..61a886a7db 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 @@ -45,10 +45,6 @@ import java.util.UUID; import static org.thingsboard.server.common.data.device.profile.MqttTopics.DEVICE_SOFTWARE_FIRMWARE_RESPONSES_TOPIC_FORMAT; - -/** - * @author Andrew Shvayka - */ @Component @Slf4j public class JsonMqttAdaptor implements MqttTransportAdaptor { @@ -122,7 +118,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { public Optional convertToGatewayPublish(MqttDeviceAwareSessionContext ctx, String deviceName, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException { return processConvertFromGatewayAttributeResponseMsg(ctx, deviceName, responseMsg); } - + @Override public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.AttributeUpdateNotificationMsg notificationMsg, String topic) { return Optional.of(createMqttPublishMsg(ctx, topic, JsonConverter.toJson(notificationMsg))); 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 012709c31f..356004246f 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 @@ -41,9 +41,6 @@ import org.thingsboard.server.transport.mqtt.session.MqttDeviceAwareSessionConte import java.util.Optional; -/** - * @author Andrew Shvayka - */ public interface MqttTransportAdaptor { ByteBufAllocator ALLOCATOR = new UnpooledByteBufAllocator(false); @@ -90,4 +87,5 @@ public interface MqttTransportAdaptor { payload.writeBytes(payloadInBytes); return new MqttPublishMessage(mqttFixedHeader, header, payload); } + } 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 fd10ad750e..4472bb77d2 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 @@ -48,7 +48,7 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { public TransportProtos.PostTelemetryMsg convertToPostTelemetry(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException { DeviceSessionCtx deviceSessionCtx = (DeviceSessionCtx) ctx; byte[] bytes = toBytes(inbound.payload()); - Descriptors.Descriptor telemetryDynamicMsgDescriptor = ProtoConverter.validateDescriptor(deviceSessionCtx.getTelemetryDynamicMsgDescriptor()); + Descriptors.Descriptor telemetryDynamicMsgDescriptor = ProtoConverter.validateDescriptor(deviceSessionCtx.getTelemetryDynamicMessageDescriptor()); try { return JsonConverter.convertToTelemetryProto(JsonParser.parseString(ProtoConverter.dynamicMsgToJson(bytes, telemetryDynamicMsgDescriptor))); } catch (Exception e) { @@ -228,4 +228,5 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor { private int getRequestId(String topicName, String topic) { return Integer.parseInt(topicName.substring(topic.length())); } + } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewayDeviceSessionContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewayDeviceSessionContext.java index 8ce0fc5c1e..5bb3c9c28f 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewayDeviceSessionContext.java +++ b/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.concurrent.ConcurrentMap; -/** - * Created by ashvayka on 19.01.17. - */ @ToString(callSuper = true) @Slf4j public abstract class AbstractGatewayDeviceSessionContext extends MqttDeviceAwareSessionContext implements SessionMsgListener { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java index fee2f04e24..6542492238 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/AbstractGatewaySessionHandler.java +++ b/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.messageName; -/** - * Created by ashvayka on 19.01.17. - */ @Slf4j public abstract class AbstractGatewaySessionHandler { @@ -116,6 +113,7 @@ public abstract class AbstractGatewaySessionHandler deviceCreationLockMap; + @Getter private final ConcurrentMap devices; private final ConcurrentMap> deviceFutures; protected final ConcurrentMap mqttQoSMap; @@ -821,11 +819,7 @@ public abstract class AbstractGatewaySessionHandler getDevices() { - return this.devices; - } - - protected TransportServiceCallback getAggregatePubAckCallback( + protected TransportServiceCallback getAggregatePubAckCallback( final ChannelHandlerContext ctx, final int msgId, final String deviceName, @@ -915,4 +909,5 @@ public abstract class AbstractGatewaySessionHandler getDefaultAdaptor(); + case V2_JSON -> context.getJsonMqttAdaptor(); + case V2_PROTO -> context.getProtoMqttAdaptor(); + default -> useJsonPayloadFormatForDefaultDownlinkTopics ? context.getJsonMqttAdaptor() : getDefaultAdaptor(); + }; } private MqttTransportAdaptor getDefaultAdaptor() { @@ -269,7 +246,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext { } } - public Collection getMsgQueueSnapshot(){ + public Collection getMsgQueueSnapshot() { return Collections.unmodifiableCollection(msgQueue); } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionContext.java index a48ee50054..24d7888a81 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionContext.java +++ b/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; -/** - * Created by nickAS21 on 26.12.22 - */ @ToString(callSuper = true) public class GatewayDeviceSessionContext extends AbstractGatewayDeviceSessionContext { 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 2980b2ad60..487fe1a541 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 @@ -28,9 +28,6 @@ import org.thingsboard.server.gen.transport.TransportProtos; import java.util.Optional; import java.util.UUID; -/** - * Created by nickAS21 on 26.12.22 - */ @Slf4j public class GatewaySessionHandler extends AbstractGatewaySessionHandler { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttDeviceAwareSessionContext.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttDeviceAwareSessionContext.java index 9db81a4d8b..24d23d7bdc 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttDeviceAwareSessionContext.java +++ b/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.stream.Collectors; -/** - * Created by ashvayka on 30.08.18. - */ @ToString(callSuper = true) public abstract class MqttDeviceAwareSessionContext extends DeviceAwareSessionContext { @@ -47,7 +44,7 @@ public abstract class MqttDeviceAwareSessionContext extends DeviceAwareSessionCo .stream() .filter(entry -> entry.getKey().matches(topic)) .map(Map.Entry::getValue) - .collect(Collectors.toList()); + .toList(); if (!qosList.isEmpty()) { return MqttQoS.valueOf(qosList.get(0)); } else { diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttTopicMatcher.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttTopicMatcher.java index 950a1f5bd9..ad3c7b9897 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/MqttTopicMatcher.java +++ b/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; +import lombok.Getter; + import java.util.regex.Pattern; public class MqttTopicMatcher { + @Getter private final String topic; private final Pattern topicRegex; public MqttTopicMatcher(String topic) { - if(topic == null){ + if (topic == null) { throw new NullPointerException("topic"); } this.topic = topic; this.topicRegex = Pattern.compile(topic.replace("+", "[^/]+").replace("#", ".+") + "$"); } - public String getTopic() { - return topic; - } - - public boolean matches(String topic){ + public boolean matches(String topic) { return this.topicRegex.matcher(topic).matches(); } @@ -52,4 +51,5 @@ public class MqttTopicMatcher { public int hashCode() { return topic.hashCode(); } + } diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java index f9aba81c90..9ba4c3f989 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/SparkplugNodeSessionHandler.java +++ b/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_STATE_REGEXP; -/** - * Created by nickAS21 on 12.12.22 - */ @Slf4j @SpecVersion(spec = "sparkplug", version = "3.0.0") public class SparkplugNodeSessionHandler extends AbstractGatewaySessionHandler { diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsEvictEvent.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsEvictEvent.java index ca265e78b1..33b410bcfe 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsEvictEvent.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsEvictEvent.java @@ -18,9 +18,9 @@ package org.thingsboard.server.dao.device; import lombok.Data; @Data -class DeviceCredentialsEvictEvent { +public class DeviceCredentialsEvictEvent { - private final String newCedentialsId; + private final String newCredentialsId; private final String oldCredentialsId; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java index 12859129a2..3423a74543 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java +++ b/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.service.validator.DeviceCredentialsDataValidator; +import java.util.Objects; + import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validateString; @@ -61,8 +63,8 @@ public class DeviceCredentialsServiceImpl extends AbstractCachedEntityService JacksonUtil.valueToTree(deviceCredentials.getCredentialsId()); + case X509_CERTIFICATE -> JacksonUtil.valueToTree(deviceCredentials.getCredentialsValue()); + default -> JacksonUtil.fromString(deviceCredentials.getCredentialsValue(), JsonNode.class); + }; } private void formatSimpleMqttCredentials(DeviceCredentials deviceCredentials) { @@ -407,4 +407,11 @@ public class DeviceCredentialsServiceImpl extends AbstractCachedEntityService Date: Tue, 28 Oct 2025 12:28:42 +0200 Subject: [PATCH 2/2] coaps: add to test -> "coap.dtls.x509.skip_validity_check_for_client_cert=true" --- .../coap/security/AbstractCoapSecurityIntegrationTest.java | 1 + 1 file changed, 1 insertion(+) diff --git a/application/src/test/java/org/thingsboard/server/transport/coap/security/AbstractCoapSecurityIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/coap/security/AbstractCoapSecurityIntegrationTest.java index 72f8cc5b4f..839b8c7cad 100644 --- a/application/src/test/java/org/thingsboard/server/transport/coap/security/AbstractCoapSecurityIntegrationTest.java +++ b/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.dtls.enabled=true", "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", "service.integrations.supported=ALL", "transport.coap.enabled=true",