From 0245ff24a8151dc5a39452cde2402ec74b09a0c1 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 4 Sep 2020 15:35:50 +0300 Subject: [PATCH 1/4] Transport Profile cache --- .../transport/TransportProfileCache.java | 35 +++++++++ .../service/DefaultTransportProfileCache.java | 77 +++++++++++++++++++ .../service/DefaultTransportService.java | 73 +++++------------- 3 files changed, 130 insertions(+), 55 deletions(-) create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportProfileCache.java create mode 100644 common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportProfileCache.java diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportProfileCache.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportProfileCache.java new file mode 100644 index 0000000000..f34b76ade9 --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportProfileCache.java @@ -0,0 +1,35 @@ +/** + * Copyright © 2016-2020 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport; + +import com.google.protobuf.ByteString; +import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.id.DeviceProfileId; + +import java.util.Optional; + +public interface TransportProfileCache { + + + DeviceProfile getOrCreate(DeviceProfileId id, ByteString profileBody); + + DeviceProfile get(DeviceProfileId id); + + void put(DeviceProfile profile); + + DeviceProfile put(ByteString profileBody); + +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportProfileCache.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportProfileCache.java new file mode 100644 index 0000000000..afa8e15e20 --- /dev/null +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportProfileCache.java @@ -0,0 +1,77 @@ +/** + * Copyright © 2016-2020 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.transport.service; + +import com.google.protobuf.ByteString; +import lombok.extern.slf4j.Slf4j; +import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.DeviceProfile; +import org.thingsboard.server.common.data.id.DeviceProfileId; +import org.thingsboard.server.common.transport.TransportProfileCache; +import org.thingsboard.server.common.transport.util.DataDecodingEncodingService; + +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; + +@Slf4j +@Component +@ConditionalOnExpression("('${service.type:null}'=='monolith' && '${transport.api_enabled:true}'=='true') || '${service.type:null}'=='tb-transport'") +public class DefaultTransportProfileCache implements TransportProfileCache { + + private final ConcurrentMap deviceProfiles = new ConcurrentHashMap<>(); + + private final DataDecodingEncodingService dataDecodingEncodingService; + + public DefaultTransportProfileCache(DataDecodingEncodingService dataDecodingEncodingService) { + this.dataDecodingEncodingService = dataDecodingEncodingService; + } + + @Override + public DeviceProfile getOrCreate(DeviceProfileId id, ByteString profileBody) { + DeviceProfile profile = deviceProfiles.get(id); + if (profile == null) { + Optional deviceProfile = dataDecodingEncodingService.decode(profileBody.toByteArray()); + if (deviceProfile.isPresent()) { + profile = deviceProfile.get(); + deviceProfiles.put(id, profile); + } + } + return profile; + } + + @Override + public DeviceProfile get(DeviceProfileId id) { + return deviceProfiles.get(id); + } + + @Override + public void put(DeviceProfile profile) { + deviceProfiles.put(profile.getId(), profile); + } + + @Override + public DeviceProfile put(ByteString profileBody) { + Optional deviceProfile = dataDecodingEncodingService.decode(profileBody.toByteArray()); + if (deviceProfile.isPresent()) { + put(deviceProfile.get()); + return deviceProfile.get(); + } else { + return null; + } + } +} diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 3bf9d78b48..b572835971 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -43,6 +43,7 @@ import org.thingsboard.server.common.msg.session.SessionMsgType; import org.thingsboard.server.common.msg.tools.TbRateLimits; import org.thingsboard.server.common.msg.tools.TbRateLimitsException; import org.thingsboard.server.common.transport.SessionMsgListener; +import org.thingsboard.server.common.transport.TransportProfileCache; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.common.transport.TransportServiceCallback; import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse; @@ -121,8 +122,7 @@ public class DefaultTransportService implements TransportService { private final PartitionService partitionService; private final TbServiceInfoProvider serviceInfoProvider; private final StatsFactory statsFactory; - private final DataDecodingEncodingService dataDecodingEncodingService; - + private final TransportProfileCache transportProfileCache; protected TbQueueRequestTemplate, TbProtoQueueMsg> transportApiRequestTemplate; protected TbQueueProducer> ruleEngineMsgProducer; @@ -141,7 +141,6 @@ public class DefaultTransportService implements TransportService { //TODO 3.2: @ybondarenko Implement cleanup of this maps. private final ConcurrentMap perTenantLimits = new ConcurrentHashMap<>(); private final ConcurrentMap perDeviceLimits = new ConcurrentHashMap<>(); - private final ConcurrentMap deviceProfiles = new ConcurrentHashMap<>(); private ExecutorService mainConsumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("transport-consumer")); private volatile boolean stopped = false; @@ -151,13 +150,13 @@ public class DefaultTransportService implements TransportService { TbQueueProducerProvider producerProvider, PartitionService partitionService, StatsFactory statsFactory, - DataDecodingEncodingService dataDecodingEncodingService) { + TransportProfileCache transportProfileCache) { this.serviceInfoProvider = serviceInfoProvider; this.queueProvider = queueProvider; this.producerProvider = producerProvider; this.partitionService = partitionService; this.statsFactory = statsFactory; - this.dataDecodingEncodingService = dataDecodingEncodingService; + this.transportProfileCache = transportProfileCache; } @PostConstruct @@ -276,14 +275,7 @@ public class DefaultTransportService implements TransportService { result.deviceInfo(tdi); ByteString profileBody = msg.getProfileBody(); if (profileBody != null && !profileBody.isEmpty()) { - DeviceProfile profile = deviceProfiles.get(tdi.getDeviceProfileId()); - if (profile == null) { - Optional deviceProfile = dataDecodingEncodingService.decode(profileBody.toByteArray()); - if (deviceProfile.isPresent()) { - profile = deviceProfile.get(); - deviceProfiles.put(tdi.getDeviceProfileId(), profile); - } - } + DeviceProfile profile = transportProfileCache.getOrCreate(tdi.getDeviceProfileId(), profileBody); if (transportType != DeviceTransportType.DEFAULT && profile != null && profile.getTransportType() != DeviceTransportType.DEFAULT && profile.getTransportType() != transportType) { log.debug("[{}] Device profile [{}] has different transport type: {}, expected: {}", tdi.getDeviceId(), tdi.getDeviceProfileId(), profile.getTransportType(), transportType); @@ -309,15 +301,7 @@ public class DefaultTransportService implements TransportService { result.deviceInfo(tdi); ByteString profileBody = msg.getProfileBody(); if (profileBody != null && !profileBody.isEmpty()) { - DeviceProfile profile = deviceProfiles.get(tdi.getDeviceProfileId()); - if (profile == null) { - Optional deviceProfile = dataDecodingEncodingService.decode(profileBody.toByteArray()); - if (deviceProfile.isPresent()) { - profile = deviceProfile.get(); - deviceProfiles.put(tdi.getDeviceProfileId(), profile); - } - } - result.deviceProfile(profile); + result.deviceProfile(transportProfileCache.getOrCreate(tdi.getDeviceProfileId(), profileBody)); } } return result.build(); @@ -629,8 +613,10 @@ public class DefaultTransportService implements TransportService { } } else { if (toSessionMsg.hasDeviceProfileUpdateMsg()) { - Optional deviceProfile = dataDecodingEncodingService.decode(toSessionMsg.getDeviceProfileUpdateMsg().getData().toByteArray()); - deviceProfile.ifPresent(this::onProfileUpdate); + DeviceProfile deviceProfile = transportProfileCache.put(toSessionMsg.getDeviceProfileUpdateMsg().getData()); + if (deviceProfile != null) { + onProfileUpdate(deviceProfile); + } } else { //TODO: should we notify the device actor about missed session? log.debug("[{}] Missing session.", sessionId); @@ -640,7 +626,7 @@ public class DefaultTransportService implements TransportService { @Override public void getDeviceProfile(DeviceProfileId deviceProfileId, TransportServiceCallback callback) { - DeviceProfile deviceProfile = deviceProfiles.get(deviceProfileId); + DeviceProfile deviceProfile = transportProfileCache.get(deviceProfileId); if (deviceProfile != null) { callback.onSuccess(deviceProfile); } else { @@ -653,14 +639,13 @@ public class DefaultTransportService implements TransportService { TransportApiRequestMsg.newBuilder().setGetDeviceProfileRequestMsg(msg).build()); AsyncCallbackTemplate.withCallback(transportApiRequestTemplate.send(protoMsg), response -> { - byte[] devProfileBody = response.getValue().getGetDeviceProfileResponseMsg().getData().toByteArray(); - if (devProfileBody != null && devProfileBody.length > 0) { - Optional deviceProfileOpt = dataDecodingEncodingService.decode(devProfileBody); - if (deviceProfileOpt.isPresent()) { - deviceProfiles.put(deviceProfileOpt.get().getId(), deviceProfile); - callback.onSuccess(deviceProfileOpt.get()); + ByteString devProfileBody = response.getValue().getGetDeviceProfileResponseMsg().getData(); + if (devProfileBody != null && !devProfileBody.isEmpty()) { + DeviceProfile profile = transportProfileCache.put(devProfileBody); + if (profile != null) { + callback.onSuccess(profile); } else { - log.warn("Failed to decode device profile: {}", Arrays.toString(devProfileBody)); + log.warn("Failed to decode device profile: {}", devProfileBody); callback.onError(new IllegalArgumentException("Failed to decode device profile!")); } } else { @@ -673,7 +658,6 @@ public class DefaultTransportService implements TransportService { @Override public void onProfileUpdate(DeviceProfile deviceProfile) { - deviceProfiles.put(deviceProfile.getId(), deviceProfile); long deviceProfileIdMSB = deviceProfile.getId().getId().getMostSignificantBits(); long deviceProfileIdLSB = deviceProfile.getId().getId().getLeastSignificantBits(); sessions.forEach((id, md) -> { @@ -736,7 +720,7 @@ public class DefaultTransportService implements TransportService { private RuleChainId resolveRuleChainId(TransportProtos.SessionInfoProto sessionInfo) { DeviceProfileId deviceProfileId = new DeviceProfileId(new UUID(sessionInfo.getDeviceProfileIdMSB(), sessionInfo.getDeviceProfileIdLSB())); - DeviceProfile deviceProfile = deviceProfiles.get(deviceProfileId); + DeviceProfile deviceProfile = transportProfileCache.get(deviceProfileId); RuleChainId ruleChainId; if (deviceProfile == null) { log.warn("[{}] Device profile is null!", deviceProfileId); @@ -747,27 +731,6 @@ public class DefaultTransportService implements TransportService { return ruleChainId; } - private ListenableFuture> extractProfile(ListenableFuture> send, - Function hasDeviceInfo, - Function deviceInfoF, - Function profileBodyF) { - return Futures.transform(send, response -> { - T value = response.getValue(); - if (hasDeviceInfo.apply(value)) { - TransportProtos.DeviceInfoProto deviceInfo = deviceInfoF.apply(value); - ByteString profileBody = profileBodyF.apply(value); - if (profileBody != null && !profileBody.isEmpty()) { - DeviceProfileId deviceProfileId = new DeviceProfileId(new UUID(deviceInfo.getDeviceProfileIdMSB(), deviceInfo.getDeviceProfileIdLSB())); - if (!deviceProfiles.containsKey(deviceProfileId)) { - Optional deviceProfile = dataDecodingEncodingService.decode(profileBody.toByteArray()); - deviceProfile.ifPresent(profile -> deviceProfiles.put(deviceProfileId, profile)); - } - } - } - return response; - }, transportCallbackExecutor); - } - private class TransportTbQueueCallback implements TbQueueCallback { private final TransportServiceCallback callback; From 73a1a79821665486408306b4455810ae72e0dc7d Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 4 Sep 2020 16:51:56 +0300 Subject: [PATCH 2/4] Device Profile updates --- .../install/SqlDatabaseUpgradeService.java | 6 ++- .../dao/device/DeviceProfileService.java | 6 +-- .../transport/mqtt/MqttTransportHandler.java | 9 +++- .../server/dao/device/DeviceDao.java | 11 ++++ .../server/dao/device/DeviceProfileDao.java | 2 + .../dao/device/DeviceProfileServiceImpl.java | 54 ++++++++++++++----- .../server/dao/device/DeviceServiceImpl.java | 20 +++++-- .../sql/device/DeviceProfileRepository.java | 4 ++ .../dao/sql/device/DeviceRepository.java | 8 +++ .../server/dao/sql/device/JpaDeviceDao.java | 10 ++++ .../dao/sql/device/JpaDeviceProfileDao.java | 4 ++ 11 files changed, 111 insertions(+), 23 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java index 4f83c672c3..be35b9aeb2 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java @@ -355,10 +355,12 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService pageData = tenantService.findTenants(pageLink); for (Tenant tenant : pageData.getData()) { List deviceTypes = deviceService.findDeviceTypesByTenantId(tenant.getId()).get(); - deviceProfileService.findOrCreateDefaultDeviceProfile(tenant.getId()); + try { + deviceProfileService.createDefaultDeviceProfile(tenant.getId()); + } catch (Exception e){} for (EntitySubtype deviceType : deviceTypes) { try { - deviceProfileService.createDeviceProfile(tenant.getId(), deviceType.getType()); + deviceProfileService.findOrCreateDeviceProfile(tenant.getId(), deviceType.getType()); } catch (Exception e) { } } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceProfileService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceProfileService.java index 8441e7e11d..e38bac68e5 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceProfileService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceProfileService.java @@ -27,6 +27,8 @@ public interface DeviceProfileService { DeviceProfile findDeviceProfileById(TenantId tenantId, DeviceProfileId deviceProfileId); + DeviceProfile findDeviceProfileByName(TenantId tenantId, String profileName); + DeviceProfileInfo findDeviceProfileInfoById(TenantId tenantId, DeviceProfileId deviceProfileId); DeviceProfile saveDeviceProfile(DeviceProfile deviceProfile); @@ -37,12 +39,10 @@ public interface DeviceProfileService { PageData findDeviceProfileInfos(TenantId tenantId, PageLink pageLink); - DeviceProfile findOrCreateDefaultDeviceProfile(TenantId tenantId); + DeviceProfile findOrCreateDeviceProfile(TenantId tenantId, String profileName); DeviceProfile createDefaultDeviceProfile(TenantId tenantId); - DeviceProfile createDeviceProfile(TenantId tenantId, String profileName); - DeviceProfile findDefaultDeviceProfile(TenantId tenantId); DeviceProfileInfo findDefaultDeviceProfileInfo(TenantId tenantId); 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 05254f612e..019e47f520 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 @@ -39,6 +39,7 @@ import io.netty.util.concurrent.Future; import io.netty.util.concurrent.GenericFutureListener; import lombok.extern.slf4j.Slf4j; import org.springframework.util.StringUtils; +import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.device.profile.MqttTopics; import org.thingsboard.server.common.msg.EncryptionUtil; @@ -574,7 +575,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement try { adaptor.convertToPublish(deviceSessionCtx, rpcResponse).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); } catch (Exception e) { - log.trace("[{}] Failed to convert device RPC commandto MQTT msg", sessionId, e); + log.trace("[{}] Failed to convert device RPC command to MQTT msg", sessionId, e); } } + + @Override + public void onProfileUpdate(DeviceProfile deviceProfile) { + deviceSessionCtx.getDeviceInfo().setDeviceType(deviceProfile.getName()); + sessionInfo = SessionInfoProto.newBuilder().mergeFrom(sessionInfo).setDeviceType(deviceProfile.getName()).build(); + } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java index bd000f3181..16f849fe73 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java @@ -184,4 +184,15 @@ public interface DeviceDao extends Dao { ListenableFuture findDeviceByTenantIdAndIdAsync(TenantId tenantId, UUID id); Long countDevicesByDeviceProfileId(TenantId tenantId, UUID deviceProfileId); + + /** + * Find devices by tenantId, profileId and page link. + * + * @param tenantId the tenantId + * @param profileId the profileId + * @param pageLink the page link + * @return the list of device objects + */ + PageData findDevicesByTenantIdAndProfileId(UUID tenantId, UUID profileId, PageLink pageLink); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileDao.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileDao.java index 34c12ab90c..267aff358e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileDao.java @@ -37,4 +37,6 @@ public interface DeviceProfileDao extends Dao { DeviceProfile findDefaultDeviceProfile(TenantId tenantId); DeviceProfileInfo findDefaultDeviceProfileInfo(TenantId tenantId); + + DeviceProfile findByName(TenantId tenantId, String profileName); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileServiceImpl.java index 529e1992ac..b40ad102eb 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileServiceImpl.java @@ -23,10 +23,12 @@ import org.springframework.cache.Cache; import org.springframework.cache.CacheManager; import org.springframework.cache.annotation.Cacheable; import org.springframework.stereotype.Service; +import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfileInfo; import org.thingsboard.server.common.data.DeviceProfileType; import org.thingsboard.server.common.data.DeviceTransportType; +import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.device.profile.DefaultDeviceProfileConfiguration; import org.thingsboard.server.common.data.device.profile.DefaultDeviceProfileTransportConfiguration; @@ -44,6 +46,7 @@ import org.thingsboard.server.dao.tenant.TenantDao; import java.util.Arrays; import java.util.Collections; +import java.util.List; import static org.thingsboard.server.common.data.CacheConstants.DEVICE_PROFILE_CACHE; import static org.thingsboard.server.dao.service.Validator.validateId; @@ -54,6 +57,7 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D private static final String INCORRECT_TENANT_ID = "Incorrect tenantId "; private static final String INCORRECT_DEVICE_PROFILE_ID = "Incorrect deviceProfileId "; + private static final String INCORRECT_DEVICE_PROFILE_NAME = "Incorrect deviceProfileName "; @Autowired private DeviceProfileDao deviceProfileDao; @@ -61,6 +65,9 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D @Autowired private DeviceDao deviceDao; + @Autowired + private DeviceService deviceService; + @Autowired private TenantDao tenantDao; @@ -75,6 +82,13 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D return deviceProfileDao.findById(tenantId, deviceProfileId.getId()); } + @Override + public DeviceProfile findDeviceProfileByName(TenantId tenantId, String profileName) { + log.trace("Executing findDeviceProfileByName [{}][{}]", tenantId, profileName); + Validator.validateString(profileName, INCORRECT_DEVICE_PROFILE_NAME + profileName); + return deviceProfileDao.findByName(tenantId, profileName); + } + @Cacheable(cacheNames = DEVICE_PROFILE_CACHE, key = "{'info', #deviceProfileId.id}") @Override public DeviceProfileInfo findDeviceProfileInfoById(TenantId tenantId, DeviceProfileId deviceProfileId) { @@ -87,6 +101,10 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D public DeviceProfile saveDeviceProfile(DeviceProfile deviceProfile) { log.trace("Executing saveDeviceProfile [{}]", deviceProfile); deviceProfileValidator.validate(deviceProfile, DeviceProfile::getTenantId); + DeviceProfile oldDeviceProfile = null; + if (deviceProfile.getId() != null) { + oldDeviceProfile = deviceProfileDao.findById(deviceProfile.getTenantId(), deviceProfile.getId().getId()); + } DeviceProfile savedDeviceProfile; try { savedDeviceProfile = deviceProfileDao.save(deviceProfile.getTenantId(), deviceProfile); @@ -101,10 +119,23 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D Cache cache = cacheManager.getCache(DEVICE_PROFILE_CACHE); cache.evict(Collections.singletonList(savedDeviceProfile.getId().getId())); cache.evict(Arrays.asList("info", savedDeviceProfile.getId().getId())); + cache.evict(Arrays.asList(deviceProfile.getTenantId().getId(), deviceProfile.getName())); if (savedDeviceProfile.isDefault()) { cache.evict(Arrays.asList("default", savedDeviceProfile.getTenantId().getId())); cache.evict(Arrays.asList("default", "info", savedDeviceProfile.getTenantId().getId())); } + if (oldDeviceProfile != null && !oldDeviceProfile.getName().equals(deviceProfile.getName())) { + PageLink pageLink = new PageLink(100); + PageData pageData; + do { + pageData = deviceDao.findDevicesByTenantIdAndProfileId(deviceProfile.getTenantId().getId(), deviceProfile.getUuidId(), pageLink); + for (Device device : pageData.getData()) { + device.setType(deviceProfile.getName()); + deviceService.saveDevice(device); + } + pageLink = pageLink.nextPageLink(); + } while (pageData.hasNext()); + } return savedDeviceProfile; } @@ -116,10 +147,11 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D if (deviceProfile != null && deviceProfile.isDefault()) { throw new DataValidationException("Deletion of Default Device Profile is prohibited!"); } - this.removeDeviceProfile(tenantId, deviceProfileId); + this.removeDeviceProfile(tenantId, deviceProfile); } - private void removeDeviceProfile(TenantId tenantId, DeviceProfileId deviceProfileId) { + private void removeDeviceProfile(TenantId tenantId, DeviceProfile deviceProfile) { + DeviceProfileId deviceProfileId = deviceProfile.getId(); try { deviceProfileDao.removeById(tenantId, deviceProfileId.getId()); } catch (Exception t) { @@ -134,6 +166,7 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D Cache cache = cacheManager.getCache(DEVICE_PROFILE_CACHE); cache.evict(Collections.singletonList(deviceProfileId.getId())); cache.evict(Arrays.asList("info", deviceProfileId.getId())); + cache.evict(Arrays.asList(tenantId.getId(), deviceProfile.getName())); } @Override @@ -152,12 +185,13 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D return deviceProfileDao.findDeviceProfileInfos(tenantId, pageLink); } + @Cacheable(cacheNames = DEVICE_PROFILE_CACHE, key = "{#tenantId.id, #name}") @Override - public DeviceProfile findOrCreateDefaultDeviceProfile(TenantId tenantId) { + public DeviceProfile findOrCreateDeviceProfile(TenantId tenantId, String name) { log.trace("Executing findOrCreateDefaultDeviceProfile"); - DeviceProfile deviceProfile = findDefaultDeviceProfile(tenantId); + DeviceProfile deviceProfile = findDeviceProfileByName(tenantId, name); if (deviceProfile == null) { - deviceProfile = this.createDefaultDeviceProfile(tenantId); + deviceProfile = this.doCreateDefaultDeviceProfile(tenantId, name, name.equals("default")); } return deviceProfile; } @@ -168,12 +202,6 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D return doCreateDefaultDeviceProfile(tenantId, "default", true); } - @Override - public DeviceProfile createDeviceProfile(TenantId tenantId, String profileName) { - log.trace("Executing createDefaultDeviceProfile tenantId [{}], profileName [{}]", tenantId, profileName); - return doCreateDefaultDeviceProfile(tenantId, profileName, false); - } - private DeviceProfile doCreateDefaultDeviceProfile(TenantId tenantId, String profileName, boolean defaultProfile) { validateId(tenantId, INCORRECT_TENANT_ID + tenantId); DeviceProfile deviceProfile = new DeviceProfile(); @@ -227,6 +255,7 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D deviceProfileDao.save(tenantId, deviceProfile); cache.evict(Collections.singletonList(previousDefaultDeviceProfile.getId().getId())); cache.evict(Arrays.asList("info", previousDefaultDeviceProfile.getId().getId())); + cache.evict(Arrays.asList(tenantId.getId(), previousDefaultDeviceProfile.getName())); changed = true; } if (changed) { @@ -234,6 +263,7 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D cache.evict(Arrays.asList("info", deviceProfile.getId().getId())); cache.evict(Arrays.asList("default", tenantId.getId())); cache.evict(Arrays.asList("default", "info", tenantId.getId())); + cache.evict(Arrays.asList(tenantId.getId(), deviceProfile.getName())); } return changed; } @@ -309,7 +339,7 @@ public class DeviceProfileServiceImpl extends AbstractEntityService implements D @Override protected void removeEntity(TenantId tenantId, DeviceProfile entity) { - removeDeviceProfile(tenantId, entity.getId()); + removeDeviceProfile(tenantId, entity); } }; diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java index 7153069849..9e52dc8f09 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java @@ -169,8 +169,14 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe deviceValidator.validate(device, Device::getTenantId); Device savedDevice; try { + DeviceProfile deviceProfile; if (device.getDeviceProfileId() == null) { - DeviceProfile deviceProfile = this.deviceProfileService.findOrCreateDefaultDeviceProfile(device.getTenantId()); + if (!StringUtils.isEmpty(device.getType())) { + deviceProfile = this.deviceProfileService.findOrCreateDeviceProfile(device.getTenantId(), device.getType()); + } else { + deviceProfile = this.deviceProfileService.findDefaultDeviceProfile(device.getTenantId()); + device.setType(deviceProfile.getName()); + } device.setDeviceProfileId(new DeviceProfileId(deviceProfile.getId().getId())); DeviceData deviceData = new DeviceData(); switch (deviceProfile.getType()) { @@ -178,7 +184,7 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe deviceData.setConfiguration(new DefaultDeviceConfiguration()); break; } - switch (deviceProfile.getTransportType()){ + switch (deviceProfile.getTransportType()) { case DEFAULT: deviceData.setTransportConfiguration(new DefaultDeviceTransportConfiguration()); break; @@ -190,7 +196,14 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe break; } device.setDeviceData(deviceData); + } else { + deviceProfile = this.deviceProfileService.findDeviceProfileById(device.getTenantId(), device.getDeviceProfileId()); + if (deviceProfile == null) { + throw new DataValidationException("Device is referencing non existing device profile!"); + } } + device.setType(deviceProfile.getName()); + savedDevice = deviceDao.save(device.getTenantId(), device); } catch (Exception t) { ConstraintViolationException e = extractConstraintViolationException(t).orElse(null); @@ -441,9 +454,6 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe @Override protected void validateDataImpl(TenantId tenantId, Device device) { - if (StringUtils.isEmpty(device.getType())) { - throw new DataValidationException("Device type should be specified!"); - } if (StringUtils.isEmpty(device.getName())) { throw new DataValidationException("Device name should be specified!"); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceProfileRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceProfileRepository.java index 6c11c328d8..caa8d1da84 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceProfileRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceProfileRepository.java @@ -20,6 +20,7 @@ import org.springframework.data.domain.Pageable; import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.PagingAndSortingRepository; import org.springframework.data.repository.query.Param; +import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfileInfo; import org.thingsboard.server.dao.model.sql.DeviceProfileEntity; @@ -53,4 +54,7 @@ public interface DeviceProfileRepository extends PagingAndSortingRepository findByTenantIdAndProfileId(@Param("tenantId") UUID tenantId, + @Param("profileId") UUID profileId, + @Param("searchText") String searchText, + Pageable pageable); + @Query("SELECT new org.thingsboard.server.dao.model.sql.DeviceInfoEntity(d, c.title, c.additionalInfo, p.name) " + "FROM DeviceEntity d " + "LEFT JOIN CustomerEntity c on c.id = d.customerId " + diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java index c68925c177..8c7ee3ae99 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java @@ -104,6 +104,16 @@ public class JpaDeviceDao extends JpaAbstractSearchTextDao DaoUtil.toPageable(pageLink))); } + @Override + public PageData findDevicesByTenantIdAndProfileId(UUID tenantId, UUID profileId, PageLink pageLink) { + return DaoUtil.toPageData( + deviceRepository.findByTenantIdAndProfileId( + tenantId, + profileId, + Objects.toString(pageLink.getTextSearch(), ""), + DaoUtil.toPageable(pageLink))); + } + @Override public PageData findDeviceInfosByTenantIdAndCustomerId(UUID tenantId, UUID customerId, PageLink pageLink) { return DaoUtil.toPageData( diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceProfileDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceProfileDao.java index 3eb911b054..6f5001a2a2 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceProfileDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceProfileDao.java @@ -80,4 +80,8 @@ public class JpaDeviceProfileDao extends JpaAbstractSearchTextDao Date: Fri, 4 Sep 2020 17:34:15 +0300 Subject: [PATCH 3/4] No more device type --- .../server/controller/BaseDeviceControllerTest.java | 5 ++--- .../server/controller/BaseDeviceProfileControllerTest.java | 2 +- .../server/dao/sql/device/DeviceProfileRepository.java | 2 +- .../server/dao/sql/device/JpaDeviceProfileDao.java | 2 +- .../server/dao/service/BaseDeviceServiceTest.java | 6 +++--- 5 files changed, 8 insertions(+), 9 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseDeviceControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseDeviceControllerTest.java index 5db29e7004..72a885b062 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseDeviceControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseDeviceControllerTest.java @@ -186,9 +186,8 @@ public abstract class BaseDeviceControllerTest extends AbstractControllerTest { public void testSaveDeviceWithEmptyType() throws Exception { Device device = new Device(); device.setName("My device"); - doPost("/api/device", device) - .andExpect(status().isBadRequest()) - .andExpect(statusReason(containsString("Device type should be specified"))); + Device savedDevice = doPost("/api/device", device, Device.class); + Assert.assertEquals("default", savedDevice.getType()); } @Test diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseDeviceProfileControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseDeviceProfileControllerTest.java index 8d2c17c25a..b2334d7c46 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseDeviceProfileControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseDeviceProfileControllerTest.java @@ -121,7 +121,7 @@ public abstract class BaseDeviceProfileControllerTest extends AbstractController Assert.assertNotNull(foundDefaultDeviceProfileInfo.getName()); Assert.assertNotNull(foundDefaultDeviceProfileInfo.getType()); Assert.assertEquals(DeviceProfileType.DEFAULT, foundDefaultDeviceProfileInfo.getType()); - Assert.assertEquals("Default", foundDefaultDeviceProfileInfo.getName()); + Assert.assertEquals("default", foundDefaultDeviceProfileInfo.getName()); } @Test diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceProfileRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceProfileRepository.java index caa8d1da84..8116b711d5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceProfileRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceProfileRepository.java @@ -55,6 +55,6 @@ public interface DeviceProfileRepository extends PagingAndSortingRepository devicesTitle2 = new ArrayList<>(); @@ -282,7 +282,7 @@ public abstract class BaseDeviceServiceTest extends AbstractServiceTest { name = i % 2 == 0 ? name.toLowerCase() : name.toUpperCase(); device.setName(name); device.setType("default"); - devicesTitle2.add(new DeviceInfo(deviceService.saveDevice(device), null, false, "Default")); + devicesTitle2.add(new DeviceInfo(deviceService.saveDevice(device), null, false, "default")); } List loadedDevicesTitle1 = new ArrayList<>(); @@ -435,7 +435,7 @@ public abstract class BaseDeviceServiceTest extends AbstractServiceTest { device.setName("Device"+i); device.setType("default"); device = deviceService.saveDevice(device); - devices.add(new DeviceInfo(deviceService.assignDeviceToCustomer(tenantId, device.getId(), customerId), customer.getTitle(), customer.isPublic(), "Default")); + devices.add(new DeviceInfo(deviceService.assignDeviceToCustomer(tenantId, device.getId(), customerId), customer.getTitle(), customer.isPublic(), "default")); } List loadedDevices = new ArrayList<>(); From 8090e5201e8840bfd645dc5858785234ed322d11 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 4 Sep 2020 17:47:47 +0300 Subject: [PATCH 4/4] Gateway Profile update --- .../transport/mqtt/session/GatewayDeviceSessionCtx.java | 9 ++++++++- .../transport/session/DeviceAwareSessionContext.java | 2 +- 2 files changed, 9 insertions(+), 2 deletions(-) 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 c1e975ff92..767eb0102c 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 @@ -16,6 +16,7 @@ package org.thingsboard.server.transport.mqtt.session; import lombok.extern.slf4j.Slf4j; +import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.transport.SessionMsgListener; import org.thingsboard.server.common.transport.auth.TransportDeviceInfo; import org.thingsboard.server.gen.transport.TransportProtos; @@ -32,7 +33,7 @@ import java.util.concurrent.ConcurrentMap; public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext implements SessionMsgListener { private final GatewaySessionHandler parent; - private final SessionInfoProto sessionInfo; + private volatile SessionInfoProto sessionInfo; public GatewayDeviceSessionCtx(GatewaySessionHandler parent, TransportDeviceInfo deviceInfo, ConcurrentMap mqttQoSMap) { super(UUID.randomUUID(), mqttQoSMap); @@ -105,4 +106,10 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse) { // This feature is not supported in the TB IoT Gateway yet. } + + @Override + public void onProfileUpdate(DeviceProfile deviceProfile) { + deviceInfo.setDeviceType(deviceProfile.getName()); + sessionInfo = SessionInfoProto.newBuilder().mergeFrom(sessionInfo).setDeviceType(deviceProfile.getName()).build(); + } } 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 fe26e53b4f..5ca5410156 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 @@ -35,7 +35,7 @@ public abstract class DeviceAwareSessionContext implements SessionContext { @Getter private volatile DeviceId deviceId; @Getter - private volatile TransportDeviceInfo deviceInfo; + protected volatile TransportDeviceInfo deviceInfo; private volatile boolean connected; public DeviceId getDeviceId() {