Browse Source

Working version for provision feature with provision credentials on device level

pull/3518/head
zbeacon 6 years ago
parent
commit
8a8695f260
  1. 121
      application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java
  2. 58
      application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java
  3. 2
      application/src/main/resources/thingsboard.yml
  4. 3
      common/data/src/main/java/org/thingsboard/server/common/data/device/data/DeviceConfiguration.java
  5. 2
      common/message/src/main/java/org/thingsboard/server/common/msg/session/FeatureType.java
  6. 4
      common/queue/src/main/proto/queue.proto
  7. 39
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java
  8. 3
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java
  9. 10
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java
  10. 28
      common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java
  11. 7
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  12. 6
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java
  13. 21
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  14. 2
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java
  15. 5
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceProfileDao.java
  16. 29
      dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceProfileRepository.java
  17. 16
      dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceRepository.java
  18. 6
      dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java
  19. 10
      dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceProfileDao.java

121
application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java

@ -34,6 +34,7 @@ import org.thingsboard.server.common.data.device.profile.ProvisionDeviceProfileC
import org.thingsboard.server.common.data.device.profile.ProvisionRequestValidationStrategyType;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry;
@ -60,6 +61,7 @@ import org.thingsboard.server.queue.TbQueueCallback;
import org.thingsboard.server.queue.TbQueueProducer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
import org.thingsboard.server.service.state.DeviceStateService;
import java.util.Collections;
@ -103,55 +105,88 @@ public class DeviceProvisionServiceImpl implements DeviceProvisionService {
@Autowired
PartitionService partitionService;
public DeviceProvisionServiceImpl(TbQueueProducerProvider producerProvider) {
ruleEngineMsgProducer = producerProvider.getRuleEngineMsgProducer();
}
@Override
public ListenableFuture<ProvisionResponse> provisionDevice(ProvisionRequest provisionRequest) {
Device targetDevice = deviceDao.findDeviceByTenantIdAndDeviceDataProvisionConfigurationPair(
TenantId.SYS_TENANT_ID,
provisionRequest.getCredentials().getProvisionDeviceKey(),
provisionRequest.getCredentials().getProvisionDeviceSecret());
String provisionRequestKey = provisionRequest.getCredentials().getProvisionDeviceKey();
String provisionRequestSecret = provisionRequest.getCredentials().getProvisionDeviceSecret();
if (targetDevice != null) {
if (targetDevice.getDeviceData().getConfiguration().getType() != DeviceProfileType.PROVISION) {
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.NOT_FOUND));
}
if (StringUtils.isEmpty(provisionRequestKey) || StringUtils.isEmpty(provisionRequestSecret)) {
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.NOT_FOUND));
}
DeviceProfile targetProfile = deviceProfileDao.findById(TenantId.SYS_TENANT_ID, targetDevice.getDeviceProfileId().getId());
Device targetDevice = deviceDao.findDeviceByProfileNameAndDeviceDataProvisionConfigurationPair(
provisionRequest.getDeviceType(),
provisionRequestKey,
provisionRequestSecret
).orElse(null);
ProvisionDeviceConfiguration currentProfileConfiguration = (ProvisionDeviceConfiguration) targetDevice.getDeviceData().getConfiguration();
if (!new ProvisionDeviceConfiguration(provisionRequest.getCredentials().getProvisionDeviceKey(), provisionRequest.getCredentials().getProvisionDeviceSecret()).equals(currentProfileConfiguration)) {
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.NOT_FOUND));
}
ProvisionRequestValidationStrategyType targetStrategy = getStrategy(targetProfile);
switch (targetStrategy) {
case CHECK_NEW_DEVICE:
log.warn("[{}] The device is present and could not be provisioned once more!", targetDevice.getName());
notify(targetDevice, provisionRequest, DataConstants.PROVISION_FAILURE, false);
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.FAILURE));
case CHECK_PRE_PROVISIONED_DEVICE:
return processProvision(targetDevice, provisionRequest);
default:
throw new RuntimeException("Strategy is not supported - " + targetStrategy.name());
}
if (targetDevice != null) {
return processProvisionDeviceWithKeySecretPairExists(provisionRequest, provisionRequestKey, provisionRequestSecret, targetDevice);
} else {
DeviceProfile targetProfile = deviceProfileDao.findProfileByTenantIdAndProfileDataProvisionConfigurationPair(
TenantId.SYS_TENANT_ID,
provisionRequest.getCredentials().getProvisionDeviceKey(),
provisionRequest.getCredentials().getProvisionDeviceSecret()
);
if (targetProfile.getProfileData().getConfiguration().getType() != DeviceProfileType.PROVISION) {
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.NOT_FOUND));
}
ProvisionRequestValidationStrategyType targetStrategy = getStrategy(targetProfile);
switch (targetStrategy) {
case CHECK_NEW_DEVICE:
return createDevice(provisionRequest, targetProfile);
case CHECK_PRE_PROVISIONED_DEVICE:
log.warn("[{}] Failed to find pre provisioned device!", provisionRequest.getDeviceName());
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.FAILURE));
default:
throw new RuntimeException("Strategy is not supported - " + targetStrategy.name());
}
return processProvisionDeviceWithKeySecretPairNotExists(provisionRequest, provisionRequestKey, provisionRequestSecret);
}
}
private ListenableFuture<ProvisionResponse> processProvisionDeviceWithKeySecretPairExists(ProvisionRequest provisionRequest, String provisionRequestKey, String provisionRequestSecret, Device targetDevice) {
if (targetDevice.getDeviceData().getConfiguration().getType() != DeviceProfileType.PROVISION) {
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.NOT_FOUND));
}
DeviceProfile targetProfile = deviceProfileDao.findById(targetDevice.getTenantId(), targetDevice.getDeviceProfileId().getId());
if (targetProfile == null || targetProfile.getProfileData().getConfiguration().getType() != DeviceProfileType.PROVISION) {
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.NOT_FOUND));
}
ProvisionDeviceConfiguration currentDeviceConfiguration = (ProvisionDeviceConfiguration) targetDevice.getDeviceData().getConfiguration();
if (!new ProvisionDeviceConfiguration(provisionRequestKey, provisionRequestSecret).equals(currentDeviceConfiguration)) {
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.NOT_FOUND));
}
ProvisionRequestValidationStrategyType targetStrategy = getStrategy(targetProfile);
switch (targetStrategy) {
case CHECK_NEW_DEVICE:
log.warn("[{}] The device is present and could not be provisioned once more!", targetDevice.getName());
notify(targetDevice, provisionRequest, DataConstants.PROVISION_FAILURE, false);
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.FAILURE));
case CHECK_PRE_PROVISIONED_DEVICE:
return processProvision(targetDevice, provisionRequest);
default:
throw new RuntimeException("Strategy is not supported - " + targetStrategy.name());
}
}
private ListenableFuture<ProvisionResponse> processProvisionDeviceWithKeySecretPairNotExists(ProvisionRequest provisionRequest, String provisionRequestKey, String provisionRequestSecret){
DeviceProfile targetProfile = deviceProfileDao.findProfileByProfileNameAndProfileDataProvisionConfigurationPair(
provisionRequest.getDeviceType(),
provisionRequestKey,
provisionRequestSecret
);
if (targetProfile == null || targetProfile.getProfileData().getConfiguration().getType() != DeviceProfileType.PROVISION) {
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.NOT_FOUND));
}
ProvisionRequestValidationStrategyType targetStrategy = getStrategy(targetProfile);
switch (targetStrategy) {
case CHECK_NEW_DEVICE:
return createDevice(provisionRequest, targetProfile);
case CHECK_PRE_PROVISIONED_DEVICE:
ProvisionDeviceProfileConfiguration currentDeviceProfileConfiguration = (ProvisionDeviceProfileConfiguration) targetProfile.getProfileData().getConfiguration();
if(new ProvisionDeviceProfileConfiguration(provisionRequestKey, provisionRequestSecret).equals(currentDeviceProfileConfiguration)) {
Optional<Device> optionalDevice = deviceDao.findDeviceByTenantIdAndName(targetProfile.getTenantId().getId(), provisionRequest.getDeviceName());
if (optionalDevice.isPresent()) {
return processProvision(optionalDevice.get(), provisionRequest);
}
}
log.warn("[{}] Failed to find pre provisioned device!", provisionRequest.getDeviceName());
return Futures.immediateFuture(new ProvisionResponse(null, ProvisionResponseStatus.FAILURE));
default:
throw new RuntimeException("Strategy is not supported - " + targetStrategy.name());
}
}
@ -281,6 +316,6 @@ public class DeviceProvisionServiceImpl implements DeviceProvisionService {
private void logAction(TenantId tenantId, CustomerId customerId, Device device, boolean success, ProvisionRequest provisionRequest) {
ActionType actionType = success ? ActionType.PROVISION_SUCCESS : ActionType.PROVISION_FAILURE;
auditLogService.logEntityAction(tenantId, customerId, null, device.getName(), device.getId(), device, actionType, null, provisionRequest);
auditLogService.logEntityAction(tenantId, customerId, new UserId(UserId.NULL_UUID), device.getName(), device.getId(), device, actionType, null, provisionRequest);
}
}

