@ -17,12 +17,15 @@ package org.thingsboard.server.service.edge.rpc.processor;
import com.fasterxml.jackson.core.JsonProcessingException ;
import com.fasterxml.jackson.databind.node.ObjectNode ;
import com.google.common.util.concurrent.FutureCallback ;
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.apache.commons.lang.RandomStringUtils ;
import org.apache.commons.lang.StringUtils ;
import org.checkerframework.checker.nullness.qual.Nullable ;
import org.springframework.security.core.parameters.P ;
import org.springframework.stereotype.Component ;
import org.thingsboard.rule.engine.api.RpcError ;
import org.thingsboard.server.common.data.DataConstants ;
@ -52,6 +55,7 @@ import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponse ;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponseActorMsg ;
import java.util.List ;
import java.util.UUID ;
import java.util.concurrent.locks.ReentrantLock ;
@ -64,36 +68,62 @@ public class DeviceProcessor extends BaseProcessor {
public ListenableFuture < Void > onDeviceUpdate ( TenantId tenantId , Edge edge , DeviceUpdateMsg deviceUpdateMsg ) {
log . trace ( "[{}] onDeviceUpdate [{}] from edge [{}]" , tenantId , deviceUpdateMsg , edge . getName ( ) ) ;
DeviceId edgeDeviceId = 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 ) {
log . info ( "[{}] Device with name '{}' already exists on the cloud. Updating id of device entity on the edge" , tenantId , deviceName ) ;
if ( ! device . getId ( ) . equals ( edgeDeviceId ) ) {
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . DEVICE , EdgeEventActionType . ENTITY_EXISTS_REQUEST , device . getId ( ) , null ) ;
}
ListenableFuture < List < EdgeId > > future = edgeService . findRelatedEdgeIdsByEntityId ( tenantId , device . getId ( ) ) ;
SettableFuture < Void > futureToSet = SettableFuture . create ( ) ;
Futures . addCallback ( future , new FutureCallback < List < EdgeId > > ( ) {
@Override
public void onSuccess ( @Nullable List < EdgeId > edgeIds ) {
boolean update = false ;
if ( edgeIds ! = null & & ! edgeIds . isEmpty ( ) ) {
if ( edgeIds . contains ( edge . getId ( ) ) ) {
update = true ;
}
}
Device device ;
if ( update ) {
log . info ( "[{}] Device with name '{}' already exists on the cloud, and related to this edge [{}]. " +
"deviceUpdateMsg [{}], Updating device" , tenantId , deviceName , edge . getId ( ) , deviceUpdateMsg ) ;
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 ( ) + "_" + RandomStringUtils . randomAlphabetic ( 15 ) ;
device = createDevice ( tenantId , edge , deviceUpdateMsg , newDeviceName ) ;
ObjectNode body = mapper . createObjectNode ( ) ;
body . put ( "conflictName" , deviceName ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . DEVICE , EdgeEventActionType . ENTITY_MERGE_REQUEST , device . getId ( ) , body ) ;
deviceService . assignDeviceToEdge ( edge . getTenantId ( ) , device . getId ( ) , edge . getId ( ) ) ;
}
futureToSet . set ( null ) ;
}
@Override
public void onFailure ( Throwable t ) {
log . error ( "[{}] Failed to get related edge ids by device id [{}], edge [{}]" , tenantId , deviceUpdateMsg , edge . getId ( ) , t ) ;
futureToSet . setException ( t ) ;
}
} , dbCallbackExecutorService ) ;
return futureToSet ;
} else {
Device deviceById = deviceService . findDeviceById ( edge . getTenantId ( ) , edgeDeviceId ) ;
if ( deviceById ! = null ) {
log . info ( "[{}] Device ID [{}] already used by other device on the cloud. Creating new device and replacing device entity on the edge" , tenantId , edgeDeviceId . getId ( ) ) ;
device = createDevice ( tenantId , edge , deviceUpdateMsg ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . DEVICE , EdgeEventActionType . ENTITY_EXISTS_REQUEST , device . getId ( ) , null ) ;
} else {
device = createDevice ( tenantId , edge , deviceUpdateMsg ) ;
}
log . info ( "[{}] Creating new device and replacing device entity on the edge [{}]" , tenantId , deviceUpdateMsg ) ;
device = createDevice ( tenantId , edge , deviceUpdateMsg , deviceUpdateMsg . getName ( ) ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . DEVICE , EdgeEventActionType . ENTITY_MERGE_REQUEST , device . getId ( ) , null ) ;
deviceService . assignDeviceToEdge ( edge . getTenantId ( ) , device . getId ( ) , edge . getId ( ) ) ;
}
// TODO: voba - assign device only in case device is not assigned yet. Missing functionality to check this relation prior assignment
deviceService . assignDeviceToEdge ( edge . getTenantId ( ) , device . getId ( ) , edge . getId ( ) ) ;
break ;
case ENTITY_UPDATED_RPC_MESSAGE :
updateDevice ( tenantId , edge , deviceUpdateMsg ) ;
break ;
case ENTITY_DELETED_RPC_MESSAGE :
Device deviceToDelete = deviceService . findDeviceById ( tenantId , edgeDeviceId ) ;
DeviceId deviceId = new DeviceId ( new UUID ( deviceUpdateMsg . getIdMSB ( ) , deviceUpdateMsg . getIdLSB ( ) ) ) ;
Device deviceToDelete = deviceService . findDeviceById ( tenantId , deviceId ) ;
if ( deviceToDelete ! = null ) {
deviceService . unassignDeviceFromEdge ( tenantId , edgeDeviceId , edge . getId ( ) ) ;
deviceService . unassignDeviceFromEdge ( tenantId , deviceId , edge . getId ( ) ) ;
}
break ;
case UNRECOGNIZED :
@ -103,7 +133,6 @@ public class DeviceProcessor extends BaseProcessor {
return Futures . immediateFuture ( null ) ;
}
public ListenableFuture < Void > onDeviceCredentialsUpdate ( TenantId tenantId , DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg ) {
log . debug ( "Executing onDeviceCredentialsUpdate, deviceCredentialsUpdateMsg [{}]" , deviceCredentialsUpdateMsg ) ;
DeviceId deviceId = new DeviceId ( new UUID ( deviceCredentialsUpdateMsg . getDeviceIdMSB ( ) , deviceCredentialsUpdateMsg . getDeviceIdLSB ( ) ) ) ;
@ -131,36 +160,34 @@ public class DeviceProcessor extends BaseProcessor {
private 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 . getLabel ( ) ) ;
device . setAdditionalInfo ( JacksonUtil . toJsonNode ( deviceUpdateMsg . getAdditionalInfo ( ) ) ) ;
deviceService . saveDevice ( device ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . DEVICE , EdgeEventActionType . CREDENTIALS_REQUEST , deviceId , null ) ;
if ( device ! = null ) {
device . setName ( deviceUpdateMsg . getName ( ) ) ;
device . setType ( deviceUpdateMsg . getType ( ) ) ;
device . setLabel ( deviceUpdateMsg . getLabel ( ) ) ;
device . setAdditionalInfo ( JacksonUtil . toJsonNode ( deviceUpdateMsg . getAdditionalInfo ( ) ) ) ;
deviceService . saveDevice ( device ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . DEVICE , EdgeEventActionType . CREDENTIALS_REQUEST , deviceId , null ) ;
} else {
log . warn ( "[{}] can't find device [{}], edge [{}]" , tenantId , deviceUpdateMsg , edge . getId ( ) ) ;
}
}
private Device createDevice ( TenantId tenantId , Edge edge , DeviceUpdateMsg deviceUpdateMsg ) {
private Device createDevice ( TenantId tenantId , Edge edge , DeviceUpdateMsg deviceUpdateMsg , String deviceName ) {
Device device ;
try {
deviceCreationLock . lock ( ) ;
log . debug ( "[{}] Creating device entity [{}] from edge [{}]" , tenantId , deviceUpdateMsg , edge . getName ( ) ) ;
DeviceId deviceId = new DeviceId ( new UUID ( deviceUpdateMsg . getIdMSB ( ) , deviceUpdateMsg . getIdLSB ( ) ) ) ;
device = new Device ( ) ;
device . setTenantId ( edge . getTenantId ( ) ) ;
device . setCustomerId ( edge . getCustomerId ( ) ) ;
device . setId ( deviceId ) ;
device . setName ( deviceUpdateMsg . getName ( ) ) ;
device . setName ( deviceName ) ;
device . setType ( deviceUpdateMsg . getType ( ) ) ;
device . setLabel ( deviceUpdateMsg . getLabel ( ) ) ;
device . setAdditionalInfo ( JacksonUtil . toJsonNode ( deviceUpdateMsg . getAdditionalInfo ( ) ) ) ;
device = deviceService . saveDevice ( device ) ;
createDeviceCredentials ( device ) ;
createRelationFromEdge ( tenantId , edge . getId ( ) , device . getId ( ) ) ;
deviceStateService . onDeviceAdded ( device ) ;
pushDeviceCreatedEventToRuleEngine ( tenantId , edge , device ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . DEVICE , EdgeEventActionType . CREDENTIALS_REQUEST , deviceId , null ) ;
} finally {
deviceCreationLock . unlock ( ) ;
}
@ -176,14 +203,6 @@ public class DeviceProcessor extends BaseProcessor {
relationService . saveRelation ( tenantId , relation ) ;
}
private void createDeviceCredentials ( Device device ) {
DeviceCredentials deviceCredentials = new DeviceCredentials ( ) ;
deviceCredentials . setDeviceId ( device . getId ( ) ) ;
deviceCredentials . setCredentialsType ( DeviceCredentialsType . ACCESS_TOKEN ) ;
deviceCredentials . setCredentialsId ( RandomStringUtils . randomAlphanumeric ( 20 ) ) ;
deviceCredentialsService . createDeviceCredentials ( device . getTenantId ( ) , deviceCredentials ) ;
}
private void pushDeviceCreatedEventToRuleEngine ( TenantId tenantId , Edge edge , Device device ) {
try {
DeviceId deviceId = device . getId ( ) ;