diff --git a/application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java b/application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java index 739af86c13..2a32e86f11 100644 --- a/application/src/main/java/org/thingsboard/server/service/device/DeviceProvisionServiceImpl.java +++ b/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 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 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 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 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); } } diff --git a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java b/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java index 16c7c894da..fb3ba8a4b4 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java +++ b/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 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 handle(ProvisionDeviceRequestMsg requestMsg) { + ListenableFuture 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 handle(GetTenantRoutingInfoRequestMsg requestMsg) { TenantId tenantId = new TenantId(new UUID(requestMsg.getTenantIdMSB(), requestMsg.getTenantIdLSB())); // TODO: Tenant Profile from cache diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 8e691ba48e..cab87e9680 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/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: diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/data/DeviceConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/data/DeviceConfiguration.java index 1ea2ee4f97..5c9a116abe 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/data/DeviceConfiguration.java +++ b/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 diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/session/FeatureType.java b/common/message/src/main/java/org/thingsboard/server/common/msg/session/FeatureType.java index 9f421c0e73..58289cc94b 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/session/FeatureType.java +++ b/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 } diff --git a/common/queue/src/main/proto/queue.proto b/common/queue/src/main/proto/queue.proto index 4caf51a7b4..007c17f3d3 100644 --- a/common/queue/src/main/proto/queue.proto +++ b/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 */ diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java index b846bac38c..63af515c00 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java +++ b/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 { + 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 { private final CoapExchange exchange; diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java index 82c0b80547..fd010a7141 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java +++ b/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; + } diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java index 5c9c471570..28292f6c20 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java +++ b/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) { diff --git a/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java b/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java index 404024b202..cdc2791839 100644 --- a/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java +++ b/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 provisionDevice(@RequestBody String json, HttpServletRequest httpRequest) { + DeferredResult responseWriter = new DeferredResult<>(); + transportContext.getTransportService().process(JsonConverter.convertToProvisionRequestMsg(json), + new DeviceProvisionCallback(responseWriter)); + return responseWriter; + } + private static class DeviceAuthCallback implements TransportServiceCallback { private final TransportContext transportContext; private final DeferredResult responseWriter; @@ -230,6 +239,25 @@ public class DeviceApiController { } } + private static class DeviceProvisionCallback implements TransportServiceCallback { + private final DeferredResult responseWriter; + + DeviceProvisionCallback(DeferredResult 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 { private final TransportService transportService; private final SessionInfoProto sessionInfo; 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 85206cf1d1..4fc41d6267 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 @@ -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 { diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java index 04da376ff3..995b9eca91 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java +++ b/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 callback); + void process(ProvisionDeviceRequestMsg msg, + TransportServiceCallback callback); + void getDeviceProfile(DeviceProfileId deviceProfileId, TransportServiceCallback callback); void onProfileUpdate(DeviceProfile deviceProfile); @@ -84,8 +88,6 @@ public interface TransportService { void process(SessionInfoProto sessionInfo, ClaimDeviceMsg msg, TransportServiceCallback callback); - void process(SessionInfoProto sessionInfo, ProvisionDeviceRequestMsg msg, TransportServiceCallback deviceProvisionCallback); - void registerAsyncSession(SessionInfoProto sessionInfo, SessionMsgListener listener); void registerSyncSession(SessionInfoProto sessionInfo, SessionMsgListener listener, long timeout); 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 35bfd00cb0..e9be0155a7 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 @@ -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 callback) { + log.trace("Processing msg: {}", requestMsg); + TbProtoQueueMsg protoMsg = new TbProtoQueueMsg<>(UUID.randomUUID(), TransportApiRequestMsg.newBuilder().setProvisionDeviceRequestMsg(requestMsg).build()); + ListenableFuture 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 callback) { if (log.isTraceEnabled()) { @@ -484,15 +496,6 @@ public class DefaultTransportService implements TransportService { } } - @Override - public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ProvisionDeviceRequestMsg msg, TransportServiceCallback 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); 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 7cf392c556..eb66de7009 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 @@ -215,6 +215,6 @@ public interface DeviceDao extends Dao { */ PageData findDevicesByTenantIdAndProfileId(UUID tenantId, UUID profileId, PageLink pageLink); - Device findDeviceByTenantIdAndDeviceDataProvisionConfigurationPair(TenantId tenantId, String provisionDeviceKey, String provisionDeviceSecret); + Optional findDeviceByProfileNameAndDeviceDataProvisionConfigurationPair(String profileName, String provisionDeviceKey, String provisionDeviceSecret); } 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 5b118180d8..3c5a9269a1 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 @@ -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 { @@ -38,9 +39,7 @@ public interface DeviceProfileDao extends Dao { 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); } 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 97736cdb46..d528650c10 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,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>{'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); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceRepository.java index aec705e519..8ad229850a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceRepository.java +++ b/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>('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); } 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 6372e23e16..1cbb9bc24c 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 @@ -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 } @Override - public Device findDeviceByTenantIdAndDeviceDataProvisionConfigurationPair(TenantId tenantId, String provisionDeviceKey, String provisionDeviceSecret) { - return DaoUtil.getData(deviceRepository.findDeviceByTenantIdAndDeviceDataProvisionConfigurationPair(tenantId.getId(), provisionDeviceKey, provisionDeviceSecret)); + public Optional findDeviceByProfileNameAndDeviceDataProvisionConfigurationPair(String profileName, String provisionDeviceKey, String provisionDeviceSecret) { + return Optional.ofNullable(DaoUtil.getData(deviceRepository.findDeviceByProfileNameAndDeviceDataProvisionConfigurationPair(profileName, provisionDeviceKey, provisionDeviceSecret))); } private List convertTenantDeviceTypesToDto(UUID tenantId, List types) { 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 b21f8a0ee6..6fa5768603 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 @@ -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