58
application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java

@ -22,6 +22,7 @@ import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.protobuf.ByteString;
import com.google.protobuf.InvalidProtocolBufferException;
import lombok.extern.slf4j.Slf4j;
import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder;
import org.springframework.stereotype.Service;
@ -31,6 +32,7 @@ import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials;
import org.thingsboard.server.common.data.device.data.ProvisionDeviceConfiguration;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.DeviceProfileId;
@ -45,17 +47,24 @@ import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.dao.device.DeviceCredentialsService;
import org.thingsboard.server.dao.device.DeviceProfileService;
import org.thingsboard.server.dao.device.DeviceProvisionService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.device.provision.ProvisionRequest;
import org.thingsboard.server.dao.device.provision.ProvisionResponse;
import org.thingsboard.server.dao.device.provision.ProvisionResponseStatus;
import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.dao.tenant.TenantProfileService;
import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.dao.util.mapping.JacksonUtil;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.CredentialsType;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceInfoProto;
import org.thingsboard.server.gen.transport.TransportProtos.GetOrCreateDeviceFromGatewayRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetOrCreateDeviceFromGatewayResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetTenantRoutingInfoRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetTenantRoutingInfoResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceCredentialsMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TransportApiResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg;
@ -94,6 +103,7 @@ public class DefaultTransportApiService implements TransportApiService {
private final DbCallbackExecutorService dbCallbackExecutorService;
private final TbClusterService tbClusterService;
private final DataDecodingEncodingService dataDecodingEncodingService;
private final DeviceProvisionService deviceProvisionService;
private final ConcurrentMap<String, ReentrantLock> deviceCreationLocks = new ConcurrentHashMap<>();
@ -102,7 +112,8 @@ public class DefaultTransportApiService implements TransportApiService {
TenantProfileService tenantProfileService, DeviceService deviceService,
RelationService relationService, DeviceCredentialsService deviceCredentialsService,
DeviceStateService deviceStateService, DbCallbackExecutorService dbCallbackExecutorService,
TbClusterService tbClusterService, DataDecodingEncodingService dataDecodingEncodingService) {
TbClusterService tbClusterService, DataDecodingEncodingService dataDecodingEncodingService,
DeviceProvisionService deviceProvisionService) {
this.deviceProfileService = deviceProfileService;
this.tenantService = tenantService;
this.tenantProfileService = tenantProfileService;
@ -113,6 +124,7 @@ public class DefaultTransportApiService implements TransportApiService {
this.dbCallbackExecutorService = dbCallbackExecutorService;
this.tbClusterService = tbClusterService;
this.dataDecodingEncodingService = dataDecodingEncodingService;
this.deviceProvisionService = deviceProvisionService;
}
@Override
@ -139,6 +151,9 @@ public class DefaultTransportApiService implements TransportApiService {
} else if (transportApiRequestMsg.hasGetDeviceProfileRequestMsg()) {
return Futures.transform(handle(transportApiRequestMsg.getGetDeviceProfileRequestMsg()),
value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
} else if (transportApiRequestMsg.hasProvisionDeviceRequestMsg()) {
return Futures.transform(handle(transportApiRequestMsg.getProvisionDeviceRequestMsg()),
value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
}
return Futures.transform(getEmptyTransportApiResponseFuture(),
value -> new TbProtoQueueMsg<>(tbProtoQueueMsg.getKey(), value, tbProtoQueueMsg.getHeaders()), MoreExecutors.directExecutor());
@ -261,6 +276,47 @@ public class DefaultTransportApiService implements TransportApiService {
}, dbCallbackExecutorService);
}
private ListenableFuture<TransportApiResponseMsg> handle(ProvisionDeviceRequestMsg requestMsg) {
ListenableFuture<ProvisionResponse> provisionResponseFuture = null;
provisionResponseFuture = deviceProvisionService.provisionDevice(
new ProvisionRequest(
requestMsg.getDeviceName(),
requestMsg.getDeviceType(),
requestMsg.getX509CertPubKey(),
new ProvisionDeviceConfiguration(
requestMsg.getProvisionDeviceCredentialsMsg().getProvisionDeviceKey(),
requestMsg.getProvisionDeviceCredentialsMsg().getProvisionDeviceSecret())));
return Futures.transform(provisionResponseFuture, provisionResponse -> {
if (provisionResponse.getResponseStatus() == ProvisionResponseStatus.NOT_FOUND) {
return getTransportApiResponseMsg(TransportProtos.DeviceCredentialsProto.getDefaultInstance(), TransportProtos.ProvisionResponseStatus.NOT_FOUND);
} else if (provisionResponse.getResponseStatus() == ProvisionResponseStatus.FAILURE) {
return getTransportApiResponseMsg(TransportProtos.DeviceCredentialsProto.getDefaultInstance(), TransportProtos.ProvisionResponseStatus.FAILURE);
} else {
return getTransportApiResponseMsg(getDeviceCredentials(provisionResponse.getDeviceCredentials()), TransportProtos.ProvisionResponseStatus.SUCCESS);
}
}, dbCallbackExecutorService);
}
private TransportApiResponseMsg getTransportApiResponseMsg(TransportProtos.DeviceCredentialsProto deviceCredentials, TransportProtos.ProvisionResponseStatus status) {
return TransportApiResponseMsg.newBuilder()
.setProvisionDeviceResponseMsg(TransportProtos.ProvisionDeviceResponseMsg.newBuilder()
.setDeviceCredentials(deviceCredentials)
.setProvisionResponseStatus(status)
.build())
.build();
}
private TransportProtos.DeviceCredentialsProto getDeviceCredentials(DeviceCredentials deviceCredentials) {
return TransportProtos.DeviceCredentialsProto.newBuilder()
.setDeviceIdMSB(deviceCredentials.getDeviceId().getId().getMostSignificantBits())
.setDeviceIdLSB(deviceCredentials.getDeviceId().getId().getLeastSignificantBits())
.setCredentialsType(deviceCredentials.getCredentialsType() == DeviceCredentialsType.ACCESS_TOKEN ?
CredentialsType.ACCESS_TOKEN : CredentialsType.X509_CERTIFICATE)
.setCredentialsId(deviceCredentials.getCredentialsId())
.setCredentialsValue(deviceCredentials.getCredentialsValue() != null ? deviceCredentials.getCredentialsValue() : "")
.build();
}
private ListenableFuture<TransportApiResponseMsg> handle(GetTenantRoutingInfoRequestMsg requestMsg) {
TenantId tenantId = new TenantId(new UUID(requestMsg.getTenantIdMSB(), requestMsg.getTenantIdLSB()));
// TODO: Tenant Profile from cache

2
application/src/main/resources/thingsboard.yml

@ -473,7 +473,9 @@ spring:
enabled: "true"
jpa:
open-in-view: "false"
show-sql: "true"
hibernate:
format_sql: "true"
ddl-auto: "none"
database-platform: "${SPRING_JPA_DATABASE_PLATFORM:org.hibernate.dialect.PostgreSQLDialect}"
datasource:

3
common/data/src/main/java/org/thingsboard/server/common/data/device/data/DeviceConfiguration.java

@ -27,7 +27,8 @@ import org.thingsboard.server.common.data.DeviceProfileType;
include = JsonTypeInfo.As.PROPERTY,
property = "type")
@JsonSubTypes({
@JsonSubTypes.Type(value = DefaultDeviceConfiguration.class, name = "DEFAULT")})
@JsonSubTypes.Type(value = DefaultDeviceConfiguration.class, name = "DEFAULT"),
@JsonSubTypes.Type(value = ProvisionDeviceConfiguration.class, name = "PROVISION")})
public interface DeviceConfiguration {
@JsonIgnore

2
common/message/src/main/java/org/thingsboard/server/common/msg/session/FeatureType.java

@ -16,5 +16,5 @@
package org.thingsboard.server.common.msg.session;
public enum FeatureType {
ATTRIBUTES, TELEMETRY, RPC, CLAIM
ATTRIBUTES, TELEMETRY, RPC, CLAIM, PROVISION
}

4
common/queue/src/main/proto/queue.proto

@ -477,7 +477,7 @@ message TransportApiRequestMsg {
GetTenantRoutingInfoRequestMsg getTenantRoutingInfoRequestMsg = 4;
GetDeviceProfileRequestMsg getDeviceProfileRequestMsg = 5;
ValidateBasicMqttCredRequestMsg validateBasicMqttCredRequestMsg = 6;
// ProvisionDeviceRequestMsg provisionDeviceRequestMsg = 7;
ProvisionDeviceRequestMsg provisionDeviceRequestMsg = 7;
}
/* Response from ThingsBoard Core Service to Transport Service */
@ -486,7 +486,7 @@ message TransportApiResponseMsg {
GetOrCreateDeviceFromGatewayResponseMsg getOrCreateDeviceResponseMsg = 2;
GetTenantRoutingInfoResponseMsg getTenantRoutingInfoResponseMsg = 4;
GetDeviceProfileResponseMsg getDeviceProfileResponseMsg = 5;
// ProvisionDeviceResponseMsg provisionDeviceResponseMsg = 6;
ProvisionDeviceResponseMsg provisionDeviceResponseMsg = 6;
}
/* Messages that are handled by ThingsBoard Core Service */

39
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java

@ -24,6 +24,7 @@ import org.eclipse.californium.core.network.ExchangeObserver;
import org.eclipse.californium.core.server.resources.CoapExchange;
import org.eclipse.californium.core.server.resources.Resource;
import org.springframework.util.ReflectionUtils;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.DeviceTransportType;
import org.thingsboard.server.common.data.security.DeviceTokenCredentials;
import org.thingsboard.server.common.msg.session.FeatureType;
@ -33,9 +34,11 @@ import org.thingsboard.server.common.transport.TransportContext;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.common.transport.adaptor.JsonConverter;
import org.thingsboard.server.common.transport.auth.SessionInfoCreator;
import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceResponseMsg;
import java.lang.reflect.Field;
import java.util.List;
@ -130,10 +133,25 @@ public class CoapTransportResource extends CoapResource {
case CLAIM:
processRequest(exchange, SessionMsgType.CLAIM_REQUEST);
break;
case PROVISION:
processProvision(exchange);
break;
}
}
}
private void processProvision(CoapExchange exchange) {
log.trace("Processing {}", exchange.advanced().getRequest());
exchange.accept();
try {
transportService.process(transportContext.getAdaptor().convertToProvisionRequestMsg(UUID.randomUUID(), exchange.advanced().getRequest()),
new DeviceProvisionCallback(exchange));
} catch (AdaptorException e) {
log.trace("Failed to decode message: ", e);
exchange.respond(ResponseCode.BAD_REQUEST);
}
}
private void processRequest(CoapExchange exchange, SessionMsgType type) {
log.trace("Processing {}", exchange.advanced().getRequest());
exchange.accept();
@ -274,6 +292,8 @@ public class CoapTransportResource extends CoapResource {
try {
if (uriPath.size() >= FEATURE_TYPE_POSITION) {
return Optional.of(FeatureType.valueOf(uriPath.get(FEATURE_TYPE_POSITION - 1).toUpperCase()));
} else if (uriPath.size() == 3 && uriPath.contains(DataConstants.PROVISION)) {
return Optional.of(FeatureType.valueOf(DataConstants.PROVISION.toUpperCase()));
}
} catch (RuntimeException e) {
log.warn("Failed to decode feature type: {}", uriPath);
@ -325,6 +345,25 @@ public class CoapTransportResource extends CoapResource {
}
}
private static class DeviceProvisionCallback implements TransportServiceCallback<ProvisionDeviceResponseMsg> {
private final CoapExchange exchange;
DeviceProvisionCallback(CoapExchange exchange) {
this.exchange = exchange;
}
@Override
public void onSuccess(TransportProtos.ProvisionDeviceResponseMsg msg) {
exchange.respond(JsonConverter.toJson(msg).toString());
}
@Override
public void onError(Throwable e) {
log.warn("Failed to process request", e);
exchange.respond(ResponseCode.INTERNAL_SERVER_ERROR);
}
}
private static class CoapOkCallback implements TransportServiceCallback<Void> {
private final CoapExchange exchange;

3
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java

@ -19,6 +19,7 @@ import org.eclipse.californium.core.coap.Request;
import org.eclipse.californium.core.coap.Response;
import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg;
import org.thingsboard.server.transport.coap.CoapTransportResource;
import java.util.UUID;
@ -45,4 +46,6 @@ public interface CoapTransportAdaptor {
Response convertToPublish(CoapTransportResource.CoapSessionListener coapSessionListener, TransportProtos.ToServerRpcResponseMsg msg) throws AdaptorException;
ProvisionDeviceRequestMsg convertToProvisionRequestMsg(UUID sessionId, Request inbound) throws AdaptorException;
}

10
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java

@ -123,6 +123,16 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor {
return response;
}
@Override
public TransportProtos.ProvisionDeviceRequestMsg convertToProvisionRequestMsg(UUID sessionId, Request inbound) throws AdaptorException {
String payload = validatePayload(sessionId, inbound, false);
try {
return JsonConverter.convertToProvisionRequestMsg(payload);
} catch (IllegalStateException | JsonSyntaxException ex) {
throw new AdaptorException(ex);
}
}
@Override
public Response convertToPublish(CoapTransportResource.CoapSessionListener session, TransportProtos.GetAttributeResponseMsg msg) throws AdaptorException {
if (msg.getClientAttributeListCount() == 0 && msg.getSharedAttributeListCount() == 0) {

28
common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java

@ -44,6 +44,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotif
import org.thingsboard.server.gen.transport.TransportProtos.DeviceInfoProto;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.SessionCloseNotificationProto;
import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto;
import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToAttributeUpdatesMsg;
@ -203,6 +204,14 @@ public class DeviceApiController {
return responseWriter;
}
@RequestMapping(value = "/provision", method = RequestMethod.POST)
public DeferredResult<ResponseEntity> provisionDevice(@RequestBody String json, HttpServletRequest httpRequest) {
DeferredResult<ResponseEntity> responseWriter = new DeferredResult<>();
transportContext.getTransportService().process(JsonConverter.convertToProvisionRequestMsg(json),
new DeviceProvisionCallback(responseWriter));
return responseWriter;
}
private static class DeviceAuthCallback implements TransportServiceCallback<ValidateDeviceCredentialsResponse> {
private final TransportContext transportContext;
private final DeferredResult<ResponseEntity> responseWriter;
@ -230,6 +239,25 @@ public class DeviceApiController {
}
}
private static class DeviceProvisionCallback implements TransportServiceCallback<ProvisionDeviceResponseMsg> {
private final DeferredResult<ResponseEntity> responseWriter;
DeviceProvisionCallback(DeferredResult<ResponseEntity> responseWriter) {
this.responseWriter = responseWriter;
}
@Override
public void onSuccess(ProvisionDeviceResponseMsg msg) {
responseWriter.setResult(new ResponseEntity<>(JsonConverter.toJson(msg).toString(), HttpStatus.OK));
}
@Override
public void onError(Throwable e) {
log.warn("Failed to process request", e);
responseWriter.setResult(new ResponseEntity<>(HttpStatus.INTERNAL_SERVER_ERROR));
}
}
private static class SessionCloseOnErrorCallback implements TransportServiceCallback<Void> {
private final TransportService transportService;
private final SessionInfoProto sessionInfo;

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

@ -154,7 +154,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
try {
if (topicName.equals(MqttTopics.DEVICE_PROVISION_REQUEST_TOPIC)) {
TransportProtos.ProvisionDeviceRequestMsg provisionRequestMsg = adaptor.convertToProvisionRequestMsg(deviceSessionCtx, mqttMsg);
transportService.process(deviceSessionCtx.getSessionInfo(), provisionRequestMsg, (TransportServiceCallback) new DeviceProvisionCallback(ctx, msgId, provisionRequestMsg));
transportService.process(provisionRequestMsg, new DeviceProvisionCallback(ctx, msgId, provisionRequestMsg));
log.trace("[{}][{}] Processing publish msg [{}][{}]!", sessionId, deviceSessionCtx.getDeviceId(), topicName, msgId);
} else {
throw new RuntimeException("Unsupported topic for provisioning requests!");
@ -167,6 +167,10 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
case PINGREQ:
ctx.writeAndFlush(new MqttMessage(new MqttFixedHeader(PINGRESP, false, AT_MOST_ONCE, false, 0)));
break;
// case SUBSCRIBE:
// deviceSessionCtx.setDeviceInfo(TransportDeviceInfo);
// processSubscribe(ctx, (MqttSubscribeMessage) msg);
// break;
case DISCONNECT:
ctx.close();
break;
@ -425,7 +429,6 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
log.info("[{}] Processing connect msg for client: {}!", sessionId, msg.payload().clientIdentifier());
String userName = msg.payload().userName();
if (DataConstants.PROVISION.equals(userName)) {
deviceSessionCtx.setDeviceInfo(new TransportDeviceInfo());
deviceSessionCtx.setProvisionOnly(true);
ctx.writeAndFlush(createMqttConnAckMsg(CONNECTION_ACCEPTED));
} else {

6
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java

@ -28,6 +28,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.GetTenantRoutingInfo
import org.thingsboard.server.gen.transport.TransportProtos.PostAttributeMsg;
import org.thingsboard.server.gen.transport.TransportProtos.PostTelemetryMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.SessionEventMsg;
import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto;
import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToAttributeUpdatesMsg;
@ -58,6 +59,9 @@ public interface TransportService {
void process(GetOrCreateDeviceFromGatewayRequestMsg msg,
TransportServiceCallback<GetOrCreateDeviceFromGatewayResponse> callback);
void process(ProvisionDeviceRequestMsg msg,
TransportServiceCallback<ProvisionDeviceResponseMsg> callback);
void getDeviceProfile(DeviceProfileId deviceProfileId, TransportServiceCallback<DeviceProfile> callback);
void onProfileUpdate(DeviceProfile deviceProfile);
@ -84,8 +88,6 @@ public interface TransportService {
void process(SessionInfoProto sessionInfo, ClaimDeviceMsg msg, TransportServiceCallback<Void> callback);
void process(SessionInfoProto sessionInfo, ProvisionDeviceRequestMsg msg, TransportServiceCallback<Void> deviceProvisionCallback);
void registerAsyncSession(SessionInfoProto sessionInfo, SessionMsgListener listener);
void registerSyncSession(SessionInfoProto sessionInfo, SessionMsgListener listener, long timeout);

21
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

@ -52,6 +52,8 @@ import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsRes
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.common.transport.util.JsonUtils;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ProvisionDeviceResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
@ -332,6 +334,16 @@ public class DefaultTransportService implements TransportService {
return tdi;
}
@Override
public void process(ProvisionDeviceRequestMsg requestMsg, TransportServiceCallback<ProvisionDeviceResponseMsg> callback) {
log.trace("Processing msg: {}", requestMsg);
TbProtoQueueMsg<TransportApiRequestMsg> protoMsg = new TbProtoQueueMsg<>(UUID.randomUUID(), TransportApiRequestMsg.newBuilder().setProvisionDeviceRequestMsg(requestMsg).build());
ListenableFuture<ProvisionDeviceResponseMsg> response = Futures.transform(transportApiRequestTemplate.send(protoMsg), tmp ->
tmp.getValue().getProvisionDeviceResponseMsg()
, MoreExecutors.directExecutor());
AsyncCallbackTemplate.withCallback(response, callback::onSuccess, callback::onError, transportCallbackExecutor);
}
@Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.SubscriptionInfoProto msg, TransportServiceCallback<Void> callback) {
if (log.isTraceEnabled()) {
@ -484,15 +496,6 @@ public class DefaultTransportService implements TransportService {
}
}
@Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ProvisionDeviceRequestMsg msg, TransportServiceCallback<Void> callback) {
if (checkLimits(sessionInfo, msg, callback)) {
reportActivityInternal(sessionInfo);
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo)
.setProvisionDevice(msg).build(), callback);
}
}
@Override
public void reportActivity(TransportProtos.SessionInfoProto sessionInfo) {
reportActivityInternal(sessionInfo);

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

@ -215,6 +215,6 @@ public interface DeviceDao extends Dao<Device> {
*/
PageData<Device> findDevicesByTenantIdAndProfileId(UUID tenantId, UUID profileId, PageLink pageLink);
Device findDeviceByTenantIdAndDeviceDataProvisionConfigurationPair(TenantId tenantId, String provisionDeviceKey, String provisionDeviceSecret);
Optional<Device> findDeviceByProfileNameAndDeviceDataProvisionConfigurationPair(String profileName, String provisionDeviceKey, String provisionDeviceSecret);
}

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

@ -22,6 +22,7 @@ import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.Dao;
import java.util.Optional;
import java.util.UUID;
public interface DeviceProfileDao extends Dao<DeviceProfile> {
@ -38,9 +39,7 @@ public interface DeviceProfileDao extends Dao<DeviceProfile> {
DeviceProfileInfo findDefaultDeviceProfileInfo(TenantId tenantId);
DeviceProfileInfo findProfileInfoByTenantIdAndProfileDataProvisionConfigurationPair(TenantId tenantId, String provisionDeviceKey, String provisionDeviceSecret);
DeviceProfile findProfileByTenantIdAndProfileDataProvisionConfigurationPair(TenantId tenantId, String provisionDeviceKey, String provisionDeviceSecret);
DeviceProfile findProfileByProfileNameAndProfileDataProvisionConfigurationPair(String profileName, String provisionDeviceKey, String provisionDeviceSecret);
DeviceProfile findByName(TenantId tenantId, String profileName);
}

29
dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceProfileRepository.java

@ -20,9 +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.common.data.DeviceProfileType;
import org.thingsboard.server.dao.model.sql.DeviceProfileEntity;
import java.util.UUID;
@ -58,23 +56,14 @@ public interface DeviceProfileRepository extends PagingAndSortingRepository<Devi
DeviceProfileEntity findByTenantIdAndName(UUID id, String profileName);
@Query(value = "SELECT d FROM DeviceProfileEntity d " +
"WHERE d.tenantId = :tenantId " +
"AND d.profileData::jsonb->>{'configuration', 'provisionDeviceKey'} = :provisionDeviceKey " +
"AND d.profileData::jsonb->>{'configuration', 'provisionDeviceSecret' = :provisionDeviceSecret}",
@Query(value = "SELECT d.* FROM device_profile as d " +
"WHERE d.name = :profileName " +
"AND d.profile_data->'configuration'->>'provisionDeviceKey' IS NOT NULL " +
"AND d.profile_data->'configuration'->>'provisionDeviceSecret' IS NOT NULL " +
"AND d.profile_data->'configuration'->>'provisionDeviceKey' = :provisionDeviceKey " +
"AND d.profile_data->'configuration'->>'provisionDeviceSecret' = :provisionDeviceSecret",
nativeQuery = true)
DeviceProfileEntity findProfileByTenantIdAndProfileDataProvisionConfigurationPair(@Param("tenantId") UUID tenantId,
@Param("provisionDeviceKey") String provisionDeviceKey,
@Param("provisionDeviceSecret") String provisionDeviceSecret);
@Query(value = "SELECT new org.thingsboard.server.common.data.DeviceProfileInfo(d.id, d.name, d.type, d.transportType) " +
" FROM DeviceProfileEntity d " +
"WHERE d.tenantId = :tenantId " +
"AND d.profileData::jsonb->>{'configuration', 'provisionDeviceKey'} = :provisionDeviceKey " +
"AND d.profileData::jsonb->>{'configuration', 'provisionDeviceSecret' = :provisionDeviceSecret}",
nativeQuery = true)
DeviceProfileInfo findProfileInfoByTenantIdAndProfileDataProvisionConfigurationPair(@Param("tenantId") UUID tenantId,
@Param("provisionDeviceKey") String provisionDeviceKey,
@Param("provisionDeviceSecret") String provisionDeviceSecret);
DeviceProfileEntity findProfileByProfileNameAndProfileDataProvisionConfigurationPair(@Param("profileName") String profileName,
@Param("provisionDeviceKey") String provisionDeviceKey,
@Param("provisionDeviceSecret") String provisionDeviceSecret);
}

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

@ -20,10 +20,8 @@ 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.DeviceProfileInfo;
import org.thingsboard.server.dao.model.sql.DeviceEntity;
import org.thingsboard.server.dao.model.sql.DeviceInfoEntity;
import org.thingsboard.server.dao.model.sql.DeviceProfileEntity;
import java.util.List;
import java.util.UUID;
@ -171,12 +169,12 @@ public interface DeviceRepository extends PagingAndSortingRepository<DeviceEntit
Long countByDeviceProfileId(UUID deviceProfileId);
@Query(value = "SELECT d FROM Device d " +
"WHERE d.tenant_id = :tenantId " +
"AND d.device_data::jsonb->>('configuration', 'provisionDeviceKey') = :provisionDeviceKey " +
"AND d.device_data::jsonb->>('configuration', 'provisionDeviceSecret') = :provisionDeviceSecret",
@Query(value = "SELECT * FROM Device as d " +
"WHERE d.device_data->'configuration'->>'provisionDeviceKey' = :provisionDeviceKey " +
"AND d.device_data->'configuration'->>'provisionDeviceSecret' = :provisionDeviceSecret " +
"AND d.type = :profileName",
nativeQuery = true)
DeviceEntity findDeviceByTenantIdAndDeviceDataProvisionConfigurationPair(@Param("tenantId") UUID tenantId,
@Param("provisionDeviceKey") String provisionDeviceKey,
@Param("provisionDeviceSecret") String provisionDeviceSecret);
DeviceEntity findDeviceByProfileNameAndDeviceDataProvisionConfigurationPair(@Param("profileName") String profileName,
@Param("provisionDeviceKey") String provisionDeviceKey,
@Param("provisionDeviceSecret") String provisionDeviceSecret);
}

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

@ -22,8 +22,6 @@ import org.springframework.stereotype.Component;
import org.springframework.util.StringUtils;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceInfo;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.DeviceProfileInfo;
import org.thingsboard.server.common.data.EntitySubtype;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.TenantId;
@ -222,8 +220,8 @@ public class JpaDeviceDao extends JpaAbstractSearchTextDao<DeviceEntity, Device>
}
@Override
public Device findDeviceByTenantIdAndDeviceDataProvisionConfigurationPair(TenantId tenantId, String provisionDeviceKey, String provisionDeviceSecret) {
return DaoUtil.getData(deviceRepository.findDeviceByTenantIdAndDeviceDataProvisionConfigurationPair(tenantId.getId(), provisionDeviceKey, provisionDeviceSecret));
public Optional<Device> findDeviceByProfileNameAndDeviceDataProvisionConfigurationPair(String profileName, String provisionDeviceKey, String provisionDeviceSecret) {
return Optional.ofNullable(DaoUtil.getData(deviceRepository.findDeviceByProfileNameAndDeviceDataProvisionConfigurationPair(profileName, provisionDeviceKey, provisionDeviceSecret)));
}
private List<EntitySubtype> convertTenantDeviceTypesToDto(UUID tenantId, List<String> types) {

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

@ -29,6 +29,7 @@ import org.thingsboard.server.dao.model.sql.DeviceProfileEntity;
import org.thingsboard.server.dao.sql.JpaAbstractSearchTextDao;
import java.util.Objects;
import java.util.Optional;
import java.util.UUID;
@Component
@ -81,13 +82,8 @@ public class JpaDeviceProfileDao extends JpaAbstractSearchTextDao<DeviceProfileE
}
@Override
public DeviceProfile findProfileByTenantIdAndProfileDataProvisionConfigurationPair(TenantId tenantId, String provisionDeviceKey, String provisionDeviceSecret) {
return DaoUtil.getData(deviceProfileRepository.findProfileByTenantIdAndProfileDataProvisionConfigurationPair(tenantId.getId(), provisionDeviceKey, provisionDeviceSecret));
}
@Override
public DeviceProfileInfo findProfileInfoByTenantIdAndProfileDataProvisionConfigurationPair(TenantId tenantId, String provisionDeviceKey, String provisionDeviceSecret) {
return deviceProfileRepository.findProfileInfoByTenantIdAndProfileDataProvisionConfigurationPair(tenantId.getId(), provisionDeviceKey, provisionDeviceSecret);
public DeviceProfile findProfileByProfileNameAndProfileDataProvisionConfigurationPair(String profileName, String provisionDeviceKey, String provisionDeviceSecret) {
return DaoUtil.getData(deviceProfileRepository.findProfileByProfileNameAndProfileDataProvisionConfigurationPair(profileName, provisionDeviceKey, provisionDeviceSecret));
}
@Override

Loading…
Cancel
Save