@ -24,13 +24,12 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value ;
import org.springframework.beans.factory.annotation.Value ;
import org.springframework.context.annotation.Lazy ;
import org.springframework.context.annotation.Lazy ;
import org.springframework.stereotype.Service ;
import org.springframework.stereotype.Service ;
import org.thingsboard.common.util.ThingsBoardThreadFactory ;
import org.thingsboard.rule.engine.api.MailService ;
import org.thingsboard.rule.engine.api.MailService ;
import org.thingsboard.server.cluster.TbClusterService ;
import org.thingsboard.server.cluster.TbClusterService ;
import org.thingsboard.server.common.data.ApiFeature ;
import org.thingsboard.server.common.data.ApiFeature ;
import org.thingsboard.server.common.data.ApiUsageRecordKey ;
import org.thingsboard.server.common.data.ApiUsageRecordKey ;
import org.thingsboard.server.common.data.ApiUsageRecordState ;
import org.thingsboard.server.common.data.ApiUsageState ;
import org.thingsboard.server.common.data.ApiUsageState ;
import org.thingsboard.server.common.data.ApiUsageStateMailMessage ;
import org.thingsboard.server.common.data.ApiUsageStateValue ;
import org.thingsboard.server.common.data.ApiUsageStateValue ;
import org.thingsboard.server.common.data.EntityType ;
import org.thingsboard.server.common.data.EntityType ;
import org.thingsboard.server.common.data.StringUtils ;
import org.thingsboard.server.common.data.StringUtils ;
@ -54,6 +53,7 @@ import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo ;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo ;
import org.thingsboard.server.common.msg.tools.SchedulerUtils ;
import org.thingsboard.server.common.msg.tools.SchedulerUtils ;
import org.thingsboard.server.dao.notification.NotificationRuleProcessingService ;
import org.thingsboard.server.dao.notification.NotificationRuleProcessingService ;
import org.thingsboard.server.dao.notification.trigger.ApiUsageLimitTrigger ;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache ;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache ;
import org.thingsboard.server.dao.tenant.TenantService ;
import org.thingsboard.server.dao.tenant.TenantService ;
import org.thingsboard.server.dao.timeseries.TimeseriesService ;
import org.thingsboard.server.dao.timeseries.TimeseriesService ;
@ -63,6 +63,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto;
import org.thingsboard.server.queue.common.TbProtoQueueMsg ;
import org.thingsboard.server.queue.common.TbProtoQueueMsg ;
import org.thingsboard.server.queue.discovery.PartitionService ;
import org.thingsboard.server.queue.discovery.PartitionService ;
import org.thingsboard.server.service.executors.DbCallbackExecutorService ;
import org.thingsboard.server.service.executors.DbCallbackExecutorService ;
import org.thingsboard.server.service.mail.MailExecutorService ;
import org.thingsboard.server.service.partition.AbstractPartitionBasedService ;
import org.thingsboard.server.service.partition.AbstractPartitionBasedService ;
import org.thingsboard.server.service.telemetry.InternalTelemetryService ;
import org.thingsboard.server.service.telemetry.InternalTelemetryService ;
@ -79,8 +80,6 @@ import java.util.Set;
import java.util.UUID ;
import java.util.UUID ;
import java.util.concurrent.ConcurrentHashMap ;
import java.util.concurrent.ConcurrentHashMap ;
import java.util.concurrent.ExecutionException ;
import java.util.concurrent.ExecutionException ;
import java.util.concurrent.ExecutorService ;
import java.util.concurrent.Executors ;
import java.util.concurrent.TimeUnit ;
import java.util.concurrent.TimeUnit ;
import java.util.concurrent.locks.Lock ;
import java.util.concurrent.locks.Lock ;
import java.util.concurrent.locks.ReentrantLock ;
import java.util.concurrent.locks.ReentrantLock ;
@ -110,6 +109,7 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService
private final MailService mailService ;
private final MailService mailService ;
private final NotificationRuleProcessingService notificationRuleProcessingService ;
private final NotificationRuleProcessingService notificationRuleProcessingService ;
private final DbCallbackExecutorService dbExecutor ;
private final DbCallbackExecutorService dbExecutor ;
private final MailExecutorService mailExecutor ;
@Lazy
@Lazy
@Autowired
@Autowired
@ -130,8 +130,6 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService
private final Lock updateLock = new ReentrantLock ( ) ;
private final Lock updateLock = new ReentrantLock ( ) ;
private final ExecutorService mailExecutor = Executors . newSingleThreadExecutor ( ThingsBoardThreadFactory . forName ( "api-usage-svc-mail" ) ) ;
@PostConstruct
@PostConstruct
public void init ( ) {
public void init ( ) {
super . init ( ) ;
super . init ( ) ;
@ -340,32 +338,35 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService
tsWsService . saveAndNotifyInternal ( state . getTenantId ( ) , state . getApiUsageState ( ) . getId ( ) , stateTelemetry , VOID_CALLBACK ) ;
tsWsService . saveAndNotifyInternal ( state . getTenantId ( ) , state . getApiUsageState ( ) . getId ( ) , stateTelemetry , VOID_CALLBACK ) ;
if ( state . getEntityType ( ) = = EntityType . TENANT & & ! state . getEntityId ( ) . equals ( TenantId . SYS_TENANT_ID ) ) {
if ( state . getEntityType ( ) = = EntityType . TENANT & & ! state . getEntityId ( ) . equals ( TenantId . SYS_TENANT_ID ) ) {
String email = tenantService . findTenantById ( state . getTenantId ( ) ) . getEmail ( ) ;
String email = tenantService . findTenantById ( state . getTenantId ( ) ) . getEmail ( ) ;
if ( StringUtils . isNotEmpty ( email ) ) {
result . forEach ( ( apiFeature , stateValue ) - > {
result . forEach ( ( apiFeature , stateValue ) - > {
ApiUsageRecordState recordState = createApiUsageRecordState ( ( TenantApiUsageState ) state , apiFeature , stateValue ) ;
notificationRuleProcessingService . process ( ApiUsageLimitTrigger . builder ( )
. tenantId ( state . getTenantId ( ) )
. state ( recordState )
. status ( stateValue )
. build ( ) ) ;
if ( StringUtils . isNotEmpty ( email ) ) {
mailExecutor . submit ( ( ) - > {
mailExecutor . submit ( ( ) - > {
try {
try {
mailService . sendApiFeatureStateEmail ( apiFeature , stateValue , email , createStateMailMessage ( ( TenantApiUsageState ) state , apiFeature , stateValue ) ) ;
mailService . sendApiFeatureStateEmail ( apiFeature , stateValue , email , recordState ) ;
} catch ( ThingsboardException e ) {
} catch ( ThingsboardException e ) {
log . warn ( "[{}] Can't send update of the API state to tenant with provided email [{}]" , state . getTenantId ( ) , email , e ) ;
log . warn ( "[{}] Can't send update of the API state to tenant with provided email [{}]" , state . getTenantId ( ) , email , e ) ;
}
}
} ) ;
} ) ;
} ) ;
}
} else {
} ) ;
log . warn ( "[{}] Can't send update of the API state to tenant with empty email!" , state . getTenantId ( ) ) ;
}
}
}
}
}
private ApiUsageStateMailMessag e createStateMailMessag e ( TenantApiUsageState state , ApiFeature apiFeature , ApiUsageStateValue stateValue ) {
private ApiUsageRecord State createApiUsageRecord State ( TenantApiUsageState state , ApiFeature apiFeature , ApiUsageStateValue stateValue ) {
StateChecker checker = getStateChecker ( stateValue ) ;
StateChecker checker = getStateChecker ( stateValue ) ;
for ( ApiUsageRecordKey apiUsageRecordKey : ApiUsageRecordKey . getKeys ( apiFeature ) ) {
for ( ApiUsageRecordKey apiUsageRecordKey : ApiUsageRecordKey . getKeys ( apiFeature ) ) {
long threshold = state . getProfileThreshold ( apiUsageRecordKey ) ;
long threshold = state . getProfileThreshold ( apiUsageRecordKey ) ;
long warnThreshold = state . getProfileWarnThreshold ( apiUsageRecordKey ) ;
long warnThreshold = state . getProfileWarnThreshold ( apiUsageRecordKey ) ;
long value = state . get ( apiUsageRecordKey ) ;
long value = state . get ( apiUsageRecordKey ) ;
if ( checker . check ( threshold , warnThreshold , value ) ) {
if ( checker . check ( threshold , warnThreshold , value ) ) {
return new ApiUsageStateMailMessage ( apiUsageRecordKey , threshold , value ) ;
return new ApiUsageRecordState ( apiFeature , apiUsageRecordKey , threshold , value ) ;
}
}
}
}
return null ;
return null ;
@ -377,7 +378,7 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService
} else if ( ApiUsageStateValue . WARNING . equals ( stateValue ) ) {
} else if ( ApiUsageStateValue . WARNING . equals ( stateValue ) ) {
return ( t , wt , v ) - > v < t & & v > = wt ;
return ( t , wt , v ) - > v < t & & v > = wt ;
} else {
} else {
return ( t , wt , v ) - > v > = t ;
return ( t , wt , v ) - > t > 0 & & v > = t ;
}
}
}
}
@ -529,8 +530,5 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService
@PreDestroy
@PreDestroy
private void destroy ( ) {
private void destroy ( ) {
super . stop ( ) ;
super . stop ( ) ;
if ( mailExecutor ! = null ) {
mailExecutor . shutdownNow ( ) ;
}
}
}
}
}