Browse Source

Merge branch 'feature/usage-records' of https://github.com/thingsboard/thingsboard into feature/usage-records

pull/3615/head
YevhenBondarenko 6 years ago
parent
commit
42dd386182
  1. 9
      application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
  2. 39
      application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
  3. 6
      application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java
  4. 123
      application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java
  5. 20
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  6. 48
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
  7. 18
      application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java
  8. 6
      common/dao-api/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateService.java
  9. 30
      common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java
  10. 8
      dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiApiUsageStateServiceImpl.java

9
application/src/main/java/org/thingsboard/server/actors/app/AppActor.java

@ -40,8 +40,6 @@ import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg;
import org.thingsboard.server.common.msg.queue.RuleEngineException; import org.thingsboard.server.common.msg.queue.RuleEngineException;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.tenant.TenantProfileService;
import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.service.profile.TbTenantProfileCache; import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper; import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper;
@ -150,15 +148,12 @@ public class AppActor extends ContextAwareActor {
private void onComponentLifecycleMsg(ComponentLifecycleMsg msg) { private void onComponentLifecycleMsg(ComponentLifecycleMsg msg) {
TbActorRef target = null; TbActorRef target = null;
if (TenantId.SYS_TENANT_ID.equals(msg.getTenantId())) { if (TenantId.SYS_TENANT_ID.equals(msg.getTenantId())) {
if (msg.getEntityId().getEntityType() == EntityType.TENANT_PROFILE) { if (!EntityType.TENANT_PROFILE.equals(msg.getEntityId().getEntityType())) {
tenantProfileCache.evict(new TenantProfileId(msg.getEntityId().getId()));
} else {
log.warn("Message has system tenant id: {}", msg); log.warn("Message has system tenant id: {}", msg);
} }
} else { } else {
if (msg.getEntityId().getEntityType() == EntityType.TENANT) { if (EntityType.TENANT.equals(msg.getEntityId().getEntityType())) {
TenantId tenantId = new TenantId(msg.getEntityId().getId()); TenantId tenantId = new TenantId(msg.getEntityId().getId());
tenantProfileCache.evict(tenantId);
if (msg.getEvent() == ComponentLifecycleEvent.DELETED) { if (msg.getEvent() == ComponentLifecycleEvent.DELETED) {
log.info("[{}] Handling tenant deleted notification: {}", msg.getTenantId(), msg); log.info("[{}] Handling tenant deleted notification: {}", msg.getTenantId(), msg);
deletedTenants.add(tenantId); deletedTenants.add(tenantId);

39
application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java

@ -22,6 +22,7 @@ import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.ApiUsageState; import org.thingsboard.server.common.data.ApiUsageState;
import org.thingsboard.server.common.data.TenantProfile; import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
@ -87,6 +88,7 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService {
TenantId tenantId = new TenantId(new UUID(statsMsg.getTenantIdMSB(), statsMsg.getTenantIdLSB())); TenantId tenantId = new TenantId(new UUID(statsMsg.getTenantIdMSB(), statsMsg.getTenantIdLSB()));
TenantApiUsageState tenantState; TenantApiUsageState tenantState;
List<TsKvEntry> updatedEntries; List<TsKvEntry> updatedEntries;
boolean stateUpdated = false;
updateLock.lock(); updateLock.lock();
try { try {
tenantState = getOrFetchState(tenantId); tenantState = getOrFetchState(tenantId);
@ -101,13 +103,21 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService {
ApiUsageRecordKey recordKey = ApiUsageRecordKey.valueOf(kvProto.getKey()); ApiUsageRecordKey recordKey = ApiUsageRecordKey.valueOf(kvProto.getKey());
long newValue = tenantState.add(recordKey, kvProto.getValue()); long newValue = tenantState.add(recordKey, kvProto.getValue());
updatedEntries.add(new BasicTsKvEntry(ts, new LongDataEntry(recordKey.name(), newValue))); updatedEntries.add(new BasicTsKvEntry(ts, new LongDataEntry(recordKey.name(), newValue)));
newValue = tenantState.addToHourly(recordKey, kvProto.getValue()); long newHourlyValue = tenantState.addToHourly(recordKey, kvProto.getValue());
updatedEntries.add(new BasicTsKvEntry(hourTs, new LongDataEntry(HOURLY + recordKey.name(), newValue))); updatedEntries.add(new BasicTsKvEntry(hourTs, new LongDataEntry(HOURLY + recordKey.name(), newHourlyValue)));
stateUpdated |= tenantState.checkStateUpdatedDueToThreshold(recordKey);
} }
} finally { } finally {
updateLock.unlock(); updateLock.unlock();
} }
tsService.save(tenantId, tenantState.getId(), updatedEntries, 0L); tsService.save(tenantId, tenantState.getApiUsageState().getId(), updatedEntries, 0L);
if (stateUpdated) {
// Save new state into the database;
apiUsageStateService.update(tenantState.getApiUsageState());
//TODO: clear cache on cluster repartition.
//TODO: update profiles on tenant and profile updates.
//TODO: broadcast to everyone notifications about enabled/disabled features.
}
} }
@Override @Override
@ -116,12 +126,26 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService {
} }
@Override @Override
public void onAddedToAllowList(TenantId tenantId) { public void onTenantProfileUpdate(TenantProfileId tenantProfileId) {
TenantProfile tenantProfile = tenantProfileCache.get(tenantProfileId);
updateLock.lock();
try {
tenantStates.values().forEach(state -> {
if (tenantProfile.getId().equals(state.getTenantProfileId())) {
state.setTenantProfileData(tenantProfile.getProfileData());
if (state.checkStateUpdatedDueToThresholds()) {
apiUsageStateService.update(state.getApiUsageState());
//TODO: send notification to cluster;
}
}
});
} finally {
updateLock.unlock();
}
} }
@Override @Override
public void onAddedToDenyList(TenantId tenantId) { public void onTenantUpdate(TenantId tenantId) {
} }
@ -150,7 +174,8 @@ public class DefaultTbApiUsageStateService implements TbApiUsageStateService {
dbStateEntity = apiUsageStateService.findTenantApiUsageState(tenantId); dbStateEntity = apiUsageStateService.findTenantApiUsageState(tenantId);
} }
} }
tenantState = new TenantApiUsageState(dbStateEntity.getId()); TenantProfile tenantProfile = tenantProfileCache.get(tenantId);
tenantState = new TenantApiUsageState(tenantProfile, dbStateEntity);
try { try {
List<TsKvEntry> dbValues = tsService.findAllLatest(tenantId, dbStateEntity.getEntityId()).get(); List<TsKvEntry> dbValues = tsService.findAllLatest(tenantId, dbStateEntity.getEntityId()).get();
for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) {

6
application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java

@ -16,6 +16,7 @@
package org.thingsboard.server.service.apiusage; package org.thingsboard.server.service.apiusage;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg;
@ -26,8 +27,7 @@ public interface TbApiUsageStateService {
TenantApiUsageState getApiUsageState(TenantId tenantId); TenantApiUsageState getApiUsageState(TenantId tenantId);
void onAddedToAllowList(TenantId tenantId); void onTenantProfileUpdate(TenantProfileId tenantProfileId);
void onAddedToDenyList(TenantId tenantId);
void onTenantUpdate(TenantId tenantId);
} }

123
application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java

@ -16,9 +16,13 @@
package org.thingsboard.server.service.apiusage; package org.thingsboard.server.service.apiusage;
import lombok.Getter; import lombok.Getter;
import lombok.Setter;
import org.thingsboard.server.common.data.ApiUsageRecordKey; import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.id.ApiUsageStateId; import org.thingsboard.server.common.data.ApiUsageState;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.TenantProfileData;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.msg.tools.SchedulerUtils; import org.thingsboard.server.common.msg.tools.SchedulerUtils;
import java.util.Map; import java.util.Map;
@ -30,7 +34,13 @@ public class TenantApiUsageState {
private final Map<ApiUsageRecordKey, Long> currentHourValues = new ConcurrentHashMap<>(); private final Map<ApiUsageRecordKey, Long> currentHourValues = new ConcurrentHashMap<>();
@Getter @Getter
private final ApiUsageStateId id; @Setter
private TenantProfileId tenantProfileId;
@Getter
@Setter
private TenantProfileData tenantProfileData;
@Getter
private final ApiUsageState apiUsageState;
@Getter @Getter
private volatile long currentCycleTs; private volatile long currentCycleTs;
@Getter @Getter
@ -38,8 +48,10 @@ public class TenantApiUsageState {
@Getter @Getter
private volatile long currentHourTs; private volatile long currentHourTs;
public TenantApiUsageState(ApiUsageStateId id) { public TenantApiUsageState(TenantProfile tenantProfile, ApiUsageState apiUsageState) {
this.id = id; this.tenantProfileId = tenantProfile.getId();
this.tenantProfileData = tenantProfile.getProfileData();
this.apiUsageState = apiUsageState;
this.currentCycleTs = SchedulerUtils.getStartOfCurrentMonth(); this.currentCycleTs = SchedulerUtils.getStartOfCurrentMonth();
this.nextCycleTs = SchedulerUtils.getStartOfNextMonth(); this.nextCycleTs = SchedulerUtils.getStartOfNextMonth();
this.currentHourTs = SchedulerUtils.getStartOfCurrentHour(); this.currentHourTs = SchedulerUtils.getStartOfCurrentHour();
@ -59,6 +71,10 @@ public class TenantApiUsageState {
return result; return result;
} }
public long get(ApiUsageRecordKey key) {
return currentCycleValues.getOrDefault(key, 0L);
}
public long addToHourly(ApiUsageRecordKey key, long value) { public long addToHourly(ApiUsageRecordKey key, long value) {
long result = currentHourValues.getOrDefault(key, 0L) + value; long result = currentHourValues.getOrDefault(key, 0L) + value;
currentHourValues.put(key, result); currentHourValues.put(key, result);
@ -80,4 +96,103 @@ public class TenantApiUsageState {
} }
} }
public long getProfileThreshold(ApiUsageRecordKey key) {
Object threshold = tenantProfileData.getProperties().get(key.name());
if (threshold != null) {
if (threshold instanceof String) {
return Long.parseLong((String) threshold);
} else if (threshold instanceof Long) {
return (Long) threshold;
}
}
return 0L;
}
public EntityId getEntityId() {
return apiUsageState.getEntityId();
}
public boolean isTransportEnabled() {
return apiUsageState.isTransportEnabled();
}
public boolean isDbStorageEnabled() {
return apiUsageState.isDbStorageEnabled();
}
public boolean isRuleEngineEnabled() {
return apiUsageState.isRuleEngineEnabled();
}
public boolean isJsExecEnabled() {
return apiUsageState.isJsExecEnabled();
}
public void setTransportEnabled(boolean transportEnabled) {
apiUsageState.setTransportEnabled(transportEnabled);
}
public void setDbStorageEnabled(boolean dbStorageEnabled) {
apiUsageState.setDbStorageEnabled(dbStorageEnabled);
}
public void setRuleEngineEnabled(boolean ruleEngineEnabled) {
apiUsageState.setRuleEngineEnabled(ruleEngineEnabled);
}
public void setJsExecEnabled(boolean jsExecEnabled) {
apiUsageState.setJsExecEnabled(jsExecEnabled);
}
public boolean isFeatureEnabled(ApiUsageRecordKey recordKey) {
switch (recordKey) {
case MSG_COUNT:
case MSG_BYTES_COUNT:
case DP_TRANSPORT_COUNT:
return isTransportEnabled();
case RE_EXEC_COUNT:
return isRuleEngineEnabled();
case DP_STORAGE_COUNT:
return isDbStorageEnabled();
case JS_EXEC_COUNT:
return isJsExecEnabled();
default:
return true;
}
}
public boolean setFeatureValue(ApiUsageRecordKey recordKey, boolean value) {
boolean currentValue = isFeatureEnabled(recordKey);
switch (recordKey) {
case MSG_COUNT:
case MSG_BYTES_COUNT:
case DP_TRANSPORT_COUNT:
setTransportEnabled(value);
break;
case RE_EXEC_COUNT:
setRuleEngineEnabled(value);
break;
case DP_STORAGE_COUNT:
setDbStorageEnabled(value);
break;
case JS_EXEC_COUNT:
setJsExecEnabled(value);
break;
}
return currentValue == value;
}
public boolean checkStateUpdatedDueToThresholds() {
boolean update = false;
for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) {
update |= checkStateUpdatedDueToThreshold(key);
}
return update;
}
public boolean checkStateUpdatedDueToThreshold(ApiUsageRecordKey recordKey) {
long value = get(recordKey);
long threshold = getProfileThreshold(recordKey);
return setFeatureValue(recordKey, threshold == 0 || value < threshold);
}
} }

20
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -54,6 +54,7 @@ import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.service.queue.processing.AbstractConsumerService; import org.thingsboard.server.service.queue.processing.AbstractConsumerService;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; import org.thingsboard.server.service.rpc.FromDeviceRpcResponse;
import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService; import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService;
@ -102,12 +103,19 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
protected volatile ExecutorService usageStatsExecutor; protected volatile ExecutorService usageStatsExecutor;
public DefaultTbCoreConsumerService(TbCoreQueueFactory tbCoreQueueFactory, ActorSystemContext actorContext, public DefaultTbCoreConsumerService(TbCoreQueueFactory tbCoreQueueFactory,
DeviceStateService stateService, TbLocalSubscriptionService localSubscriptionService, ActorSystemContext actorContext,
SubscriptionManagerService subscriptionManagerService, DataDecodingEncodingService encodingService, DeviceStateService stateService,
TbCoreDeviceRpcService tbCoreDeviceRpcService, StatsFactory statsFactory, TbDeviceProfileCache deviceProfileCache, TbLocalSubscriptionService localSubscriptionService,
TbApiUsageStateService statsService) { SubscriptionManagerService subscriptionManagerService,
super(actorContext, encodingService, deviceProfileCache, tbCoreQueueFactory.createToCoreNotificationsMsgConsumer()); DataDecodingEncodingService encodingService,
TbCoreDeviceRpcService tbCoreDeviceRpcService,
StatsFactory statsFactory,
TbDeviceProfileCache deviceProfileCache,
TbApiUsageStateService statsService,
TbTenantProfileCache tenantProfileCache,
TbApiUsageStateService apiUsageStateService) {
super(actorContext, encodingService, tenantProfileCache, deviceProfileCache, apiUsageStateService, tbCoreQueueFactory.createToCoreNotificationsMsgConsumer());
this.mainConsumer = tbCoreQueueFactory.createToCoreMsgConsumer(); this.mainConsumer = tbCoreQueueFactory.createToCoreMsgConsumer();
this.usageStatsConsumer = tbCoreQueueFactory.createToUsageStatsServiceMsgConsumer(); this.usageStatsConsumer = tbCoreQueueFactory.createToUsageStatsServiceMsgConsumer();
this.stateService = stateService; this.stateService = stateService;

48
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java

@ -15,7 +15,6 @@
*/ */
package org.thingsboard.server.service.queue; package org.thingsboard.server.service.queue;
import com.google.protobuf.ByteString;
import com.google.protobuf.ProtocolStringList; import com.google.protobuf.ProtocolStringList;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
@ -24,10 +23,16 @@ import org.springframework.stereotype.Service;
import org.thingsboard.rule.engine.api.RpcError; import org.thingsboard.rule.engine.api.RpcError;
import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.queue.*; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg;
import org.thingsboard.server.common.msg.queue.RuleEngineException;
import org.thingsboard.server.common.msg.queue.RuleNodeInfo;
import org.thingsboard.server.common.msg.queue.ServiceQueue;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TbMsgCallback;
import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
@ -38,17 +43,33 @@ import org.thingsboard.server.queue.provider.TbRuleEngineQueueFactory;
import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings; import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
import org.thingsboard.server.queue.settings.TbRuleEngineQueueConfiguration; import org.thingsboard.server.queue.settings.TbRuleEngineQueueConfiguration;
import org.thingsboard.server.queue.util.TbRuleEngineComponent; import org.thingsboard.server.queue.util.TbRuleEngineComponent;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService; import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.queue.processing.*; import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.service.queue.processing.AbstractConsumerService;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingDecision;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingResult;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingStrategy;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingStrategyFactory;
import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategy;
import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategyFactory;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponse; import org.thingsboard.server.service.rpc.FromDeviceRpcResponse;
import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService; import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService;
import org.thingsboard.server.service.stats.RuleEngineStatisticsService; import org.thingsboard.server.service.stats.RuleEngineStatisticsService;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import java.util.*; import java.util.Collections;
import java.util.concurrent.*; import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
@Service @Service
@TbRuleEngineComponent @TbRuleEngineComponent
@ -79,11 +100,16 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
public DefaultTbRuleEngineConsumerService(TbRuleEngineProcessingStrategyFactory processingStrategyFactory, public DefaultTbRuleEngineConsumerService(TbRuleEngineProcessingStrategyFactory processingStrategyFactory,
TbRuleEngineSubmitStrategyFactory submitStrategyFactory, TbRuleEngineSubmitStrategyFactory submitStrategyFactory,
TbQueueRuleEngineSettings ruleEngineSettings, TbQueueRuleEngineSettings ruleEngineSettings,
TbRuleEngineQueueFactory tbRuleEngineQueueFactory, RuleEngineStatisticsService statisticsService, TbRuleEngineQueueFactory tbRuleEngineQueueFactory,
ActorSystemContext actorContext, DataDecodingEncodingService encodingService, RuleEngineStatisticsService statisticsService,
ActorSystemContext actorContext,
DataDecodingEncodingService encodingService,
TbRuleEngineDeviceRpcService tbDeviceRpcService, TbRuleEngineDeviceRpcService tbDeviceRpcService,
StatsFactory statsFactory, TbDeviceProfileCache deviceProfileCache) { StatsFactory statsFactory,
super(actorContext, encodingService, deviceProfileCache, tbRuleEngineQueueFactory.createToRuleEngineNotificationsMsgConsumer()); TbDeviceProfileCache deviceProfileCache,
TbTenantProfileCache tenantProfileCache,
TbApiUsageStateService apiUsageStateService) {
super(actorContext, encodingService, tenantProfileCache, deviceProfileCache, apiUsageStateService, tbRuleEngineQueueFactory.createToRuleEngineNotificationsMsgConsumer());
this.statisticsService = statisticsService; this.statisticsService = statisticsService;
this.ruleEngineSettings = ruleEngineSettings; this.ruleEngineSettings = ruleEngineSettings;
this.tbRuleEngineQueueFactory = tbRuleEngineQueueFactory; this.tbRuleEngineQueueFactory = tbRuleEngineQueueFactory;

18
application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java

@ -25,6 +25,7 @@ import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.msg.TbActorMsg; import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
@ -33,7 +34,9 @@ import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService; import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache; import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.service.queue.TbPackCallback; import org.thingsboard.server.service.queue.TbPackCallback;
import org.thingsboard.server.service.queue.TbPackProcessingContext; import org.thingsboard.server.service.queue.TbPackProcessingContext;
@ -59,15 +62,19 @@ public abstract class AbstractConsumerService<N extends com.google.protobuf.Gene
protected final ActorSystemContext actorContext; protected final ActorSystemContext actorContext;
protected final DataDecodingEncodingService encodingService; protected final DataDecodingEncodingService encodingService;
protected final TbTenantProfileCache tenantProfileCache;
protected final TbDeviceProfileCache deviceProfileCache; protected final TbDeviceProfileCache deviceProfileCache;
protected final TbApiUsageStateService apiUsageStateService;
protected final TbQueueConsumer<TbProtoQueueMsg<N>> nfConsumer; protected final TbQueueConsumer<TbProtoQueueMsg<N>> nfConsumer;
public AbstractConsumerService(ActorSystemContext actorContext, DataDecodingEncodingService encodingService, public AbstractConsumerService(ActorSystemContext actorContext, DataDecodingEncodingService encodingService,
TbDeviceProfileCache deviceProfileCache, TbQueueConsumer<TbProtoQueueMsg<N>> nfConsumer) { TbTenantProfileCache tenantProfileCache, TbDeviceProfileCache deviceProfileCache, TbApiUsageStateService apiUsageStateService, TbQueueConsumer<TbProtoQueueMsg<N>> nfConsumer) {
this.actorContext = actorContext; this.actorContext = actorContext;
this.encodingService = encodingService; this.encodingService = encodingService;
this.tenantProfileCache = tenantProfileCache;
this.deviceProfileCache = deviceProfileCache; this.deviceProfileCache = deviceProfileCache;
this.apiUsageStateService = apiUsageStateService;
this.nfConsumer = nfConsumer; this.nfConsumer = nfConsumer;
} }
@ -143,7 +150,14 @@ public abstract class AbstractConsumerService<N extends com.google.protobuf.Gene
TbActorMsg actorMsg = actorMsgOpt.get(); TbActorMsg actorMsg = actorMsgOpt.get();
if (actorMsg instanceof ComponentLifecycleMsg) { if (actorMsg instanceof ComponentLifecycleMsg) {
ComponentLifecycleMsg componentLifecycleMsg = (ComponentLifecycleMsg) actorMsg; ComponentLifecycleMsg componentLifecycleMsg = (ComponentLifecycleMsg) actorMsg;
if (EntityType.DEVICE_PROFILE.equals(componentLifecycleMsg.getEntityId().getEntityType())) { if (EntityType.TENANT_PROFILE.equals(componentLifecycleMsg.getEntityId().getEntityType())) {
TenantProfileId tenantProfileId = new TenantProfileId(componentLifecycleMsg.getEntityId().getId());
tenantProfileCache.evict(tenantProfileId);
apiUsageStateService.onTenantProfileUpdate(tenantProfileId);
} else if (EntityType.TENANT.equals(componentLifecycleMsg.getEntityId().getEntityType())) {
tenantProfileCache.evict(componentLifecycleMsg.getTenantId());
apiUsageStateService.onTenantUpdate(componentLifecycleMsg.getTenantId());
} else if (EntityType.DEVICE_PROFILE.equals(componentLifecycleMsg.getEntityId().getEntityType())) {
deviceProfileCache.evict(componentLifecycleMsg.getTenantId(), new DeviceProfileId(componentLifecycleMsg.getEntityId().getId())); deviceProfileCache.evict(componentLifecycleMsg.getTenantId(), new DeviceProfileId(componentLifecycleMsg.getEntityId().getId()));
} else if (EntityType.DEVICE.equals(componentLifecycleMsg.getEntityId().getEntityType())) { } else if (EntityType.DEVICE.equals(componentLifecycleMsg.getEntityId().getEntityType())) {
deviceProfileCache.evict(new DeviceId(componentLifecycleMsg.getEntityId().getId())); deviceProfileCache.evict(new DeviceId(componentLifecycleMsg.getEntityId().getId()));

6
common/dao-api/src/main/java/org/thingsboard/server/dao/usagerecord/ApiUsageStateService.java

@ -21,11 +21,13 @@ import org.thingsboard.server.common.data.id.TenantId;
public interface ApiUsageStateService { public interface ApiUsageStateService {
ApiUsageState createDefaultApiUsageState(TenantId id);
ApiUsageState update(ApiUsageState apiUsageState);
ApiUsageState findTenantApiUsageState(TenantId tenantId); ApiUsageState findTenantApiUsageState(TenantId tenantId);
void deleteApiUsageStateByTenantId(TenantId tenantId); void deleteApiUsageStateByTenantId(TenantId tenantId);
ApiUsageState createDefaultApiUsageState(TenantId id);
ApiUsageState findApiUsageStateById(TenantId tenantId, ApiUsageStateId id); ApiUsageState findApiUsageStateById(TenantId tenantId, ApiUsageStateId id);
} }

30
common/data/src/main/java/org/thingsboard/server/common/data/ApiUsageState.java

@ -16,6 +16,8 @@
package org.thingsboard.server.common.data; package org.thingsboard.server.common.data;
import lombok.EqualsAndHashCode; import lombok.EqualsAndHashCode;
import lombok.Getter;
import lombok.Setter;
import lombok.ToString; import lombok.ToString;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -27,8 +29,18 @@ public class ApiUsageState extends BaseData<ApiUsageStateId> implements HasTenan
private static final long serialVersionUID = 8250339805336035966L; private static final long serialVersionUID = 8250339805336035966L;
@Getter @Setter
private TenantId tenantId; private TenantId tenantId;
@Getter @Setter
private EntityId entityId; private EntityId entityId;
@Getter @Setter
private boolean transportEnabled;
@Getter @Setter
private boolean dbStorageEnabled;
@Getter @Setter
private boolean ruleEngineEnabled;
@Getter @Setter
private boolean jsExecEnabled;
public ApiUsageState() { public ApiUsageState() {
super(); super();
@ -43,22 +55,4 @@ public class ApiUsageState extends BaseData<ApiUsageStateId> implements HasTenan
this.tenantId = ur.getTenantId(); this.tenantId = ur.getTenantId();
this.entityId = ur.getEntityId(); this.entityId = ur.getEntityId();
} }
@Override
public TenantId getTenantId() {
return tenantId;
}
public void setTenantId(TenantId tenantId) {
this.tenantId = tenantId;
}
public EntityId getEntityId() {
return entityId;
}
public void setEntityId(EntityId entityId) {
this.entityId = entityId;
}
} }

8
dao/src/main/java/org/thingsboard/server/dao/usagerecord/ApiApiUsageStateServiceImpl.java

@ -60,6 +60,14 @@ public class ApiApiUsageStateServiceImpl extends AbstractEntityService implement
return apiUsageStateDao.save(apiUsageState.getTenantId(), apiUsageState); return apiUsageStateDao.save(apiUsageState.getTenantId(), apiUsageState);
} }
@Override
public ApiUsageState update(ApiUsageState apiUsageState) {
log.trace("Executing save [{}]", apiUsageState.getTenantId());
validateId(apiUsageState.getTenantId(), INCORRECT_TENANT_ID + apiUsageState.getTenantId());
validateId(apiUsageState.getId(), "Can't save new usage state. Only update is allowed!");
return apiUsageStateDao.save(apiUsageState.getTenantId(), apiUsageState);
}
@Override @Override
public ApiUsageState findTenantApiUsageState(TenantId tenantId) { public ApiUsageState findTenantApiUsageState(TenantId tenantId) {
log.trace("Executing findTenantUsageRecord, tenantId [{}]", tenantId); log.trace("Executing findTenantUsageRecord, tenantId [{}]", tenantId);

Loading…
Cancel
Save