Browse Source

Merge branch 'develop/3.2' of github.com:thingsboard/thingsboard into develop/3.2

pull/3477/head
Vladyslav_Prykhodko 6 years ago
parent
commit
43596ec5ce
  1. 6
      application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java
  2. 5
      application/src/test/java/org/thingsboard/server/controller/BaseDeviceControllerTest.java
  3. 2
      application/src/test/java/org/thingsboard/server/controller/BaseDeviceProfileControllerTest.java
  4. 6
      common/dao-api/src/main/java/org/thingsboard/server/dao/device/DeviceProfileService.java
  5. 9
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  6. 9
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java
  7. 35
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportProfileCache.java
  8. 77
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportProfileCache.java
  9. 73
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  10. 2
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/session/DeviceAwareSessionContext.java
  11. 11
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java
  12. 2
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileDao.java
  13. 54
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileServiceImpl.java
  14. 20
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java
  15. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceProfileRepository.java
  16. 8
      dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceRepository.java
  17. 10
      dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java
  18. 4
      dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceProfileDao.java
  19. 6
      dao/src/test/java/org/thingsboard/server/dao/service/BaseDeviceServiceTest.java

6
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<EntitySubtype> 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) {
}
}

5
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

2
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

6
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<DeviceProfileInfo> 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);

9
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();
}
}

9
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<MqttTopicMatcher, Integer> 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();
}
}

35
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);
}

77
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<DeviceProfileId, DeviceProfile> 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> 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> deviceProfile = dataDecodingEncodingService.decode(profileBody.toByteArray());
if (deviceProfile.isPresent()) {
put(deviceProfile.get());
return deviceProfile.get();
} else {
return null;
}
}
}

73
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<TransportApiRequestMsg>, TbProtoQueueMsg<TransportApiResponseMsg>> transportApiRequestTemplate;
protected TbQueueProducer<TbProtoQueueMsg<ToRuleEngineMsg>> ruleEngineMsgProducer;
@ -141,7 +141,6 @@ public class DefaultTransportService implements TransportService {
//TODO 3.2: @ybondarenko Implement cleanup of this maps.
private final ConcurrentMap<TenantId, TbRateLimits> perTenantLimits = new ConcurrentHashMap<>();
private final ConcurrentMap<DeviceId, TbRateLimits> perDeviceLimits = new ConcurrentHashMap<>();
private final ConcurrentMap<DeviceProfileId, DeviceProfile> 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> 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> 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> 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<DeviceProfile> 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<DeviceProfile> 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 <T extends com.google.protobuf.GeneratedMessageV3> ListenableFuture<TbProtoQueueMsg<T>> extractProfile(ListenableFuture<TbProtoQueueMsg<T>> send,
Function<T, Boolean> hasDeviceInfo,
Function<T, TransportProtos.DeviceInfoProto> deviceInfoF,
Function<T, ByteString> 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> deviceProfile = dataDecodingEncodingService.decode(profileBody.toByteArray());
deviceProfile.ifPresent(profile -> deviceProfiles.put(deviceProfileId, profile));
}
}
}
return response;
}, transportCallbackExecutor);
}
private class TransportTbQueueCallback implements TbQueueCallback {
private final TransportServiceCallback<Void> callback;

2
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() {

11
dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java

@ -184,4 +184,15 @@ public interface DeviceDao extends Dao<Device> {
ListenableFuture<Device> 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<Device> findDevicesByTenantIdAndProfileId(UUID tenantId, UUID profileId, PageLink pageLink);
}

2
dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileDao.java

@ -37,4 +37,6 @@ public interface DeviceProfileDao extends Dao<DeviceProfile> {
DeviceProfile findDefaultDeviceProfile(TenantId tenantId);
DeviceProfileInfo findDefaultDeviceProfileInfo(TenantId tenantId);
DeviceProfile findByName(TenantId tenantId, String profileName);
}

54
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<Device> 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);
}
};

20
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!");
}

4
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<Devi
"FROM DeviceProfileEntity d " +
"WHERE d.tenantId = :tenantId AND d.isDefault = true")
DeviceProfileInfo findDefaultDeviceProfileInfo(@Param("tenantId") UUID tenantId);
DeviceProfileEntity findByTenantIdAndName(UUID id, String profileName);
}

8
dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceRepository.java

@ -46,6 +46,14 @@ public interface DeviceRepository extends PagingAndSortingRepository<DeviceEntit
@Param("searchText") String searchText,
Pageable pageable);
@Query("SELECT d FROM DeviceEntity d WHERE d.tenantId = :tenantId " +
"AND d.deviceProfileId = :profileId " +
"AND LOWER(d.searchText) LIKE LOWER(CONCAT(:searchText, '%'))")
Page<DeviceEntity> 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 " +

10
dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java

@ -104,6 +104,16 @@ public class JpaDeviceDao extends JpaAbstractSearchTextDao<DeviceEntity, Device>
DaoUtil.toPageable(pageLink)));
}
@Override
public PageData<Device> findDevicesByTenantIdAndProfileId(UUID tenantId, UUID profileId, PageLink pageLink) {
return DaoUtil.toPageData(
deviceRepository.findByTenantIdAndProfileId(
tenantId,
profileId,
Objects.toString(pageLink.getTextSearch(), ""),
DaoUtil.toPageable(pageLink)));
}
@Override
public PageData<DeviceInfo> findDeviceInfosByTenantIdAndCustomerId(UUID tenantId, UUID customerId, PageLink pageLink) {
return DaoUtil.toPageData(

4
dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceProfileDao.java

@ -80,4 +80,8 @@ public class JpaDeviceProfileDao extends JpaAbstractSearchTextDao<DeviceProfileE
return deviceProfileRepository.findDefaultDeviceProfileInfo(tenantId.getId());
}
@Override
public DeviceProfile findByName(TenantId tenantId, String profileName) {
return DaoUtil.getData(deviceProfileRepository.findByTenantIdAndName(tenantId.getId(), profileName));
}
}

6
dao/src/test/java/org/thingsboard/server/dao/service/BaseDeviceServiceTest.java

@ -270,7 +270,7 @@ public abstract class BaseDeviceServiceTest extends AbstractServiceTest {
name = i % 2 == 0 ? name.toLowerCase() : name.toUpperCase();
device.setName(name);
device.setType("default");
devicesTitle1.add(new DeviceInfo(deviceService.saveDevice(device), null, false, "Default"));
devicesTitle1.add(new DeviceInfo(deviceService.saveDevice(device), null, false, "default"));
}
String title2 = "Device title 2";
List<DeviceInfo> 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<DeviceInfo> 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<DeviceInfo> loadedDevices = new ArrayList<>();

Loading…
Cancel
Save