@ -148,35 +148,35 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
private TbClusterService tbClusterService ;
@Override
public void sync ( Edge edge ) {
log . trace ( "[{}][{}] Staring edge sync process" , edge . ge tT enantId( ) , edge . getId ( ) ) ;
public void sync ( TenantId tenantId , Edge edge ) {
log . trace ( "[{}][{}] Staring edge sync process" , tenantId , edge . getId ( ) ) ;
try {
syncWidgetsBundleAndWidgetTypes ( edge ) ;
syncWidgetsBundleAndWidgetTypes ( tenantId , edge ) ;
// TODO: voba - implement this functionality
// syncAdminSettings(edge);
syncRuleChains ( edge ) ;
syncDeviceProfiles ( edge ) ;
syncUsers ( edge ) ;
syncDevices ( edge ) ;
syncAssets ( edge ) ;
syncEntityViews ( edge ) ;
syncDashboards ( edge ) ;
syncRuleChains ( tenantId , edge ) ;
syncDeviceProfiles ( tenantId , edge ) ;
syncUsers ( tenantId , edge ) ;
syncDevices ( tenantId , edge ) ;
syncAssets ( tenantId , edge ) ;
syncEntityViews ( tenantId , edge ) ;
syncDashboards ( tenantId , edge ) ;
} catch ( Exception e ) {
log . error ( "[{}][{}] Exception during sync process" , edge . ge tT enantId( ) , edge . getId ( ) , e ) ;
log . error ( "[{}][{}] Exception during sync process" , tenantId , edge . getId ( ) , e ) ;
}
}
private void syncRuleChains ( Edge edge ) {
log . trace ( "[{}] syncRuleChains [{}]" , edge . ge tT enantId( ) , edge . getName ( ) ) ;
private void syncRuleChains ( TenantId tenantId , Edge edge ) {
log . trace ( "[{}] syncRuleChains [{}]" , tenantId , edge . getName ( ) ) ;
try {
TimePageLink pageLink = new TimePageLink ( DEFAULT_LIMIT ) ;
PageData < RuleChain > pageData ;
do {
pageData = ruleChainService . findRuleChainsByTenantIdAndEdgeId ( edge . ge tT enantId( ) , edge . getId ( ) , pageLink ) ;
pageData = ruleChainService . findRuleChainsByTenantIdAndEdgeId ( tenantId , edge . getId ( ) , pageLink ) ;
if ( pageData ! = null & & pageData . getData ( ) ! = null & & ! pageData . getData ( ) . isEmpty ( ) ) {
log . trace ( "[{}] [{}] rule chains(s) are going to be pushed to edge." , edge . getId ( ) , pageData . getData ( ) . size ( ) ) ;
for ( RuleChain ruleChain : pageData . getData ( ) ) {
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . RULE_CHAIN , EdgeEventActionType . ADDED , ruleChain . getId ( ) , null ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . RULE_CHAIN , EdgeEventActionType . ADDED , ruleChain . getId ( ) , null ) ;
}
if ( pageData . hasNext ( ) ) {
pageLink = pageLink . nextPageLink ( ) ;
@ -188,17 +188,17 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
private void syncDevices ( Edge edge ) {
log . trace ( "[{}] syncDevices [{}]" , edge . ge tT enantId( ) , edge . getName ( ) ) ;
private void syncDevices ( TenantId tenantId , Edge edge ) {
log . trace ( "[{}] syncDevices [{}]" , tenantId , edge . getName ( ) ) ;
try {
TimePageLink pageLink = new TimePageLink ( DEFAULT_LIMIT ) ;
PageData < Device > pageData ;
do {
pageData = deviceService . findDevicesByTenantIdAndEdgeId ( edge . ge tT enantId( ) , edge . getId ( ) , pageLink ) ;
pageData = deviceService . findDevicesByTenantIdAndEdgeId ( tenantId , edge . getId ( ) , pageLink ) ;
if ( pageData ! = null & & pageData . getData ( ) ! = null & & ! pageData . getData ( ) . isEmpty ( ) ) {
log . trace ( "[{}] [{}] device(s) are going to be pushed to edge." , edge . getId ( ) , pageData . getData ( ) . size ( ) ) ;
for ( Device device : pageData . getData ( ) ) {
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . DEVICE , EdgeEventActionType . ADDED , device . getId ( ) , null ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . DEVICE , EdgeEventActionType . ADDED , device . getId ( ) , null ) ;
}
if ( pageData . hasNext ( ) ) {
pageLink = pageLink . nextPageLink ( ) ;
@ -210,17 +210,17 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
private void syncDeviceProfiles ( Edge edge ) {
log . trace ( "[{}] syncDeviceProfiles [{}]" , edge . ge tT enantId( ) , edge . getName ( ) ) ;
private void syncDeviceProfiles ( TenantId tenantId , Edge edge ) {
log . trace ( "[{}] syncDeviceProfiles [{}]" , tenantId , edge . getName ( ) ) ;
try {
TimePageLink pageLink = new TimePageLink ( DEFAULT_LIMIT ) ;
PageData < DeviceProfile > pageData ;
do {
pageData = deviceProfileService . findDeviceProfiles ( edge . ge tT enantId( ) , pageLink ) ;
pageData = deviceProfileService . findDeviceProfiles ( tenantId , pageLink ) ;
if ( pageData ! = null & & pageData . getData ( ) ! = null & & ! pageData . getData ( ) . isEmpty ( ) ) {
log . trace ( "[{}] [{}] user(s) are going to be pushed to edge." , edge . getId ( ) , pageData . getData ( ) . size ( ) ) ;
for ( DeviceProfile deviceProfile : pageData . getData ( ) ) {
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . DEVICE_PROFILE , EdgeEventActionType . ADDED , deviceProfile . getId ( ) , null ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . DEVICE_PROFILE , EdgeEventActionType . ADDED , deviceProfile . getId ( ) , null ) ;
}
if ( pageData . hasNext ( ) ) {
pageLink = pageLink . nextPageLink ( ) ;
@ -232,17 +232,17 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
private void syncAssets ( Edge edge ) {
log . trace ( "[{}] syncAssets [{}]" , edge . ge tT enantId( ) , edge . getName ( ) ) ;
private void syncAssets ( TenantId tenantId , Edge edge ) {
log . trace ( "[{}] syncAssets [{}]" , tenantId , edge . getName ( ) ) ;
try {
TimePageLink pageLink = new TimePageLink ( DEFAULT_LIMIT ) ;
PageData < Asset > pageData ;
do {
pageData = assetService . findAssetsByTenantIdAndEdgeId ( edge . ge tT enantId( ) , edge . getId ( ) , pageLink ) ;
pageData = assetService . findAssetsByTenantIdAndEdgeId ( tenantId , edge . getId ( ) , pageLink ) ;
if ( pageData ! = null & & pageData . getData ( ) ! = null & & ! pageData . getData ( ) . isEmpty ( ) ) {
log . trace ( "[{}] [{}] asset(s) are going to be pushed to edge." , edge . getId ( ) , pageData . getData ( ) . size ( ) ) ;
for ( Asset asset : pageData . getData ( ) ) {
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . ASSET , EdgeEventActionType . ADDED , asset . getId ( ) , null ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . ASSET , EdgeEventActionType . ADDED , asset . getId ( ) , null ) ;
}
if ( pageData . hasNext ( ) ) {
pageLink = pageLink . nextPageLink ( ) ;
@ -254,17 +254,17 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
private void syncEntityViews ( Edge edge ) {
log . trace ( "[{}] syncEntityViews [{}]" , edge . ge tT enantId( ) , edge . getName ( ) ) ;
private void syncEntityViews ( TenantId tenantId , Edge edge ) {
log . trace ( "[{}] syncEntityViews [{}]" , tenantId , edge . getName ( ) ) ;
try {
TimePageLink pageLink = new TimePageLink ( DEFAULT_LIMIT ) ;
PageData < EntityView > pageData ;
do {
pageData = entityViewService . findEntityViewsByTenantIdAndEdgeId ( edge . ge tT enantId( ) , edge . getId ( ) , pageLink ) ;
pageData = entityViewService . findEntityViewsByTenantIdAndEdgeId ( tenantId , edge . getId ( ) , pageLink ) ;
if ( pageData ! = null & & pageData . getData ( ) ! = null & & ! pageData . getData ( ) . isEmpty ( ) ) {
log . trace ( "[{}] [{}] entity view(s) are going to be pushed to edge." , edge . getId ( ) , pageData . getData ( ) . size ( ) ) ;
for ( EntityView entityView : pageData . getData ( ) ) {
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . ENTITY_VIEW , EdgeEventActionType . ADDED , entityView . getId ( ) , null ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . ENTITY_VIEW , EdgeEventActionType . ADDED , entityView . getId ( ) , null ) ;
}
if ( pageData . hasNext ( ) ) {
pageLink = pageLink . nextPageLink ( ) ;
@ -276,17 +276,17 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
private void syncDashboards ( Edge edge ) {
log . trace ( "[{}] syncDashboards [{}]" , edge . ge tT enantId( ) , edge . getName ( ) ) ;
private void syncDashboards ( TenantId tenantId , Edge edge ) {
log . trace ( "[{}] syncDashboards [{}]" , tenantId , edge . getName ( ) ) ;
try {
TimePageLink pageLink = new TimePageLink ( DEFAULT_LIMIT ) ;
PageData < DashboardInfo > pageData ;
do {
pageData = dashboardService . findDashboardsByTenantIdAndEdgeId ( edge . ge tT enantId( ) , edge . getId ( ) , pageLink ) ;
pageData = dashboardService . findDashboardsByTenantIdAndEdgeId ( tenantId , edge . getId ( ) , pageLink ) ;
if ( pageData ! = null & & pageData . getData ( ) ! = null & & ! pageData . getData ( ) . isEmpty ( ) ) {
log . trace ( "[{}] [{}] dashboard(s) are going to be pushed to edge." , edge . getId ( ) , pageData . getData ( ) . size ( ) ) ;
for ( DashboardInfo dashboardInfo : pageData . getData ( ) ) {
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . DASHBOARD , EdgeEventActionType . ADDED , dashboardInfo . getId ( ) , null ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . DASHBOARD , EdgeEventActionType . ADDED , dashboardInfo . getId ( ) , null ) ;
}
if ( pageData . hasNext ( ) ) {
pageLink = pageLink . nextPageLink ( ) ;
@ -298,32 +298,32 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
private void syncUsers ( Edge edge ) {
log . trace ( "[{}] syncUsers [{}]" , edge . ge tT enantId( ) , edge . getName ( ) ) ;
private void syncUsers ( TenantId tenantId , Edge edge ) {
log . trace ( "[{}] syncUsers [{}]" , tenantId , edge . getName ( ) ) ;
try {
TimePageLink pageLink = new TimePageLink ( DEFAULT_LIMIT ) ;
PageData < User > pageData ;
do {
pageData = userService . findTenantAdmins ( edge . ge tT enantId( ) , pageLink ) ;
pushUsersToEdge ( pageData , edge ) ;
pageData = userService . findTenantAdmins ( tenantId , pageLink ) ;
pushUsersToEdge ( tenantId , pageData , edge ) ;
if ( pageData . hasNext ( ) ) {
pageLink = pageLink . nextPageLink ( ) ;
}
} while ( pageData . hasNext ( ) ) ;
syncCustomerUsers ( edge ) ;
syncCustomerUsers ( tenantId , edge ) ;
} catch ( Exception e ) {
log . error ( "Exception during loading edge user(s) on sync!" , e ) ;
}
}
private void syncCustomerUsers ( Edge edge ) {
private void syncCustomerUsers ( TenantId tenantId , Edge edge ) {
if ( edge . getCustomerId ( ) ! = null & & ! EntityId . NULL_UUID . equals ( edge . getCustomerId ( ) . getId ( ) ) ) {
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . CUSTOMER , EdgeEventActionType . ADDED , edge . getCustomerId ( ) , null ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . CUSTOMER , EdgeEventActionType . ADDED , edge . getCustomerId ( ) , null ) ;
TimePageLink pageLink = new TimePageLink ( DEFAULT_LIMIT ) ;
PageData < User > pageData ;
do {
pageData = userService . findCustomerUsers ( edge . ge tT enantId( ) , edge . getCustomerId ( ) , pageLink ) ;
pushUsersToEdge ( pageData , edge ) ;
pageData = userService . findCustomerUsers ( tenantId , edge . getCustomerId ( ) , pageLink ) ;
pushUsersToEdge ( tenantId , pageData , edge ) ;
if ( pageData ! = null & & pageData . hasNext ( ) ) {
pageLink = pageLink . nextPageLink ( ) ;
}
@ -331,45 +331,45 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
}
private void pushUsersToEdge ( PageData < User > pageData , Edge edge ) {
private void pushUsersToEdge ( TenantId tenantId , PageData < User > pageData , Edge edge ) {
if ( pageData ! = null & & pageData . getData ( ) ! = null & & ! pageData . getData ( ) . isEmpty ( ) ) {
log . trace ( "[{}] [{}] user(s) are going to be pushed to edge." , edge . getId ( ) , pageData . getData ( ) . size ( ) ) ;
for ( User user : pageData . getData ( ) ) {
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . USER , EdgeEventActionType . ADDED , user . getId ( ) , null ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . USER , EdgeEventActionType . ADDED , user . getId ( ) , null ) ;
}
}
}
private void syncWidgetsBundleAndWidgetTypes ( Edge edge ) {
log . trace ( "[{}] syncWidgetsBundleAndWidgetTypes [{}]" , edge . ge tT enantId( ) , edge . getName ( ) ) ;
private void syncWidgetsBundleAndWidgetTypes ( TenantId tenantId , Edge edge ) {
log . trace ( "[{}] syncWidgetsBundleAndWidgetTypes [{}]" , tenantId , edge . getName ( ) ) ;
List < WidgetsBundle > widgetsBundlesToPush = new ArrayList < > ( ) ;
List < WidgetType > widgetTypesToPush = new ArrayList < > ( ) ;
widgetsBundlesToPush . addAll ( widgetsBundleService . findAllTenantWidgetsBundlesByTenantId ( edge . ge tT enantId( ) ) ) ;
widgetsBundlesToPush . addAll ( widgetsBundleService . findSystemWidgetsBundles ( edge . ge tT enantId( ) ) ) ;
widgetsBundlesToPush . addAll ( widgetsBundleService . findAllTenantWidgetsBundlesByTenantId ( tenantId ) ) ;
widgetsBundlesToPush . addAll ( widgetsBundleService . findSystemWidgetsBundles ( tenantId ) ) ;
try {
for ( WidgetsBundle widgetsBundle : widgetsBundlesToPush ) {
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . WIDGETS_BUNDLE , EdgeEventActionType . ADDED , widgetsBundle . getId ( ) , null ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . WIDGETS_BUNDLE , EdgeEventActionType . ADDED , widgetsBundle . getId ( ) , null ) ;
widgetTypesToPush . addAll ( widgetTypeService . findWidgetTypesByTenantIdAndBundleAlias ( widgetsBundle . getTenantId ( ) , widgetsBundle . getAlias ( ) ) ) ;
}
for ( WidgetType widgetType : widgetTypesToPush ) {
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . WIDGET_TYPE , EdgeEventActionType . ADDED , widgetType . getId ( ) , null ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . WIDGET_TYPE , EdgeEventActionType . ADDED , widgetType . getId ( ) , null ) ;
}
} catch ( Exception e ) {
log . error ( "Exception during loading widgets bundle(s) and widget type(s) on sync!" , e ) ;
}
}
private void syncAdminSettings ( Edge edge ) {
log . trace ( "[{}] syncAdminSettings [{}]" , edge . ge tT enantId( ) , edge . getName ( ) ) ;
private void syncAdminSettings ( TenantId tenantId , Edge edge ) {
log . trace ( "[{}] syncAdminSettings [{}]" , tenantId , edge . getName ( ) ) ;
try {
AdminSettings systemMailSettings = adminSettingsService . findAdminSettingsByKey ( TenantId . SYS_TENANT_ID , "mail" ) ;
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . ADMIN_SETTINGS , EdgeEventActionType . UPDATED , null , mapper . valueToTree ( systemMailSettings ) ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . ADMIN_SETTINGS , EdgeEventActionType . UPDATED , null , mapper . valueToTree ( systemMailSettings ) ) ;
AdminSettings tenantMailSettings = convertToTenantAdminSettings ( systemMailSettings . getKey ( ) , ( ObjectNode ) systemMailSettings . getJsonValue ( ) ) ;
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . ADMIN_SETTINGS , EdgeEventActionType . UPDATED , null , mapper . valueToTree ( tenantMailSettings ) ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . ADMIN_SETTINGS , EdgeEventActionType . UPDATED , null , mapper . valueToTree ( tenantMailSettings ) ) ;
AdminSettings systemMailTemplates = loadMailTemplates ( ) ;
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . ADMIN_SETTINGS , EdgeEventActionType . UPDATED , null , mapper . valueToTree ( systemMailTemplates ) ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . ADMIN_SETTINGS , EdgeEventActionType . UPDATED , null , mapper . valueToTree ( systemMailTemplates ) ) ;
AdminSettings tenantMailTemplates = convertToTenantAdminSettings ( systemMailTemplates . getKey ( ) , ( ObjectNode ) systemMailTemplates . getJsonValue ( ) ) ;
saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . ADMIN_SETTINGS , EdgeEventActionType . UPDATED , null , mapper . valueToTree ( tenantMailTemplates ) ) ;
saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . ADMIN_SETTINGS , EdgeEventActionType . UPDATED , null , mapper . valueToTree ( tenantMailTemplates ) ) ;
} catch ( Exception e ) {
log . error ( "Can't load admin settings" , e ) ;
}
@ -428,13 +428,13 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
@Override
public ListenableFuture < Void > processRuleChainMetadataRequestMsg ( Edge edge , RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg ) {
log . trace ( "[{}] processRuleChainMetadataRequestMsg [{}][{}]" , edge . ge tT enantId( ) , edge . getName ( ) , ruleChainMetadataRequestMsg ) ;
public ListenableFuture < Void > processRuleChainMetadataRequestMsg ( TenantId tenantId , Edge edge , RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg ) {
log . trace ( "[{}] processRuleChainMetadataRequestMsg [{}][{}]" , tenantId , edge . getName ( ) , ruleChainMetadataRequestMsg ) ;
SettableFuture < Void > futureToSet = SettableFuture . create ( ) ;
if ( ruleChainMetadataRequestMsg . getRuleChainIdMSB ( ) ! = 0 & & ruleChainMetadataRequestMsg . getRuleChainIdLSB ( ) ! = 0 ) {
RuleChainId ruleChainId =
new RuleChainId ( new UUID ( ruleChainMetadataRequestMsg . getRuleChainIdMSB ( ) , ruleChainMetadataRequestMsg . getRuleChainIdLSB ( ) ) ) ;
ListenableFuture < EdgeEvent > future = saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . RULE_CHAIN_METADATA , EdgeEventActionType . ADDED , ruleChainId , null ) ;
ListenableFuture < EdgeEvent > future = saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . RULE_CHAIN_METADATA , EdgeEventActionType . ADDED , ruleChainId , null ) ;
Futures . addCallback ( future , new FutureCallback < EdgeEvent > ( ) {
@Override
public void onSuccess ( @Nullable EdgeEvent result ) {
@ -452,8 +452,8 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
@Override
public ListenableFuture < Void > processAttributesRequestMsg ( Edge edge , AttributesRequestMsg attributesRequestMsg ) {
log . trace ( "[{}] processAttributesRequestMsg [{}][{}]" , edge . ge tT enantId( ) , edge . getName ( ) , attributesRequestMsg ) ;
public ListenableFuture < Void > processAttributesRequestMsg ( TenantId tenantId , Edge edge , AttributesRequestMsg attributesRequestMsg ) {
log . trace ( "[{}] processAttributesRequestMsg [{}][{}]" , tenantId , edge . getName ( ) , attributesRequestMsg ) ;
EntityId entityId = EntityIdFactory . getByTypeAndUuid (
EntityType . valueOf ( attributesRequestMsg . getEntityType ( ) ) ,
new UUID ( attributesRequestMsg . getEntityIdMSB ( ) , attributesRequestMsg . getEntityIdLSB ( ) ) ) ;
@ -461,7 +461,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
if ( type ! = null ) {
SettableFuture < Void > futureToSet = SettableFuture . create ( ) ;
String scope = attributesRequestMsg . getScope ( ) ;
ListenableFuture < List < AttributeKvEntry > > ssAttrFuture = attributesService . findAll ( edge . ge tT enantId( ) , entityId , scope ) ;
ListenableFuture < List < AttributeKvEntry > > ssAttrFuture = attributesService . findAll ( tenantId , entityId , scope ) ;
Futures . addCallback ( ssAttrFuture , new FutureCallback < List < AttributeKvEntry > > ( ) {
@Override
public void onSuccess ( @Nullable List < AttributeKvEntry > ssAttributes ) {
@ -484,7 +484,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
entityData . put ( "scope" , scope ) ;
JsonNode body = mapper . valueToTree ( entityData ) ;
log . debug ( "Sending attributes data msg, entityId [{}], attributes [{}]" , entityId , body ) ;
saveEdgeEvent ( edge . ge tT enantId( ) ,
saveEdgeEvent ( tenantId ,
edge . getId ( ) ,
type ,
EdgeEventActionType . ATTRIBUTES_UPDATED ,
@ -495,7 +495,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
throw new RuntimeException ( "[" + edge . getName ( ) + "] Failed to send attribute updates to the edge" , e ) ;
}
} else {
log . trace ( "[{}][{}] No attributes found for entity {} [{}]" , edge . ge tT enantId( ) ,
log . trace ( "[{}][{}] No attributes found for entity {} [{}]" , tenantId ,
edge . getName ( ) ,
entityId . getEntityType ( ) ,
entityId . getId ( ) ) ;
@ -511,21 +511,21 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} , dbCallbackExecutorService ) ;
return futureToSet ;
} else {
log . warn ( "[{}] Type doesn't supported {}" , edge . ge tT enantId( ) , entityId . getEntityType ( ) ) ;
log . warn ( "[{}] Type doesn't supported {}" , tenantId , entityId . getEntityType ( ) ) ;
return Futures . immediateFuture ( null ) ;
}
}
@Override
public ListenableFuture < Void > processRelationRequestMsg ( Edge edge , RelationRequestMsg relationRequestMsg ) {
log . trace ( "[{}] processRelationRequestMsg [{}][{}]" , edge . ge tT enantId( ) , edge . getName ( ) , relationRequestMsg ) ;
public ListenableFuture < Void > processRelationRequestMsg ( TenantId tenantId , Edge edge , RelationRequestMsg relationRequestMsg ) {
log . trace ( "[{}] processRelationRequestMsg [{}][{}]" , tenantId , edge . getName ( ) , relationRequestMsg ) ;
EntityId entityId = EntityIdFactory . getByTypeAndUuid (
EntityType . valueOf ( relationRequestMsg . getEntityType ( ) ) ,
new UUID ( relationRequestMsg . getEntityIdMSB ( ) , relationRequestMsg . getEntityIdLSB ( ) ) ) ;
List < ListenableFuture < List < EntityRelation > > > futures = new ArrayList < > ( ) ;
futures . add ( findRelationByQuery ( edge , entityId , EntitySearchDirection . FROM ) ) ;
futures . add ( findRelationByQuery ( edge , entityId , EntitySearchDirection . TO ) ) ;
futures . add ( findRelationByQuery ( tenantId , edge , entityId , EntitySearchDirection . FROM ) ) ;
futures . add ( findRelationByQuery ( tenantId , edge , entityId , EntitySearchDirection . TO ) ) ;
ListenableFuture < List < List < EntityRelation > > > relationsListFuture = Futures . allAsList ( futures ) ;
SettableFuture < Void > futureToSet = SettableFuture . create ( ) ;
Futures . addCallback ( relationsListFuture , new FutureCallback < List < List < EntityRelation > > > ( ) {
@ -539,7 +539,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
try {
if ( ! relation . getFrom ( ) . getEntityType ( ) . equals ( EntityType . EDGE ) & &
! relation . getTo ( ) . getEntityType ( ) . equals ( EntityType . EDGE ) ) {
saveEdgeEvent ( edge . ge tT enantId( ) ,
saveEdgeEvent ( tenantId ,
edge . getId ( ) ,
EdgeEventType . RELATION ,
EdgeEventActionType . ADDED ,
@ -563,26 +563,26 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
@Override
public void onFailure ( Throwable t ) {
log . error ( "[{}] Can't find relation by query. Entity id [{}]" , edge . ge tT enantId( ) , entityId , t ) ;
log . error ( "[{}] Can't find relation by query. Entity id [{}]" , tenantId , entityId , t ) ;
futureToSet . setException ( t ) ;
}
} , dbCallbackExecutorService ) ;
return futureToSet ;
}
private ListenableFuture < List < EntityRelation > > findRelationByQuery ( Edge edge , EntityId entityId , EntitySearchDirection direction ) {
private ListenableFuture < List < EntityRelation > > findRelationByQuery ( TenantId tenantId , Edge edge , EntityId entityId , EntitySearchDirection direction ) {
EntityRelationsQuery query = new EntityRelationsQuery ( ) ;
query . setParameters ( new RelationsSearchParameters ( entityId , direction , - 1 , false ) ) ;
return relationService . findByQuery ( edge . ge tT enantId( ) , query ) ;
return relationService . findByQuery ( tenantId , query ) ;
}
@Override
public ListenableFuture < Void > processDeviceCredentialsRequestMsg ( Edge edge , DeviceCredentialsRequestMsg deviceCredentialsRequestMsg ) {
log . trace ( "[{}] processDeviceCredentialsRequestMsg [{}][{}]" , edge . ge tT enantId( ) , edge . getName ( ) , deviceCredentialsRequestMsg ) ;
public ListenableFuture < Void > processDeviceCredentialsRequestMsg ( TenantId tenantId , Edge edge , DeviceCredentialsRequestMsg deviceCredentialsRequestMsg ) {
log . trace ( "[{}] processDeviceCredentialsRequestMsg [{}][{}]" , tenantId , edge . getName ( ) , deviceCredentialsRequestMsg ) ;
SettableFuture < Void > futureToSet = SettableFuture . create ( ) ;
if ( deviceCredentialsRequestMsg . getDeviceIdMSB ( ) ! = 0 & & deviceCredentialsRequestMsg . getDeviceIdLSB ( ) ! = 0 ) {
DeviceId deviceId = new DeviceId ( new UUID ( deviceCredentialsRequestMsg . getDeviceIdMSB ( ) , deviceCredentialsRequestMsg . getDeviceIdLSB ( ) ) ) ;
ListenableFuture < EdgeEvent > future = saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . DEVICE , EdgeEventActionType . CREDENTIALS_UPDATED , deviceId , null ) ;
ListenableFuture < EdgeEvent > future = saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . DEVICE , EdgeEventActionType . CREDENTIALS_UPDATED , deviceId , null ) ;
Futures . addCallback ( future , new FutureCallback < EdgeEvent > ( ) {
@Override
public void onSuccess ( @Nullable EdgeEvent result ) {
@ -600,12 +600,12 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}
@Override
public ListenableFuture < Void > processUserCredentialsRequestMsg ( Edge edge , UserCredentialsRequestMsg userCredentialsRequestMsg ) {
log . trace ( "[{}] processUserCredentialsRequestMsg [{}][{}]" , edge . ge tT enantId( ) , edge . getName ( ) , userCredentialsRequestMsg ) ;
public ListenableFuture < Void > processUserCredentialsRequestMsg ( TenantId tenantId , Edge edge , UserCredentialsRequestMsg userCredentialsRequestMsg ) {
log . trace ( "[{}] processUserCredentialsRequestMsg [{}][{}]" , tenantId , edge . getName ( ) , userCredentialsRequestMsg ) ;
SettableFuture < Void > futureToSet = SettableFuture . create ( ) ;
if ( userCredentialsRequestMsg . getUserIdMSB ( ) ! = 0 & & userCredentialsRequestMsg . getUserIdLSB ( ) ! = 0 ) {
UserId userId = new UserId ( new UUID ( userCredentialsRequestMsg . getUserIdMSB ( ) , userCredentialsRequestMsg . getUserIdLSB ( ) ) ) ;
ListenableFuture < EdgeEvent > future = saveEdgeEvent ( edge . ge tT enantId( ) , edge . getId ( ) , EdgeEventType . USER , EdgeEventActionType . CREDENTIALS_UPDATED , userId , null ) ;
ListenableFuture < EdgeEvent > future = saveEdgeEvent ( tenantId , edge . getId ( ) , EdgeEventType . USER , EdgeEventActionType . CREDENTIALS_UPDATED , userId , null ) ;
Futures . addCallback ( future , new FutureCallback < EdgeEvent > ( ) {
@Override
public void onSuccess ( @Nullable EdgeEvent result ) {