committed by
GitHub
170 changed files with 3870 additions and 913 deletions
@ -0,0 +1,29 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.apiusage; |
|||
|
|||
import lombok.Getter; |
|||
|
|||
public enum ApiFeature { |
|||
TRANSPORT("transportApiState"), DB("dbApiState"), RE("ruleEngineApiState"), JS("jsExecutionApiState"); |
|||
|
|||
@Getter |
|||
private final String apiStateKey; |
|||
|
|||
ApiFeature(String apiStateKey) { |
|||
this.apiStateKey = apiStateKey; |
|||
} |
|||
} |
|||
@ -0,0 +1,357 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.apiusage; |
|||
|
|||
import com.google.common.util.concurrent.FutureCallback; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.checkerframework.checker.nullness.qual.Nullable; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.beans.factory.annotation.Value; |
|||
import org.springframework.context.annotation.Lazy; |
|||
import org.springframework.data.util.Pair; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.common.data.ApiUsageRecordKey; |
|||
import org.thingsboard.server.common.data.ApiUsageState; |
|||
import org.thingsboard.server.common.data.Tenant; |
|||
import org.thingsboard.server.common.data.TenantProfile; |
|||
import org.thingsboard.server.common.data.id.ApiUsageStateId; |
|||
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.BooleanDataEntry; |
|||
import org.thingsboard.server.common.data.kv.LongDataEntry; |
|||
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|||
import org.thingsboard.server.common.data.page.PageDataIterable; |
|||
import org.thingsboard.server.common.data.tenant.profile.TenantProfileConfiguration; |
|||
import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; |
|||
import org.thingsboard.server.common.msg.queue.ServiceType; |
|||
import org.thingsboard.server.common.msg.queue.TbCallback; |
|||
import org.thingsboard.server.common.msg.tools.SchedulerUtils; |
|||
import org.thingsboard.server.dao.tenant.TenantService; |
|||
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|||
import org.thingsboard.server.dao.usagerecord.ApiUsageStateService; |
|||
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg; |
|||
import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto; |
|||
import org.thingsboard.server.queue.common.TbProtoQueueMsg; |
|||
import org.thingsboard.server.queue.discovery.PartitionChangeEvent; |
|||
import org.thingsboard.server.queue.discovery.PartitionService; |
|||
import org.thingsboard.server.queue.scheduler.SchedulerComponent; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.profile.TbTenantProfileCache; |
|||
import org.thingsboard.server.service.queue.TbClusterService; |
|||
import org.thingsboard.server.service.telemetry.InternalTelemetryService; |
|||
|
|||
import javax.annotation.PostConstruct; |
|||
import java.util.ArrayList; |
|||
import java.util.HashMap; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.ConcurrentHashMap; |
|||
import java.util.concurrent.ExecutionException; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.concurrent.locks.Lock; |
|||
import java.util.concurrent.locks.ReentrantLock; |
|||
|
|||
@Slf4j |
|||
@TbCoreComponent |
|||
@Service |
|||
public class DefaultTbApiUsageStateService implements TbApiUsageStateService { |
|||
|
|||
public static final String HOURLY = "Hourly"; |
|||
public static final FutureCallback<Integer> VOID_CALLBACK = new FutureCallback<Integer>() { |
|||
@Override |
|||
public void onSuccess(@Nullable Integer result) { |
|||
} |
|||
|
|||
@Override |
|||
public void onFailure(Throwable t) { |
|||
} |
|||
}; |
|||
private final TbClusterService clusterService; |
|||
private final PartitionService partitionService; |
|||
private final TenantService tenantService; |
|||
private final TimeseriesService tsService; |
|||
private final ApiUsageStateService apiUsageStateService; |
|||
private final SchedulerComponent scheduler; |
|||
private final TbTenantProfileCache tenantProfileCache; |
|||
|
|||
@Lazy |
|||
@Autowired |
|||
private InternalTelemetryService tsWsService; |
|||
|
|||
// Tenants that should be processed on this server
|
|||
private final Map<TenantId, TenantApiUsageState> myTenantStates = new ConcurrentHashMap<>(); |
|||
// Tenants that should be processed on other servers
|
|||
private final Map<TenantId, ApiUsageState> otherTenantStates = new ConcurrentHashMap<>(); |
|||
|
|||
@Value("${usage.stats.report.enabled:true}") |
|||
private boolean enabled; |
|||
|
|||
@Value("${usage.stats.check.cycle:60000}") |
|||
private long nextCycleCheckInterval; |
|||
|
|||
private final Lock updateLock = new ReentrantLock(); |
|||
|
|||
public DefaultTbApiUsageStateService(TbClusterService clusterService, |
|||
PartitionService partitionService, |
|||
TenantService tenantService, |
|||
TimeseriesService tsService, |
|||
ApiUsageStateService apiUsageStateService, |
|||
SchedulerComponent scheduler, |
|||
TbTenantProfileCache tenantProfileCache) { |
|||
this.clusterService = clusterService; |
|||
this.partitionService = partitionService; |
|||
this.tenantService = tenantService; |
|||
this.tsService = tsService; |
|||
this.apiUsageStateService = apiUsageStateService; |
|||
this.scheduler = scheduler; |
|||
this.tenantProfileCache = tenantProfileCache; |
|||
} |
|||
|
|||
@PostConstruct |
|||
public void init() { |
|||
if (enabled) { |
|||
log.info("Starting api usage service."); |
|||
initStatesFromDataBase(); |
|||
scheduler.scheduleAtFixedRate(this::checkStartOfNextCycle, nextCycleCheckInterval, nextCycleCheckInterval, TimeUnit.MILLISECONDS); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void process(TbProtoQueueMsg<ToUsageStatsServiceMsg> msg, TbCallback callback) { |
|||
ToUsageStatsServiceMsg statsMsg = msg.getValue(); |
|||
TenantId tenantId = new TenantId(new UUID(statsMsg.getTenantIdMSB(), statsMsg.getTenantIdLSB())); |
|||
TenantApiUsageState tenantState; |
|||
List<TsKvEntry> updatedEntries; |
|||
Map<ApiFeature, Boolean> result = new HashMap<>(); |
|||
updateLock.lock(); |
|||
try { |
|||
tenantState = getOrFetchState(tenantId); |
|||
long ts = tenantState.getCurrentCycleTs(); |
|||
long hourTs = tenantState.getCurrentHourTs(); |
|||
long newHourTs = SchedulerUtils.getStartOfCurrentHour(); |
|||
if (newHourTs != hourTs) { |
|||
tenantState.setHour(newHourTs); |
|||
} |
|||
updatedEntries = new ArrayList<>(ApiUsageRecordKey.values().length); |
|||
for (UsageStatsKVProto kvProto : statsMsg.getValuesList()) { |
|||
ApiUsageRecordKey recordKey = ApiUsageRecordKey.valueOf(kvProto.getKey()); |
|||
long newValue = tenantState.add(recordKey, kvProto.getValue()); |
|||
updatedEntries.add(new BasicTsKvEntry(ts, new LongDataEntry(recordKey.getApiCountKey(), newValue))); |
|||
long newHourlyValue = tenantState.addToHourly(recordKey, kvProto.getValue()); |
|||
updatedEntries.add(new BasicTsKvEntry(hourTs, new LongDataEntry(recordKey.getApiCountKey() + HOURLY, newHourlyValue))); |
|||
Pair<ApiFeature, Boolean> update = tenantState.checkStateUpdatedDueToThreshold(recordKey); |
|||
if (update != null) { |
|||
result.put(update.getFirst(), update.getSecond()); |
|||
} |
|||
} |
|||
} finally { |
|||
updateLock.unlock(); |
|||
} |
|||
tsWsService.saveAndNotifyInternal(tenantId, tenantState.getApiUsageState().getId(), updatedEntries, VOID_CALLBACK); |
|||
if (!result.isEmpty()) { |
|||
persistAndNotify(tenantState, result); |
|||
} |
|||
callback.onSuccess(); |
|||
} |
|||
|
|||
@Override |
|||
public void onApplicationEvent(PartitionChangeEvent partitionChangeEvent) { |
|||
if (partitionChangeEvent.getServiceType().equals(ServiceType.TB_CORE)) { |
|||
myTenantStates.entrySet().removeIf(entry -> !partitionService.resolve(ServiceType.TB_CORE, entry.getKey(), entry.getKey()).isMyPartition()); |
|||
otherTenantStates.entrySet().removeIf(entry -> partitionService.resolve(ServiceType.TB_CORE, entry.getKey(), entry.getKey()).isMyPartition()); |
|||
initStatesFromDataBase(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public ApiUsageState getApiUsageState(TenantId tenantId) { |
|||
TenantApiUsageState tenantState = myTenantStates.get(tenantId); |
|||
if (tenantState != null) { |
|||
return tenantState.getApiUsageState(); |
|||
} else { |
|||
ApiUsageState state = otherTenantStates.get(tenantId); |
|||
if (state != null) { |
|||
return state; |
|||
} else { |
|||
if (partitionService.resolve(ServiceType.TB_CORE, tenantId, tenantId).isMyPartition()) { |
|||
return getOrFetchState(tenantId).getApiUsageState(); |
|||
} else { |
|||
updateLock.lock(); |
|||
try { |
|||
state = otherTenantStates.get(tenantId); |
|||
if (state == null) { |
|||
state = apiUsageStateService.findTenantApiUsageState(tenantId); |
|||
otherTenantStates.put(tenantId, state); |
|||
} |
|||
} finally { |
|||
updateLock.unlock(); |
|||
} |
|||
return state; |
|||
} |
|||
} |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void onApiUsageStateUpdate(TenantId tenantId) { |
|||
otherTenantStates.remove(tenantId); |
|||
} |
|||
|
|||
@Override |
|||
public void onTenantProfileUpdate(TenantProfileId tenantProfileId) { |
|||
TenantProfile tenantProfile = tenantProfileCache.get(tenantProfileId); |
|||
updateLock.lock(); |
|||
try { |
|||
myTenantStates.values().forEach(state -> { |
|||
if (tenantProfile.getId().equals(state.getTenantProfileId())) { |
|||
updateTenantState(state, tenantProfile); |
|||
} |
|||
}); |
|||
} finally { |
|||
updateLock.unlock(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void onTenantUpdate(TenantId tenantId) { |
|||
TenantProfile tenantProfile = tenantProfileCache.get(tenantId); |
|||
updateLock.lock(); |
|||
try { |
|||
TenantApiUsageState state = myTenantStates.get(tenantId); |
|||
if (state != null && !state.getTenantProfileId().equals(tenantProfile.getId())) { |
|||
updateTenantState(state, tenantProfile); |
|||
} |
|||
} finally { |
|||
updateLock.unlock(); |
|||
} |
|||
} |
|||
|
|||
private void updateTenantState(TenantApiUsageState state, TenantProfile tenantProfile) { |
|||
TenantProfileData oldProfileData = state.getTenantProfileData(); |
|||
state.setTenantProfileId(tenantProfile.getId()); |
|||
state.setTenantProfileData(tenantProfile.getProfileData()); |
|||
Map<ApiFeature, Boolean> result = state.checkStateUpdatedDueToThresholds(); |
|||
if (!result.isEmpty()) { |
|||
persistAndNotify(state, result); |
|||
} |
|||
updateProfileThresholds(state.getTenantId(), state.getApiUsageState().getId(), |
|||
oldProfileData.getConfiguration(), tenantProfile.getProfileData().getConfiguration()); |
|||
} |
|||
|
|||
private void updateProfileThresholds(TenantId tenantId, ApiUsageStateId id, |
|||
TenantProfileConfiguration oldData, TenantProfileConfiguration newData) { |
|||
long ts = System.currentTimeMillis(); |
|||
List<TsKvEntry> profileThresholds = new ArrayList<>(); |
|||
for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { |
|||
long newProfileThreshold = newData.getProfileThreshold(key); |
|||
if (oldData == null || oldData.getProfileThreshold(key) != newProfileThreshold) { |
|||
log.info("[{}] Updating profile threshold [{}]:[{}]", tenantId, key, newProfileThreshold); |
|||
profileThresholds.add(new BasicTsKvEntry(ts, new LongDataEntry(key.getApiLimitKey(), newProfileThreshold))); |
|||
} |
|||
} |
|||
if (!profileThresholds.isEmpty()) { |
|||
tsWsService.saveAndNotifyInternal(tenantId, id, profileThresholds, VOID_CALLBACK); |
|||
} |
|||
} |
|||
|
|||
private void persistAndNotify(TenantApiUsageState state, Map<ApiFeature, Boolean> result) { |
|||
log.info("[{}] Detected update of the API state: {}", state.getTenantId(), result); |
|||
apiUsageStateService.update(state.getApiUsageState()); |
|||
clusterService.onApiStateChange(state.getApiUsageState(), null); |
|||
long ts = System.currentTimeMillis(); |
|||
List<TsKvEntry> stateTelemetry = new ArrayList<>(); |
|||
result.forEach(((apiFeature, aState) -> stateTelemetry.add(new BasicTsKvEntry(ts, new BooleanDataEntry(apiFeature.getApiStateKey(), aState))))); |
|||
tsWsService.saveAndNotifyInternal(state.getTenantId(), state.getApiUsageState().getId(), stateTelemetry, VOID_CALLBACK); |
|||
} |
|||
|
|||
private void checkStartOfNextCycle() { |
|||
updateLock.lock(); |
|||
try { |
|||
long now = System.currentTimeMillis(); |
|||
myTenantStates.values().forEach(state -> { |
|||
if ((state.getNextCycleTs() > now) && (state.getNextCycleTs() - now < TimeUnit.HOURS.toMillis(1))) { |
|||
state.setCycles(state.getNextCycleTs(), SchedulerUtils.getStartOfNextNextMonth()); |
|||
} |
|||
}); |
|||
} finally { |
|||
updateLock.unlock(); |
|||
} |
|||
} |
|||
|
|||
private TenantApiUsageState getOrFetchState(TenantId tenantId) { |
|||
TenantApiUsageState tenantState = myTenantStates.get(tenantId); |
|||
if (tenantState == null) { |
|||
ApiUsageState dbStateEntity = apiUsageStateService.findTenantApiUsageState(tenantId); |
|||
if (dbStateEntity == null) { |
|||
try { |
|||
dbStateEntity = apiUsageStateService.createDefaultApiUsageState(tenantId); |
|||
} catch (Exception e) { |
|||
dbStateEntity = apiUsageStateService.findTenantApiUsageState(tenantId); |
|||
} |
|||
} |
|||
TenantProfile tenantProfile = tenantProfileCache.get(tenantId); |
|||
tenantState = new TenantApiUsageState(tenantProfile, dbStateEntity); |
|||
try { |
|||
List<TsKvEntry> dbValues = tsService.findAllLatest(tenantId, dbStateEntity.getId()).get(); |
|||
for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { |
|||
boolean cycleEntryFound = false; |
|||
boolean hourlyEntryFound = false; |
|||
for (TsKvEntry tsKvEntry : dbValues) { |
|||
if (tsKvEntry.getKey().equals(key.getApiCountKey())) { |
|||
cycleEntryFound = true; |
|||
tenantState.put(key, tsKvEntry.getTs() == tenantState.getCurrentCycleTs() ? tsKvEntry.getLongValue().get() : 0L); |
|||
} else if (tsKvEntry.getKey().equals(key.getApiCountKey() + HOURLY)) { |
|||
hourlyEntryFound = true; |
|||
tenantState.putHourly(key, tsKvEntry.getTs() == tenantState.getCurrentHourTs() ? tsKvEntry.getLongValue().get() : 0L); |
|||
} |
|||
if (cycleEntryFound && hourlyEntryFound) { |
|||
break; |
|||
} |
|||
} |
|||
} |
|||
log.debug("[{}] Initialized state: {}", tenantId, dbStateEntity); |
|||
myTenantStates.put(tenantId, tenantState); |
|||
} catch (InterruptedException | ExecutionException e) { |
|||
log.warn("[{}] Failed to fetch api usage state from db.", tenantId, e); |
|||
} |
|||
} |
|||
return tenantState; |
|||
} |
|||
|
|||
private void initStatesFromDataBase() { |
|||
try { |
|||
PageDataIterable<Tenant> tenantIterator = new PageDataIterable<>(tenantService::findTenants, 1024); |
|||
for (Tenant tenant : tenantIterator) { |
|||
if (!myTenantStates.containsKey(tenant.getId()) && partitionService.resolve(ServiceType.TB_CORE, tenant.getId(), tenant.getId()).isMyPartition()) { |
|||
updateLock.lock(); |
|||
try { |
|||
updateTenantState(getOrFetchState(tenant.getId()), tenantProfileCache.get(tenant.getTenantProfileId())); |
|||
} catch (Exception e) { |
|||
log.warn("[{}] Failed to initialize tenant API state", tenant.getId(), e); |
|||
} finally { |
|||
updateLock.unlock(); |
|||
} |
|||
} |
|||
} |
|||
log.info("Api usage service started."); |
|||
} catch (Exception e) { |
|||
log.warn("Unknown failure", e); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,38 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.apiusage; |
|||
|
|||
import org.springframework.context.ApplicationListener; |
|||
import org.thingsboard.server.common.data.ApiUsageState; |
|||
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.gen.transport.TransportProtos.ToUsageStatsServiceMsg; |
|||
import org.thingsboard.server.queue.common.TbProtoQueueMsg; |
|||
import org.thingsboard.server.queue.discovery.PartitionChangeEvent; |
|||
|
|||
public interface TbApiUsageStateService extends ApplicationListener<PartitionChangeEvent> { |
|||
|
|||
void process(TbProtoQueueMsg<ToUsageStatsServiceMsg> msg, TbCallback callback); |
|||
|
|||
ApiUsageState getApiUsageState(TenantId tenantId); |
|||
|
|||
void onTenantProfileUpdate(TenantProfileId tenantProfileId); |
|||
|
|||
void onTenantUpdate(TenantId tenantId); |
|||
|
|||
void onApiUsageStateUpdate(TenantId tenantId); |
|||
} |
|||
@ -0,0 +1,201 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.apiusage; |
|||
|
|||
import lombok.Getter; |
|||
import lombok.Setter; |
|||
import org.springframework.data.util.Pair; |
|||
import org.thingsboard.server.common.data.ApiUsageRecordKey; |
|||
import org.thingsboard.server.common.data.ApiUsageState; |
|||
import org.thingsboard.server.common.data.TenantProfile; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.id.TenantProfileId; |
|||
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; |
|||
import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; |
|||
import org.thingsboard.server.common.msg.tools.SchedulerUtils; |
|||
|
|||
import java.util.HashMap; |
|||
import java.util.Map; |
|||
import java.util.concurrent.ConcurrentHashMap; |
|||
|
|||
public class TenantApiUsageState { |
|||
|
|||
private final Map<ApiUsageRecordKey, Long> currentCycleValues = new ConcurrentHashMap<>(); |
|||
private final Map<ApiUsageRecordKey, Long> currentHourValues = new ConcurrentHashMap<>(); |
|||
|
|||
@Getter |
|||
@Setter |
|||
private TenantProfileId tenantProfileId; |
|||
@Getter |
|||
@Setter |
|||
private TenantProfileData tenantProfileData; |
|||
@Getter |
|||
private final ApiUsageState apiUsageState; |
|||
@Getter |
|||
private volatile long currentCycleTs; |
|||
@Getter |
|||
private volatile long nextCycleTs; |
|||
@Getter |
|||
private volatile long currentHourTs; |
|||
|
|||
public TenantApiUsageState(TenantProfile tenantProfile, ApiUsageState apiUsageState) { |
|||
this.tenantProfileId = tenantProfile.getId(); |
|||
this.tenantProfileData = tenantProfile.getProfileData(); |
|||
this.apiUsageState = apiUsageState; |
|||
this.currentCycleTs = SchedulerUtils.getStartOfCurrentMonth(); |
|||
this.nextCycleTs = SchedulerUtils.getStartOfNextMonth(); |
|||
this.currentHourTs = SchedulerUtils.getStartOfCurrentHour(); |
|||
} |
|||
|
|||
public void put(ApiUsageRecordKey key, Long value) { |
|||
currentCycleValues.put(key, value); |
|||
} |
|||
|
|||
public void putHourly(ApiUsageRecordKey key, Long value) { |
|||
currentHourValues.put(key, value); |
|||
} |
|||
|
|||
public long add(ApiUsageRecordKey key, long value) { |
|||
long result = currentCycleValues.getOrDefault(key, 0L) + value; |
|||
currentCycleValues.put(key, result); |
|||
return result; |
|||
} |
|||
|
|||
public long get(ApiUsageRecordKey key) { |
|||
return currentCycleValues.getOrDefault(key, 0L); |
|||
} |
|||
|
|||
public long addToHourly(ApiUsageRecordKey key, long value) { |
|||
long result = currentHourValues.getOrDefault(key, 0L) + value; |
|||
currentHourValues.put(key, result); |
|||
return result; |
|||
} |
|||
|
|||
public void setHour(long currentHourTs) { |
|||
this.currentHourTs = currentHourTs; |
|||
for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { |
|||
currentHourValues.put(key, 0L); |
|||
} |
|||
} |
|||
|
|||
public void setCycles(long currentCycleTs, long nextCycleTs) { |
|||
this.currentCycleTs = currentCycleTs; |
|||
this.nextCycleTs = nextCycleTs; |
|||
for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { |
|||
currentCycleValues.put(key, 0L); |
|||
} |
|||
} |
|||
|
|||
public long getProfileThreshold(ApiUsageRecordKey key) { |
|||
return tenantProfileData.getConfiguration().getProfileThreshold(key); |
|||
} |
|||
|
|||
public TenantId getTenantId() { |
|||
return apiUsageState.getTenantId(); |
|||
} |
|||
|
|||
public boolean isTransportEnabled() { |
|||
return apiUsageState.isTransportEnabled(); |
|||
} |
|||
|
|||
public boolean isDbStorageEnabled() { |
|||
return apiUsageState.isDbStorageEnabled(); |
|||
} |
|||
|
|||
public boolean isRuleEngineEnabled() { |
|||
return apiUsageState.isReExecEnabled(); |
|||
} |
|||
|
|||
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.setReExecEnabled(ruleEngineEnabled); |
|||
} |
|||
|
|||
public void setJsExecEnabled(boolean jsExecEnabled) { |
|||
apiUsageState.setJsExecEnabled(jsExecEnabled); |
|||
} |
|||
|
|||
public boolean isFeatureEnabled(ApiUsageRecordKey recordKey) { |
|||
switch (recordKey) { |
|||
case TRANSPORT_MSG_COUNT: |
|||
case TRANSPORT_DP_COUNT: |
|||
return isTransportEnabled(); |
|||
case RE_EXEC_COUNT: |
|||
return isRuleEngineEnabled(); |
|||
case STORAGE_DP_COUNT: |
|||
return isDbStorageEnabled(); |
|||
case JS_EXEC_COUNT: |
|||
return isJsExecEnabled(); |
|||
default: |
|||
return true; |
|||
} |
|||
} |
|||
|
|||
public ApiFeature setFeatureValue(ApiUsageRecordKey recordKey, boolean value) { |
|||
ApiFeature feature = null; |
|||
boolean currentValue = isFeatureEnabled(recordKey); |
|||
switch (recordKey) { |
|||
case TRANSPORT_MSG_COUNT: |
|||
case TRANSPORT_DP_COUNT: |
|||
feature = ApiFeature.TRANSPORT; |
|||
setTransportEnabled(value); |
|||
break; |
|||
case RE_EXEC_COUNT: |
|||
feature = ApiFeature.RE; |
|||
setRuleEngineEnabled(value); |
|||
break; |
|||
case STORAGE_DP_COUNT: |
|||
feature = ApiFeature.DB; |
|||
setDbStorageEnabled(value); |
|||
break; |
|||
case JS_EXEC_COUNT: |
|||
feature = ApiFeature.JS; |
|||
setJsExecEnabled(value); |
|||
break; |
|||
} |
|||
return currentValue == value ? null : feature; |
|||
} |
|||
|
|||
public Map<ApiFeature, Boolean> checkStateUpdatedDueToThresholds() { |
|||
Map<ApiFeature, Boolean> result = new HashMap<>(); |
|||
for (ApiUsageRecordKey key : ApiUsageRecordKey.values()) { |
|||
Pair<ApiFeature, Boolean> featureUpdate = checkStateUpdatedDueToThreshold(key); |
|||
if (featureUpdate != null) { |
|||
result.put(featureUpdate.getFirst(), featureUpdate.getSecond()); |
|||
} |
|||
} |
|||
return result; |
|||
} |
|||
|
|||
public Pair<ApiFeature, Boolean> checkStateUpdatedDueToThreshold(ApiUsageRecordKey recordKey) { |
|||
long value = get(recordKey); |
|||
long threshold = getProfileThreshold(recordKey); |
|||
boolean featureValue = threshold == 0 || value < threshold; |
|||
ApiFeature feature = setFeatureValue(recordKey, featureValue); |
|||
return feature != null ? Pair.of(feature, featureValue) : null; |
|||
} |
|||
} |
|||
@ -0,0 +1,46 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.telemetry; |
|||
|
|||
import com.google.common.util.concurrent.FutureCallback; |
|||
import org.thingsboard.rule.engine.api.RuleEngineTelemetryService; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|||
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|||
|
|||
import java.util.List; |
|||
|
|||
/** |
|||
* Created by ashvayka on 27.03.18. |
|||
*/ |
|||
public interface InternalTelemetryService extends RuleEngineTelemetryService { |
|||
|
|||
void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, FutureCallback<Integer> callback); |
|||
|
|||
void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, long ttl, FutureCallback<Integer> callback); |
|||
|
|||
void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback); |
|||
|
|||
void saveLatestAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, FutureCallback<Void> callback); |
|||
|
|||
void deleteAndNotifyInternal(TenantId tenantId, EntityId entityId, String scope, List<String> keys, FutureCallback<Void> callback); |
|||
|
|||
void deleteLatestInternal(TenantId tenantId, EntityId entityId, List<String> keys, FutureCallback<Void> callback); |
|||
|
|||
|
|||
|
|||
} |
|||
@ -1,250 +0,0 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.script; |
|||
|
|||
import com.datastax.oss.driver.api.core.uuid.Uuids; |
|||
import com.google.common.collect.Sets; |
|||
import org.junit.After; |
|||
import org.junit.Before; |
|||
import org.junit.Test; |
|||
import org.thingsboard.rule.engine.api.ScriptEngine; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.RuleNodeId; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
import org.thingsboard.server.common.msg.TbMsgDataType; |
|||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|||
|
|||
import javax.script.ScriptException; |
|||
import java.util.Map; |
|||
import java.util.Set; |
|||
import java.util.UUID; |
|||
import java.util.concurrent.*; |
|||
import java.util.concurrent.atomic.AtomicInteger; |
|||
|
|||
import static org.junit.Assert.*; |
|||
|
|||
public class RuleNodeJsScriptEngineTest { |
|||
|
|||
private ScriptEngine scriptEngine; |
|||
private TestNashornJsInvokeService jsSandboxService; |
|||
|
|||
private EntityId ruleNodeId = new RuleNodeId(Uuids.timeBased()); |
|||
|
|||
@Before |
|||
public void beforeTest() throws Exception { |
|||
jsSandboxService = new TestNashornJsInvokeService(false, 1, 100, 3); |
|||
} |
|||
|
|||
@After |
|||
public void afterTest() throws Exception { |
|||
jsSandboxService.stop(); |
|||
} |
|||
|
|||
@Test |
|||
public void msgCanBeUpdated() throws ScriptException { |
|||
String function = "metadata.temp = metadata.temp * 10; return {metadata: metadata};"; |
|||
scriptEngine = new RuleNodeJsScriptEngine(jsSandboxService, ruleNodeId, function); |
|||
|
|||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|||
metaData.putValue("temp", "7"); |
|||
metaData.putValue("humidity", "99"); |
|||
String rawJson = "{\"name\": \"Vit\", \"passed\": 5, \"bigObj\": {\"prop\":42}}"; |
|||
|
|||
TbMsg msg = TbMsg.newMsg( "USER", null, metaData, TbMsgDataType.JSON, rawJson); |
|||
|
|||
TbMsg actual = scriptEngine.executeUpdate(msg); |
|||
assertEquals("70", actual.getMetaData().getValue("temp")); |
|||
scriptEngine.destroy(); |
|||
} |
|||
|
|||
@Test |
|||
public void newAttributesCanBeAddedInMsg() throws ScriptException { |
|||
String function = "metadata.newAttr = metadata.humidity - msg.passed; return {metadata: metadata};"; |
|||
scriptEngine = new RuleNodeJsScriptEngine(jsSandboxService, ruleNodeId, function); |
|||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|||
metaData.putValue("temp", "7"); |
|||
metaData.putValue("humidity", "99"); |
|||
String rawJson = "{\"name\": \"Vit\", \"passed\": 5, \"bigObj\": {\"prop\":42}}"; |
|||
|
|||
TbMsg msg = TbMsg.newMsg( "USER", null, metaData, TbMsgDataType.JSON, rawJson); |
|||
|
|||
TbMsg actual = scriptEngine.executeUpdate(msg); |
|||
assertEquals("94", actual.getMetaData().getValue("newAttr")); |
|||
scriptEngine.destroy(); |
|||
} |
|||
|
|||
@Test |
|||
public void payloadCanBeUpdated() throws ScriptException { |
|||
String function = "msg.passed = msg.passed * metadata.temp; msg.bigObj.newProp = 'Ukraine'; return {msg: msg};"; |
|||
scriptEngine = new RuleNodeJsScriptEngine(jsSandboxService, ruleNodeId, function); |
|||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|||
metaData.putValue("temp", "7"); |
|||
metaData.putValue("humidity", "99"); |
|||
String rawJson = "{\"name\":\"Vit\",\"passed\": 5,\"bigObj\":{\"prop\":42}}"; |
|||
|
|||
TbMsg msg =TbMsg.newMsg("USER", null, metaData, TbMsgDataType.JSON, rawJson); |
|||
|
|||
TbMsg actual = scriptEngine.executeUpdate(msg); |
|||
|
|||
String expectedJson = "{\"name\":\"Vit\",\"passed\":35,\"bigObj\":{\"prop\":42,\"newProp\":\"Ukraine\"}}"; |
|||
assertEquals(expectedJson, actual.getData()); |
|||
scriptEngine.destroy(); |
|||
} |
|||
|
|||
@Test |
|||
public void metadataAccessibleForFilter() throws ScriptException { |
|||
String function = "return metadata.humidity < 15;"; |
|||
scriptEngine = new RuleNodeJsScriptEngine(jsSandboxService, ruleNodeId, function); |
|||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|||
metaData.putValue("temp", "7"); |
|||
metaData.putValue("humidity", "99"); |
|||
String rawJson = "{\"name\": \"Vit\", \"passed\": 5, \"bigObj\": {\"prop\":42}}"; |
|||
|
|||
TbMsg msg = TbMsg.newMsg("USER", null, metaData, TbMsgDataType.JSON, rawJson); |
|||
assertFalse(scriptEngine.executeFilter(msg)); |
|||
scriptEngine.destroy(); |
|||
} |
|||
|
|||
@Test |
|||
public void dataAccessibleForFilter() throws ScriptException { |
|||
String function = "return msg.passed < 15 && msg.name === 'Vit' && metadata.temp == 7 && msg.bigObj.prop == 42;"; |
|||
scriptEngine = new RuleNodeJsScriptEngine(jsSandboxService, ruleNodeId, function); |
|||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|||
metaData.putValue("temp", "7"); |
|||
metaData.putValue("humidity", "99"); |
|||
String rawJson = "{\"name\": \"Vit\", \"passed\": 5, \"bigObj\": {\"prop\":42}}"; |
|||
|
|||
TbMsg msg = TbMsg.newMsg( "USER", null, metaData,TbMsgDataType.JSON, rawJson); |
|||
assertTrue(scriptEngine.executeFilter(msg)); |
|||
scriptEngine.destroy(); |
|||
} |
|||
|
|||
@Test |
|||
public void dataAccessibleForSwitch() throws ScriptException { |
|||
String jsCode = "function nextRelation(metadata, msg) {\n" + |
|||
" if(msg.passed == 5 && metadata.temp == 10)\n" + |
|||
" return 'one'\n" + |
|||
" else\n" + |
|||
" return 'two';\n" + |
|||
"};\n" + |
|||
"\n" + |
|||
"return nextRelation(metadata, msg);"; |
|||
scriptEngine = new RuleNodeJsScriptEngine(jsSandboxService, ruleNodeId, jsCode); |
|||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|||
metaData.putValue("temp", "10"); |
|||
metaData.putValue("humidity", "99"); |
|||
String rawJson = "{\"name\": \"Vit\", \"passed\": 5, \"bigObj\": {\"prop\":42}}"; |
|||
|
|||
TbMsg msg = TbMsg.newMsg( "USER", null, metaData, TbMsgDataType.JSON, rawJson); |
|||
Set<String> actual = scriptEngine.executeSwitch(msg); |
|||
assertEquals(Sets.newHashSet("one"), actual); |
|||
scriptEngine.destroy(); |
|||
} |
|||
|
|||
@Test |
|||
public void multipleRelationsReturnedFromSwitch() throws ScriptException { |
|||
String jsCode = "function nextRelation(metadata, msg) {\n" + |
|||
" if(msg.passed == 5 && metadata.temp == 10)\n" + |
|||
" return ['three', 'one']\n" + |
|||
" else\n" + |
|||
" return 'two';\n" + |
|||
"};\n" + |
|||
"\n" + |
|||
"return nextRelation(metadata, msg);"; |
|||
scriptEngine = new RuleNodeJsScriptEngine(jsSandboxService, ruleNodeId, jsCode); |
|||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|||
metaData.putValue("temp", "10"); |
|||
metaData.putValue("humidity", "99"); |
|||
String rawJson = "{\"name\": \"Vit\", \"passed\": 5, \"bigObj\": {\"prop\":42}}"; |
|||
|
|||
TbMsg msg = TbMsg.newMsg( "USER", null, metaData, TbMsgDataType.JSON, rawJson); |
|||
Set<String> actual = scriptEngine.executeSwitch(msg); |
|||
assertEquals(Sets.newHashSet("one", "three"), actual); |
|||
scriptEngine.destroy(); |
|||
} |
|||
|
|||
@Test |
|||
public void concurrentReleasedCorrectly() throws InterruptedException, ExecutionException { |
|||
String code = "metadata.temp = metadata.temp * 10; return {metadata: metadata};"; |
|||
|
|||
int repeat = 1000; |
|||
ExecutorService service = Executors.newFixedThreadPool(repeat); |
|||
Map<UUID, Object> scriptIds = new ConcurrentHashMap<>(); |
|||
CountDownLatch startLatch = new CountDownLatch(repeat); |
|||
CountDownLatch finishLatch = new CountDownLatch(repeat); |
|||
AtomicInteger failedCount = new AtomicInteger(0); |
|||
|
|||
for (int i = 0; i < repeat; i++) { |
|||
service.submit(() -> runScript(startLatch, finishLatch, failedCount, scriptIds, code)); |
|||
} |
|||
|
|||
finishLatch.await(); |
|||
assertTrue(scriptIds.size() == 1); |
|||
assertTrue(failedCount.get() == 0); |
|||
|
|||
CountDownLatch nextStart = new CountDownLatch(repeat); |
|||
CountDownLatch nextFinish = new CountDownLatch(repeat); |
|||
for (int i = 0; i < repeat; i++) { |
|||
service.submit(() -> runScript(nextStart, nextFinish, failedCount, scriptIds, code)); |
|||
} |
|||
|
|||
nextFinish.await(); |
|||
assertTrue(scriptIds.size() == 1); |
|||
assertTrue(failedCount.get() == 0); |
|||
service.shutdownNow(); |
|||
} |
|||
|
|||
@Test |
|||
public void concurrentFailedEvaluationShouldThrowException() throws InterruptedException { |
|||
String code = "metadata.temp = metadata.temp * 10; urn {metadata: metadata};"; |
|||
|
|||
int repeat = 10000; |
|||
ExecutorService service = Executors.newFixedThreadPool(repeat); |
|||
Map<UUID, Object> scriptIds = new ConcurrentHashMap<>(); |
|||
CountDownLatch startLatch = new CountDownLatch(repeat); |
|||
CountDownLatch finishLatch = new CountDownLatch(repeat); |
|||
AtomicInteger failedCount = new AtomicInteger(0); |
|||
for (int i = 0; i < repeat; i++) { |
|||
service.submit(() -> { |
|||
service.submit(() -> runScript(startLatch, finishLatch, failedCount, scriptIds, code)); |
|||
}); |
|||
} |
|||
|
|||
finishLatch.await(); |
|||
assertTrue(scriptIds.isEmpty()); |
|||
assertEquals(repeat, failedCount.get()); |
|||
service.shutdownNow(); |
|||
} |
|||
|
|||
private void runScript(CountDownLatch startLatch, CountDownLatch finishLatch, AtomicInteger failedCount, |
|||
Map<UUID, Object> scriptIds, String code) { |
|||
try { |
|||
for (int k = 0; k < 10; k++) { |
|||
startLatch.countDown(); |
|||
startLatch.await(); |
|||
UUID scriptId = jsSandboxService.eval(JsScriptType.RULE_NODE_SCRIPT, code).get(); |
|||
scriptIds.put(scriptId, new Object()); |
|||
jsSandboxService.invokeFunction(scriptId, "{}", "{}", "TEXT").get(); |
|||
jsSandboxService.release(scriptId).get(); |
|||
} |
|||
} catch (Throwable th) { |
|||
failedCount.incrementAndGet(); |
|||
} finally { |
|||
finishLatch.countDown(); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -1,57 +0,0 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.script; |
|||
|
|||
public class TestNashornJsInvokeService extends AbstractNashornJsInvokeService { |
|||
|
|||
private boolean useJsSandbox; |
|||
private final int monitorThreadPoolSize; |
|||
private final long maxCpuTime; |
|||
private final int maxErrors; |
|||
|
|||
public TestNashornJsInvokeService(boolean useJsSandbox, int monitorThreadPoolSize, long maxCpuTime, int maxErrors) { |
|||
this.useJsSandbox = useJsSandbox; |
|||
this.monitorThreadPoolSize = monitorThreadPoolSize; |
|||
this.maxCpuTime = maxCpuTime; |
|||
this.maxErrors = maxErrors; |
|||
init(); |
|||
} |
|||
|
|||
@Override |
|||
protected boolean useJsSandbox() { |
|||
return useJsSandbox; |
|||
} |
|||
|
|||
@Override |
|||
protected int getMonitorThreadPoolSize() { |
|||
return monitorThreadPoolSize; |
|||
} |
|||
|
|||
@Override |
|||
protected long getMaxCpuTime() { |
|||
return maxCpuTime; |
|||
} |
|||
|
|||
@Override |
|||
protected int getMaxErrors() { |
|||
return maxErrors; |
|||
} |
|||
|
|||
@Override |
|||
protected long getMaxBlacklistDuration() { |
|||
return 100000; |
|||
} |
|||
} |
|||
@ -0,0 +1,33 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.dao.usagerecord; |
|||
|
|||
import org.thingsboard.server.common.data.ApiUsageState; |
|||
import org.thingsboard.server.common.data.id.ApiUsageStateId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
|
|||
public interface ApiUsageStateService { |
|||
|
|||
ApiUsageState createDefaultApiUsageState(TenantId id); |
|||
|
|||
ApiUsageState update(ApiUsageState apiUsageState); |
|||
|
|||
ApiUsageState findTenantApiUsageState(TenantId tenantId); |
|||
|
|||
void deleteApiUsageStateByTenantId(TenantId tenantId); |
|||
|
|||
ApiUsageState findApiUsageStateById(TenantId tenantId, ApiUsageStateId id); |
|||
} |
|||
@ -0,0 +1,38 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.common.data; |
|||
|
|||
import lombok.Getter; |
|||
|
|||
public enum ApiUsageRecordKey { |
|||
|
|||
TRANSPORT_MSG_COUNT("transportMsgCount", "transportMsgLimit"), |
|||
TRANSPORT_DP_COUNT("transportDataPointsCount", "transportDataPointsLimit"), |
|||
STORAGE_DP_COUNT("storageDataPointsCount", "storageDataPointsLimit"), |
|||
RE_EXEC_COUNT("ruleEngineExecutionCount", "ruleEngineExecutionLimit"), |
|||
JS_EXEC_COUNT("jsExecutionCount", "jsExecutionLimit"); |
|||
|
|||
@Getter |
|||
private final String apiCountKey; |
|||
@Getter |
|||
private final String apiLimitKey; |
|||
|
|||
ApiUsageRecordKey(String apiCountKey, String apiLimitKey) { |
|||
this.apiCountKey = apiCountKey; |
|||
this.apiLimitKey = apiLimitKey; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,68 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.common.data; |
|||
|
|||
import lombok.EqualsAndHashCode; |
|||
import lombok.Getter; |
|||
import lombok.Setter; |
|||
import lombok.ToString; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.id.ApiUsageStateId; |
|||
|
|||
@ToString |
|||
@EqualsAndHashCode(callSuper = true) |
|||
public class ApiUsageState extends BaseData<ApiUsageStateId> implements HasTenantId { |
|||
|
|||
private static final long serialVersionUID = 8250339805336035966L; |
|||
|
|||
@Getter |
|||
@Setter |
|||
private TenantId tenantId; |
|||
@Getter |
|||
@Setter |
|||
private EntityId entityId; |
|||
@Getter |
|||
@Setter |
|||
private boolean transportEnabled = true; |
|||
@Getter |
|||
@Setter |
|||
private boolean dbStorageEnabled = true; |
|||
@Getter |
|||
@Setter |
|||
private boolean reExecEnabled = true; |
|||
@Getter |
|||
@Setter |
|||
private boolean jsExecEnabled = true; |
|||
|
|||
public ApiUsageState() { |
|||
super(); |
|||
} |
|||
|
|||
public ApiUsageState(ApiUsageStateId id) { |
|||
super(id); |
|||
} |
|||
|
|||
public ApiUsageState(ApiUsageState ur) { |
|||
super(ur); |
|||
this.tenantId = ur.getTenantId(); |
|||
this.entityId = ur.getEntityId(); |
|||
this.transportEnabled = ur.isTransportEnabled(); |
|||
this.dbStorageEnabled = ur.isDbStorageEnabled(); |
|||
this.reExecEnabled = ur.isReExecEnabled(); |
|||
this.jsExecEnabled = ur.isJsExecEnabled(); |
|||
} |
|||
} |
|||
@ -0,0 +1,42 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.common.data.id; |
|||
|
|||
import com.fasterxml.jackson.annotation.JsonCreator; |
|||
import com.fasterxml.jackson.annotation.JsonIgnore; |
|||
import com.fasterxml.jackson.annotation.JsonProperty; |
|||
import org.thingsboard.server.common.data.EntityType; |
|||
|
|||
import java.util.UUID; |
|||
|
|||
public class ApiUsageStateId extends UUIDBased implements EntityId { |
|||
|
|||
@JsonCreator |
|||
public ApiUsageStateId(@JsonProperty("id") UUID id) { |
|||
super(id); |
|||
} |
|||
|
|||
public static ApiUsageStateId fromString(String userId) { |
|||
return new ApiUsageStateId(UUID.fromString(userId)); |
|||
} |
|||
|
|||
@JsonIgnore |
|||
@Override |
|||
public EntityType getEntityType() { |
|||
return EntityType.API_USAGE_STATE; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,28 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.common.data.query; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
|
|||
@Data |
|||
public class ApiUsageStateFilter implements EntityFilter { |
|||
@Override |
|||
public EntityFilterType getType() { |
|||
return EntityFilterType.API_USAGE_STATE; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,69 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.common.data.tenant.profile; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.ApiUsageRecordKey; |
|||
import org.thingsboard.server.common.data.TenantProfileType; |
|||
|
|||
@Data |
|||
public class DefaultTenantProfileConfiguration implements TenantProfileConfiguration { |
|||
|
|||
private long maxDevices; |
|||
private long maxAssets; |
|||
|
|||
private String transportTenantMsgRateLimit; |
|||
private String transportTenantTelemetryMsgRateLimit; |
|||
private String transportTenantTelemetryDataPointsRateLimit; |
|||
private String transportDeviceMsgRateLimit; |
|||
private String transportDeviceTelemetryMsgRateLimit; |
|||
private String transportDeviceTelemetryDataPointsRateLimit; |
|||
|
|||
private long maxTransportMessages; |
|||
private long maxTransportDataPoints; |
|||
private long maxREExecutions; |
|||
private long maxJSExecutions; |
|||
private long maxDPStorageDays; |
|||
private int maxRuleNodeExecutionsPerMessage; |
|||
|
|||
@Override |
|||
public long getProfileThreshold(ApiUsageRecordKey key) { |
|||
switch (key) { |
|||
case TRANSPORT_MSG_COUNT: |
|||
return maxTransportMessages; |
|||
case TRANSPORT_DP_COUNT: |
|||
return maxTransportDataPoints; |
|||
case JS_EXEC_COUNT: |
|||
return maxJSExecutions; |
|||
case RE_EXEC_COUNT: |
|||
return maxREExecutions; |
|||
case STORAGE_DP_COUNT: |
|||
return maxDPStorageDays; |
|||
} |
|||
return 0L; |
|||
} |
|||
|
|||
|
|||
@Override |
|||
public TenantProfileType getType() { |
|||
return TenantProfileType.DEFAULT; |
|||
} |
|||
|
|||
@Override |
|||
public int getMaxRuleNodeExecsPerMessage() { |
|||
return maxRuleNodeExecutionsPerMessage; |
|||
} |
|||
} |
|||
@ -0,0 +1,43 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.common.data.tenant.profile; |
|||
|
|||
import com.fasterxml.jackson.annotation.JsonIgnore; |
|||
import com.fasterxml.jackson.annotation.JsonIgnoreProperties; |
|||
import com.fasterxml.jackson.annotation.JsonSubTypes; |
|||
import com.fasterxml.jackson.annotation.JsonTypeInfo; |
|||
import org.thingsboard.server.common.data.ApiUsageRecordKey; |
|||
import org.thingsboard.server.common.data.TenantProfileType; |
|||
|
|||
@JsonIgnoreProperties(ignoreUnknown = true) |
|||
@JsonTypeInfo( |
|||
use = JsonTypeInfo.Id.NAME, |
|||
include = JsonTypeInfo.As.PROPERTY, |
|||
property = "type") |
|||
@JsonSubTypes({ |
|||
@JsonSubTypes.Type(value = DefaultTenantProfileConfiguration.class, name = "DEFAULT")}) |
|||
public interface TenantProfileConfiguration { |
|||
|
|||
@JsonIgnore |
|||
TenantProfileType getType(); |
|||
|
|||
@JsonIgnore |
|||
long getProfileThreshold(ApiUsageRecordKey key); |
|||
|
|||
@JsonIgnore |
|||
int getMaxRuleNodeExecsPerMessage(); |
|||
|
|||
} |
|||
@ -0,0 +1,45 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.queue.common; |
|||
|
|||
import org.thingsboard.server.common.msg.queue.RuleEngineException; |
|||
import org.thingsboard.server.queue.TbQueueCallback; |
|||
import org.thingsboard.server.queue.TbQueueMsgMetadata; |
|||
|
|||
import java.util.concurrent.atomic.AtomicInteger; |
|||
|
|||
public class MultipleTbQueueCallbackWrapper implements TbQueueCallback { |
|||
|
|||
private final AtomicInteger tbQueueCallbackCount; |
|||
private final TbQueueCallback callback; |
|||
|
|||
public MultipleTbQueueCallbackWrapper(int tbQueueCallbackCount, TbQueueCallback callback) { |
|||
this.tbQueueCallbackCount = new AtomicInteger(tbQueueCallbackCount); |
|||
this.callback = callback; |
|||
} |
|||
|
|||
@Override |
|||
public void onSuccess(TbQueueMsgMetadata metadata) { |
|||
if (tbQueueCallbackCount.decrementAndGet() <= 0) { |
|||
callback.onSuccess(metadata); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void onFailure(Throwable t) { |
|||
callback.onFailure(new RuleEngineException(t.getMessage())); |
|||
} |
|||
} |
|||
@ -0,0 +1,26 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.queue.provider; |
|||
|
|||
import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg; |
|||
import org.thingsboard.server.queue.TbQueueProducer; |
|||
import org.thingsboard.server.queue.common.TbProtoQueueMsg; |
|||
|
|||
public interface TbUsageStatsClientQueueFactory { |
|||
|
|||
TbQueueProducer<TbProtoQueueMsg<ToUsageStatsServiceMsg>> createToUsageStatsServiceMsgProducer(); |
|||
|
|||
} |
|||
@ -0,0 +1,61 @@ |
|||
/** |
|||
* Copyright © 2016-2020 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.queue.scheduler; |
|||
|
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|||
|
|||
import javax.annotation.PostConstruct; |
|||
import javax.annotation.PreDestroy; |
|||
import java.util.concurrent.Callable; |
|||
import java.util.concurrent.Executors; |
|||
import java.util.concurrent.ScheduledExecutorService; |
|||
import java.util.concurrent.ScheduledFuture; |
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
@Component |
|||
public class DefaultSchedulerComponent implements SchedulerComponent{ |
|||
|
|||
protected ScheduledExecutorService schedulerExecutor; |
|||
|
|||
@PostConstruct |
|||
public void init(){ |
|||
this.schedulerExecutor = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("queue-scheduler")); |
|||
} |
|||
|
|||
@PreDestroy |
|||
public void destroy() { |
|||
if (schedulerExecutor != null) { |
|||
schedulerExecutor.shutdownNow(); |
|||
} |
|||
} |
|||
|
|||
public ScheduledFuture<?> schedule(Runnable command, long delay, TimeUnit unit) { |
|||
return schedulerExecutor.schedule(command, delay, unit); |
|||
} |
|||
|
|||
public <V> ScheduledFuture<V> schedule(Callable<V> callable, long delay, TimeUnit unit) { |
|||
return schedulerExecutor.schedule(callable, delay, unit); |
|||
} |
|||
|
|||
public ScheduledFuture<?> scheduleAtFixedRate(Runnable command, long initialDelay, long period, TimeUnit unit) { |
|||
return schedulerExecutor.scheduleAtFixedRate(command, initialDelay, period, unit); |
|||
} |
|||
|
|||
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable command, long initialDelay, long delay, TimeUnit unit) { |
|||
return schedulerExecutor.scheduleWithFixedDelay(command, initialDelay, delay, unit); |
|||
} |
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue