|
|
|
@ -22,7 +22,7 @@ import com.google.common.util.concurrent.Futures; |
|
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
|
import com.google.common.util.concurrent.SettableFuture; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.springframework.beans.factory.annotation.Autowired; |
|
|
|
import org.apache.commons.lang3.RandomStringUtils; |
|
|
|
import org.springframework.stereotype.Component; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
import org.thingsboard.server.common.data.DataConstants; |
|
|
|
@ -42,8 +42,6 @@ import org.thingsboard.server.common.data.id.EdgeId; |
|
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
|
import org.thingsboard.server.common.data.id.OtaPackageId; |
|
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
import org.thingsboard.server.common.data.page.PageData; |
|
|
|
import org.thingsboard.server.common.data.page.PageLink; |
|
|
|
import org.thingsboard.server.common.data.relation.EntityRelation; |
|
|
|
import org.thingsboard.server.common.data.relation.RelationTypeGroup; |
|
|
|
import org.thingsboard.server.common.data.rpc.RpcError; |
|
|
|
@ -64,9 +62,7 @@ import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
|
|
|
import org.thingsboard.server.gen.transport.TransportProtos; |
|
|
|
import org.thingsboard.server.queue.TbQueueCallback; |
|
|
|
import org.thingsboard.server.queue.TbQueueMsgMetadata; |
|
|
|
import org.thingsboard.server.queue.util.DataDecodingEncodingService; |
|
|
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
|
|
|
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
|
|
|
import org.thingsboard.server.service.rpc.FromDeviceRpcResponseActorMsg; |
|
|
|
|
|
|
|
import java.util.Optional; |
|
|
|
@ -75,54 +71,17 @@ import java.util.UUID; |
|
|
|
@Component |
|
|
|
@Slf4j |
|
|
|
@TbCoreComponent |
|
|
|
public class DeviceEdgeProcessor extends BaseEdgeProcessor { |
|
|
|
|
|
|
|
@Autowired |
|
|
|
private DataDecodingEncodingService dataDecodingEncodingService; |
|
|
|
public class DeviceEdgeProcessor extends BaseDeviceEdgeProcessor { |
|
|
|
|
|
|
|
public ListenableFuture<Void> processDeviceFromEdge(TenantId tenantId, Edge edge, DeviceUpdateMsg deviceUpdateMsg) { |
|
|
|
log.trace("[{}] executing processDeviceFromEdge [{}] from edge [{}]", tenantId, deviceUpdateMsg, edge.getName()); |
|
|
|
DeviceId deviceId = new DeviceId(new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB())); |
|
|
|
try { |
|
|
|
log.trace("[{}] processDeviceFromEdge [{}] from edge [{}]", tenantId, deviceUpdateMsg, edge.getName()); |
|
|
|
DeviceId deviceId = new DeviceId(new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB())); |
|
|
|
switch (deviceUpdateMsg.getMsgType()) { |
|
|
|
case ENTITY_CREATED_RPC_MESSAGE: |
|
|
|
String deviceName = deviceUpdateMsg.getName(); |
|
|
|
Device device = deviceService.findDeviceByTenantIdAndName(tenantId, deviceName); |
|
|
|
if (device != null) { |
|
|
|
boolean deviceAlreadyExistsForThisEdge = isDeviceAlreadyExistsOnCloudForThisEdge(tenantId, edge, device); |
|
|
|
if (deviceAlreadyExistsForThisEdge) { |
|
|
|
log.info("[{}] Device with name '{}' already exists on the cloud, and related to this edge [{}]. " + |
|
|
|
"deviceUpdateMsg [{}], Updating device", tenantId, deviceName, edge.getId(), deviceUpdateMsg); |
|
|
|
return updateDevice(tenantId, edge, deviceUpdateMsg); |
|
|
|
} else { |
|
|
|
log.info("[{}] Device with name '{}' already exists on the cloud, but not related to this edge [{}]. deviceUpdateMsg [{}]." + |
|
|
|
"Creating a new device with random prefix and relate to this edge", tenantId, deviceName, edge.getId(), deviceUpdateMsg); |
|
|
|
String newDeviceName = deviceUpdateMsg.getName() + "_" + StringUtils.randomAlphabetic(15); |
|
|
|
try { |
|
|
|
createDevice(tenantId, deviceId, edge, deviceUpdateMsg, newDeviceName); |
|
|
|
} catch (DataValidationException e) { |
|
|
|
log.error("[{}] Device update msg can't be processed due to data validation [{}]", tenantId, deviceUpdateMsg, e); |
|
|
|
return Futures.immediateFuture(null); |
|
|
|
} |
|
|
|
ObjectNode body = JacksonUtil.OBJECT_MAPPER.createObjectNode(); |
|
|
|
body.put("conflictName", deviceName); |
|
|
|
ListenableFuture<Void> input = saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ENTITY_MERGE_REQUEST, deviceId, body); |
|
|
|
return Futures.transformAsync(input, unused -> |
|
|
|
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_REQUEST, deviceId, null), |
|
|
|
dbCallbackExecutorService); |
|
|
|
} |
|
|
|
} else { |
|
|
|
log.info("[{}] Creating new device on the cloud [{}]", tenantId, deviceUpdateMsg); |
|
|
|
try { |
|
|
|
createDevice(tenantId, deviceId, edge, deviceUpdateMsg, deviceUpdateMsg.getName()); |
|
|
|
} catch (DataValidationException e) { |
|
|
|
log.error("[{}] Device update msg can't be processed due to data validation [{}]", tenantId, deviceUpdateMsg, e); |
|
|
|
return Futures.immediateFuture(null); |
|
|
|
} |
|
|
|
return saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_REQUEST, deviceId, null); |
|
|
|
} |
|
|
|
case ENTITY_UPDATED_RPC_MESSAGE: |
|
|
|
return updateDevice(tenantId, edge, deviceUpdateMsg); |
|
|
|
saveOrUpdateDevice(tenantId, deviceId, deviceUpdateMsg, edge); |
|
|
|
return saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_REQUEST, deviceId, null); |
|
|
|
case ENTITY_DELETED_RPC_MESSAGE: |
|
|
|
Device deviceToDelete = deviceService.findDeviceById(tenantId, deviceId); |
|
|
|
if (deviceToDelete != null) { |
|
|
|
@ -133,64 +92,35 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { |
|
|
|
default: |
|
|
|
return handleUnsupportedMsgType(deviceUpdateMsg.getMsgType()); |
|
|
|
} |
|
|
|
} catch (Exception e) { |
|
|
|
log.error("Failed to process device message from edge, {}", deviceUpdateMsg, e); |
|
|
|
return Futures.immediateFailedFuture(e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private boolean isDeviceAlreadyExistsOnCloudForThisEdge(TenantId tenantId, Edge edge, Device device) { |
|
|
|
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); |
|
|
|
PageData<EdgeId> pageData; |
|
|
|
do { |
|
|
|
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, device.getId(), pageLink); |
|
|
|
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { |
|
|
|
if (pageData.getData().contains(edge.getId())) { |
|
|
|
return true; |
|
|
|
} |
|
|
|
if (pageData.hasNext()) { |
|
|
|
pageLink = pageLink.nextPageLink(); |
|
|
|
} |
|
|
|
} catch (DataValidationException e) { |
|
|
|
if (e.getMessage().contains("Can't create more then")) { |
|
|
|
log.warn("[{}] Number of allowed devices violated {}", tenantId, deviceUpdateMsg, e); |
|
|
|
return Futures.immediateFuture(null); |
|
|
|
} else { |
|
|
|
return Futures.immediateFailedFuture(e); |
|
|
|
} |
|
|
|
} while (pageData != null && pageData.hasNext()); |
|
|
|
return false; |
|
|
|
} |
|
|
|
|
|
|
|
public ListenableFuture<Void> processDeviceCredentialsFromEdge(TenantId tenantId, DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg) { |
|
|
|
log.debug("[{}] Executing processDeviceCredentialsFromEdge, deviceCredentialsUpdateMsg [{}]", tenantId, deviceCredentialsUpdateMsg); |
|
|
|
DeviceId deviceId = new DeviceId(new UUID(deviceCredentialsUpdateMsg.getDeviceIdMSB(), deviceCredentialsUpdateMsg.getDeviceIdLSB())); |
|
|
|
ListenableFuture<Device> deviceFuture = deviceService.findDeviceByIdAsync(tenantId, deviceId); |
|
|
|
return Futures.transform(deviceFuture, device -> { |
|
|
|
if (device != null) { |
|
|
|
log.debug("Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]", |
|
|
|
device.getName(), deviceCredentialsUpdateMsg.getCredentialsId(), deviceCredentialsUpdateMsg.getCredentialsValue()); |
|
|
|
try { |
|
|
|
DeviceCredentials deviceCredentials = deviceCredentialsService.findDeviceCredentialsByDeviceId(tenantId, device.getId()); |
|
|
|
deviceCredentials.setCredentialsType(DeviceCredentialsType.valueOf(deviceCredentialsUpdateMsg.getCredentialsType())); |
|
|
|
deviceCredentials.setCredentialsId(deviceCredentialsUpdateMsg.getCredentialsId()); |
|
|
|
if (deviceCredentialsUpdateMsg.hasCredentialsValue()) { |
|
|
|
deviceCredentials.setCredentialsValue(deviceCredentialsUpdateMsg.getCredentialsValue()); |
|
|
|
} |
|
|
|
deviceCredentialsService.updateDeviceCredentials(tenantId, deviceCredentials); |
|
|
|
} catch (Exception e) { |
|
|
|
log.error("Can't update device credentials for device [{}], deviceCredentialsUpdateMsg [{}]", device.getName(), deviceCredentialsUpdateMsg, e); |
|
|
|
throw new RuntimeException(e); |
|
|
|
} |
|
|
|
} |
|
|
|
return null; |
|
|
|
}, dbCallbackExecutorService); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void createDevice(TenantId tenantId, DeviceId deviceId, Edge edge, DeviceUpdateMsg deviceUpdateMsg, String deviceName) { |
|
|
|
private void saveOrUpdateDevice(TenantId tenantId, DeviceId deviceId, DeviceUpdateMsg deviceUpdateMsg, Edge edge) { |
|
|
|
deviceCreationLock.lock(); |
|
|
|
try { |
|
|
|
Device device = deviceService.findDeviceById(tenantId, deviceId); |
|
|
|
boolean created = false; |
|
|
|
boolean deviceNameUpdated = false; |
|
|
|
String deviceName = deviceUpdateMsg.getName(); |
|
|
|
if (device == null) { |
|
|
|
created = true; |
|
|
|
device = new Device(); |
|
|
|
device.setTenantId(tenantId); |
|
|
|
device.setCreatedTime(Uuids.unixTimestamp(deviceId.getId())); |
|
|
|
created = true; |
|
|
|
Device deviceByName = deviceService.findDeviceByTenantIdAndName(tenantId, deviceName); |
|
|
|
if (deviceByName != null) { |
|
|
|
deviceName = deviceName + "_" + RandomStringUtils.randomAlphabetic(15); |
|
|
|
log.warn("Device with name {} already exists on the cloud. Renaming device name to {}", |
|
|
|
deviceUpdateMsg.getName(), deviceName); |
|
|
|
deviceNameUpdated = true; |
|
|
|
} |
|
|
|
} |
|
|
|
device.setName(deviceName); |
|
|
|
device.setType(deviceUpdateMsg.getType()); |
|
|
|
@ -212,7 +142,6 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { |
|
|
|
|
|
|
|
UUID softwareUUID = safeGetUUID(deviceUpdateMsg.getSoftwareIdMSB(), deviceUpdateMsg.getSoftwareIdLSB()); |
|
|
|
device.setSoftwareId(softwareUUID != null ? new OtaPackageId(softwareUUID) : null); |
|
|
|
|
|
|
|
deviceValidator.validate(device, Device::getTenantId); |
|
|
|
if (created) { |
|
|
|
device.setId(deviceId); |
|
|
|
@ -224,45 +153,21 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { |
|
|
|
deviceCredentials.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN); |
|
|
|
deviceCredentials.setCredentialsId(StringUtils.randomAlphanumeric(20)); |
|
|
|
deviceCredentialsService.createDeviceCredentials(device.getTenantId(), deviceCredentials); |
|
|
|
|
|
|
|
createRelationFromEdge(tenantId, edge.getId(), device.getId()); |
|
|
|
pushDeviceCreatedEventToRuleEngine(tenantId, edge, device); |
|
|
|
deviceService.assignDeviceToEdge(tenantId, device.getId(), edge.getId()); |
|
|
|
} |
|
|
|
createRelationFromEdge(tenantId, edge.getId(), device.getId()); |
|
|
|
pushDeviceCreatedEventToRuleEngine(tenantId, edge, device); |
|
|
|
deviceService.assignDeviceToEdge(edge.getTenantId(), device.getId(), edge.getId()); |
|
|
|
tbClusterService.onDeviceUpdated(savedDevice, created ? null : device, false); |
|
|
|
|
|
|
|
if (deviceNameUpdated) { |
|
|
|
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.UPDATED, deviceId, null); |
|
|
|
} |
|
|
|
} finally { |
|
|
|
deviceCreationLock.unlock(); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<Void> updateDevice(TenantId tenantId, Edge edge, DeviceUpdateMsg deviceUpdateMsg) { |
|
|
|
DeviceId deviceId = new DeviceId(new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB())); |
|
|
|
Device device = deviceService.findDeviceById(tenantId, deviceId); |
|
|
|
device.setName(deviceUpdateMsg.getName()); |
|
|
|
device.setType(deviceUpdateMsg.getType()); |
|
|
|
device.setLabel(deviceUpdateMsg.hasLabel() ? deviceUpdateMsg.getLabel() : null); |
|
|
|
device.setAdditionalInfo(deviceUpdateMsg.hasAdditionalInfo() |
|
|
|
? JacksonUtil.toJsonNode(deviceUpdateMsg.getAdditionalInfo()) : null); |
|
|
|
|
|
|
|
UUID deviceProfileUUID = safeGetUUID(deviceUpdateMsg.getDeviceProfileIdMSB(), deviceUpdateMsg.getDeviceProfileIdLSB()); |
|
|
|
device.setDeviceProfileId(deviceProfileUUID != null ? new DeviceProfileId(deviceProfileUUID) : null); |
|
|
|
|
|
|
|
device.setCustomerId(safeGetCustomerId(deviceUpdateMsg.getCustomerIdMSB(), deviceUpdateMsg.getCustomerIdLSB())); |
|
|
|
|
|
|
|
Optional<DeviceData> deviceDataOpt = |
|
|
|
dataDecodingEncodingService.decode(deviceUpdateMsg.getDeviceDataBytes().toByteArray()); |
|
|
|
device.setDeviceData(deviceDataOpt.orElse(null)); |
|
|
|
|
|
|
|
UUID firmwareUUID = safeGetUUID(deviceUpdateMsg.getFirmwareIdMSB(), deviceUpdateMsg.getFirmwareIdLSB()); |
|
|
|
device.setFirmwareId(firmwareUUID != null ? new OtaPackageId(firmwareUUID) : null); |
|
|
|
|
|
|
|
UUID softwareUUID = safeGetUUID(deviceUpdateMsg.getSoftwareIdMSB(), deviceUpdateMsg.getSoftwareIdLSB()); |
|
|
|
device.setSoftwareId(softwareUUID != null ? new OtaPackageId(softwareUUID) : null); |
|
|
|
|
|
|
|
Device savedDevice = deviceService.saveDevice(device); |
|
|
|
tbClusterService.onDeviceUpdated(savedDevice, device, false); |
|
|
|
return saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_REQUEST, deviceId, null); |
|
|
|
} |
|
|
|
|
|
|
|
private void createRelationFromEdge(TenantId tenantId, EdgeId edgeId, EntityId entityId) { |
|
|
|
EntityRelation relation = new EntityRelation(); |
|
|
|
relation.setFrom(edgeId); |
|
|
|
@ -403,7 +308,7 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { |
|
|
|
if (device != null) { |
|
|
|
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); |
|
|
|
DeviceUpdateMsg deviceUpdateMsg = |
|
|
|
deviceMsgConstructor.constructDeviceUpdatedMsg(msgType, device, null); |
|
|
|
deviceMsgConstructor.constructDeviceUpdatedMsg(msgType, device); |
|
|
|
DownlinkMsg.Builder builder = DownlinkMsg.newBuilder() |
|
|
|
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|
|
|
.addDeviceUpdateMsg(deviceUpdateMsg); |
|
|
|
@ -438,8 +343,6 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { |
|
|
|
return convertRpcCallEventToDownlink(edgeEvent); |
|
|
|
case CREDENTIALS_REQUEST: |
|
|
|
return convertCredentialsRequestEventToDownlink(edgeEvent); |
|
|
|
case ENTITY_MERGE_REQUEST: |
|
|
|
return convertEntityMergeRequestEventToDownlink(edgeEvent); |
|
|
|
} |
|
|
|
return downlinkMsg; |
|
|
|
} |
|
|
|
@ -463,21 +366,6 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor { |
|
|
|
return builder.build(); |
|
|
|
} |
|
|
|
|
|
|
|
public DownlinkMsg convertEntityMergeRequestEventToDownlink(EdgeEvent edgeEvent) { |
|
|
|
DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); |
|
|
|
Device device = deviceService.findDeviceById(edgeEvent.getTenantId(), deviceId); |
|
|
|
String conflictName = null; |
|
|
|
if(edgeEvent.getBody() != null) { |
|
|
|
conflictName = edgeEvent.getBody().get("conflictName").asText(); |
|
|
|
} |
|
|
|
DeviceUpdateMsg deviceUpdateMsg = deviceMsgConstructor |
|
|
|
.constructDeviceUpdatedMsg(UpdateMsgType.ENTITY_MERGE_RPC_MESSAGE, device, conflictName); |
|
|
|
return DownlinkMsg.newBuilder() |
|
|
|
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|
|
|
.addDeviceUpdateMsg(deviceUpdateMsg) |
|
|
|
.build(); |
|
|
|
} |
|
|
|
|
|
|
|
public ListenableFuture<Void> processDeviceNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { |
|
|
|
return processEntityNotification(tenantId, edgeNotificationMsg); |
|
|
|
} |
|
|
|
|