@ -15,6 +15,8 @@
* /
package org.thingsboard.server.service.install.update ;
import com.fasterxml.jackson.databind.JsonNode ;
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 ;
@ -23,9 +25,13 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired ;
import org.springframework.context.annotation.Profile ;
import org.springframework.stereotype.Service ;
import org.springframework.util.StringUtils ;
import org.thingsboard.rule.engine.profile.TbDeviceProfileNode ;
import org.thingsboard.rule.engine.profile.TbDeviceProfileNodeConfiguration ;
import org.thingsboard.server.common.data.EntityView ;
import org.thingsboard.server.common.data.SearchTextBased ;
import org.thingsboard.server.common.data.Tenant ;
import org.thingsboard.server.common.data.id.EntityId ;
import org.thingsboard.server.common.data.id.EntityViewId ;
import org.thingsboard.server.common.data.id.TenantId ;
import org.thingsboard.server.common.data.id.UUIDBased ;
@ -35,10 +41,13 @@ import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.page.PageData ;
import org.thingsboard.server.common.data.page.PageLink ;
import org.thingsboard.server.common.data.rule.RuleChain ;
import org.thingsboard.server.common.data.rule.RuleChainMetaData ;
import org.thingsboard.server.common.data.rule.RuleNode ;
import org.thingsboard.server.dao.entityview.EntityViewService ;
import org.thingsboard.server.dao.rule.RuleChainService ;
import org.thingsboard.server.dao.tenant.TenantService ;
import org.thingsboard.server.dao.timeseries.TimeseriesService ;
import org.thingsboard.server.dao.util.mapping.JacksonUtil ;
import org.thingsboard.server.service.install.InstallScripts ;
import javax.annotation.Nullable ;
@ -49,6 +58,7 @@ import java.util.concurrent.ExecutionException;
import java.util.stream.Collectors ;
import static org.apache.commons.lang.StringUtils.isBlank ;
import static org.thingsboard.server.service.install.DatabaseHelper.objectMapper ;
@Service
@Profile ( "install" )
@ -81,6 +91,10 @@ public class DefaultDataUpdateService implements DataUpdateService {
log . info ( "Updating data from version 3.0.1 to 3.1.0 ..." ) ;
tenantsEntityViewsUpdater . updateEntities ( null ) ;
break ;
case "3.1.1" :
log . info ( "Updating data from version 3.1.1 to 3.2.0 ..." ) ;
tenantsRootRuleChainUpdater . updateEntities ( null ) ;
break ;
default :
throw new RuntimeException ( "Unable to update data, unsupported fromVersion: " + fromVersion ) ;
}
@ -107,6 +121,60 @@ public class DefaultDataUpdateService implements DataUpdateService {
}
} ;
private PaginatedUpdater < String , Tenant > tenantsRootRuleChainUpdater =
new PaginatedUpdater < String , Tenant > ( ) {
@Override
protected PageData < Tenant > findEntities ( String region , PageLink pageLink ) {
return tenantService . findTenants ( pageLink ) ;
}
@Override
protected void updateEntity ( Tenant tenant ) {
try {
RuleChain ruleChain = ruleChainService . getRootTenantRuleChain ( tenant . getId ( ) ) ;
if ( ruleChain = = null ) {
installScripts . createDefaultRuleChains ( tenant . getId ( ) ) ;
} else {
RuleChainMetaData md = ruleChainService . loadRuleChainMetaData ( tenant . getId ( ) , ruleChain . getId ( ) ) ;
int oldIdx = md . getFirstNodeIndex ( ) ;
int newIdx = md . getNodes ( ) . size ( ) ;
if ( md . getNodes ( ) . size ( ) < oldIdx ) {
// Skip invalid rule chains
return ;
}
RuleNode oldFirstNode = md . getNodes ( ) . get ( oldIdx ) ;
if ( oldFirstNode . getType ( ) . equals ( TbDeviceProfileNode . class . getName ( ) ) ) {
// No need to update the rule node twice.
return ;
}
RuleNode ruleNode = new RuleNode ( ) ;
ruleNode . setRuleChainId ( ruleChain . getId ( ) ) ;
ruleNode . setName ( "Device Profile Node" ) ;
ruleNode . setType ( TbDeviceProfileNode . class . getName ( ) ) ;
ruleNode . setDebugMode ( false ) ;
TbDeviceProfileNodeConfiguration ruleNodeConfiguration = new TbDeviceProfileNodeConfiguration ( ) . defaultConfiguration ( ) ;
ruleNode . setConfiguration ( JacksonUtil . valueToTree ( ruleNodeConfiguration ) ) ;
ObjectNode additionalInfo = JacksonUtil . newObjectNode ( ) ;
additionalInfo . put ( "description" , "Process incoming messages from devices with the alarm rules defined in the device profile. Dispatch all incoming messages with \"Success\" relation type." ) ;
additionalInfo . put ( "layoutX" , 204 ) ;
additionalInfo . put ( "layoutY" , 240 ) ;
ruleNode . setAdditionalInfo ( additionalInfo ) ;
md . getNodes ( ) . add ( ruleNode ) ;
md . setFirstNodeIndex ( newIdx ) ;
md . addConnectionInfo ( newIdx , oldIdx , "Success" ) ;
ruleChainService . saveRuleChainMetaData ( tenant . getId ( ) , md ) ;
}
} catch ( Exception e ) {
log . error ( "Unable to update Tenant" , e ) ;
}
}
} ;
private PaginatedUpdater < String , Tenant > tenantsEntityViewsUpdater =
new PaginatedUpdater < String , Tenant > ( ) {
@ -121,30 +189,30 @@ public class DefaultDataUpdateService implements DataUpdateService {
}
} ;
private void updateTenantEntityViews ( TenantId tenantId ) {
PageLink pageLink = new PageLink ( 100 ) ;
PageData < EntityView > pageData = entityViewService . findEntityViewByTenantId ( tenantId , pageLink ) ;
boolean hasNext = true ;
while ( hasNext ) {
List < ListenableFuture < List < Void > > > updateFutures = new ArrayList < > ( ) ;
for ( EntityView entityView : pageData . getData ( ) ) {
updateFutures . add ( updateEntityViewLatestTelemetry ( entityView ) ) ;
}
try {
Futures . allAsList ( updateFutures ) . get ( ) ;
} catch ( InterruptedException | ExecutionException e ) {
log . error ( "Failed to copy latest telemetry to entity view" , e ) ;
}
if ( pageData . hasNext ( ) ) {
pageLink = pageLink . nextPageLink ( ) ;
pageData = entityViewService . findEntityViewByTenantId ( tenantId , pageLink ) ;
} else {
hasNext = false ;
}
}
}
private void updateTenantEntityViews ( TenantId tenantId ) {
PageLink pageLink = new PageLink ( 100 ) ;
PageData < EntityView > pageData = entityViewService . findEntityViewByTenantId ( tenantId , pageLink ) ;
boolean hasNext = true ;
while ( hasNext ) {
List < ListenableFuture < List < Void > > > updateFutures = new ArrayList < > ( ) ;
for ( EntityView entityView : pageData . getData ( ) ) {
updateFutures . add ( updateEntityViewLatestTelemetry ( entityView ) ) ;
}
try {
Futures . allAsList ( updateFutures ) . get ( ) ;
} catch ( InterruptedException | ExecutionException e ) {
log . error ( "Failed to copy latest telemetry to entity view" , e ) ;
}
if ( pageData . hasNext ( ) ) {
pageLink = pageLink . nextPageLink ( ) ;
pageData = entityViewService . findEntityViewByTenantId ( tenantId , pageLink ) ;
} else {
hasNext = false ;
}
}
}
private ListenableFuture < List < Void > > updateEntityViewLatestTelemetry ( EntityView entityView ) {
EntityViewId entityId = entityView . getId ( ) ;
@ -160,13 +228,13 @@ public class DefaultDataUpdateService implements DataUpdateService {
keysFuture = Futures . immediateFuture ( keys ) ;
}
ListenableFuture < List < TsKvEntry > > latestFuture = Futures . transformAsync ( keysFuture , fetchKeys - > {
List < ReadTsKvQuery > queries = fetchKeys . stream ( ) . filter ( key - > ! isBlank ( key ) ) . map ( key - > new BaseReadTsKvQuery ( key , startTs , endTs , 1 , "DESC" ) ) . collect ( Collectors . toList ( ) ) ;
if ( ! queries . isEmpty ( ) ) {
return tsService . findAll ( TenantId . SYS_TENANT_ID , entityView . getEntityId ( ) , queries ) ;
} else {
return Futures . immediateFuture ( null ) ;
}
} , MoreExecutors . directExecutor ( ) ) ;
List < ReadTsKvQuery > queries = fetchKeys . stream ( ) . filter ( key - > ! isBlank ( key ) ) . map ( key - > new BaseReadTsKvQuery ( key , startTs , endTs , 1 , "DESC" ) ) . collect ( Collectors . toList ( ) ) ;
if ( ! queries . isEmpty ( ) ) {
return tsService . findAll ( TenantId . SYS_TENANT_ID , entityView . getEntityId ( ) , queries ) ;
} else {
return Futures . immediateFuture ( null ) ;
}
} , MoreExecutors . directExecutor ( ) ) ;
return Futures . transformAsync ( latestFuture , latestValues - > {
if ( latestValues ! = null & & ! latestValues . isEmpty ( ) ) {
ListenableFuture < List < Void > > saveFuture = tsService . saveLatest ( TenantId . SYS_TENANT_ID , entityId , latestValues ) ;