diff --git a/application/src/main/data/json/system/widget_bundles/input_widgets.json b/application/src/main/data/json/system/widget_bundles/input_widgets.json
index ac3486223a..097c31dac4 100644
--- a/application/src/main/data/json/system/widget_bundles/input_widgets.json
+++ b/application/src/main/data/json/system/widget_bundles/input_widgets.json
@@ -391,16 +391,16 @@
},
{
"alias": "web_camera_input",
- "name": "Web Camera Input",
+ "name": "Photo camera input",
"descriptor": {
"type": "latest",
"sizeX": 7.5,
"sizeY": 3,
"resources": [],
- "templateHtml": "\n",
+ "templateHtml": "\n",
"templateCss": "",
- "controllerScript": "self.onInit = function() {\n}\n\nself.onDataUpdated = function() {\n self.ctx.$scope.webCameraInputWidget.onDataUpdated();\n}\n\nself.typeParameters = function() {\n return {\n maxDatasources: 1,\n maxDataKeys: 1,\n singleEntity: true\n }\n}\n\nself.onDestroy = function() {\n}\n",
- "settingsSchema": "{\n \"schema\": {\n \"type\": \"object\",\n \"title\": \"Web Camera\",\n \"properties\": {\n \"widgetTitle\": {\n \"title\": \"Widget title\",\n \"type\": \"string\",\n \"default\": \"\"\n },\n \"imageFormat\": {\n \"title\": \"Image Format\",\n \"type\": \"string\",\n \"default\": \"image/png\"\n },\n \"imageQuality\":{\n \"title\":\"Image quality that use lossy compression such as jpeg and webp\",\n \"type\":\"number\",\n \"default\": 0.92,\n \"min\": 0,\n \"max\": 1\n },\n \"maxWidth\": {\n \"title\": \"The maximal image width\",\n \"type\": \"number\",\n \"default\": 640\n }, \n \"maxHeight\": {\n \"title\": \"The maximal image heigth\",\n \"type\": \"number\",\n \"default\": 480\n }\n },\n \"required\": []\n },\n \"form\": [\n \"widgetTitle\",\n {\n \"key\": \"imageFormat\",\n \"type\": \"rc-select\",\n \"multiple\": false,\n \"items\": [\n {\n \"value\": \"image/jpeg\",\n \"label\": \"JPEG\"\n },\n {\n \"value\": \"image/png\",\n \"label\": \"PNG\"\n },\n {\n \"value\": \"image/webp\",\n \"label\": \"WEBP\"\n }\n ]\n },\n \"imageQuality\",\n \"maxWidth\",\n \"maxHeight\"\n ]\n}",
+ "controllerScript": "self.onInit = function() {\n}\n\nself.onDataUpdated = function() {\n self.ctx.$scope.photoCameraInputWidget.onDataUpdated();\n}\n\nself.typeParameters = function() {\n return {\n maxDatasources: 1,\n maxDataKeys: 1,\n singleEntity: true\n }\n}\n\nself.onDestroy = function() {\n}\n",
+ "settingsSchema": "{\n \"schema\": {\n \"type\": \"object\",\n \"title\": \"Photo Camera\",\n \"properties\": {\n \"widgetTitle\": {\n \"title\": \"Widget title\",\n \"type\": \"string\",\n \"default\": \"\"\n },\n \"imageFormat\": {\n \"title\": \"Image Format\",\n \"type\": \"string\",\n \"default\": \"image/png\"\n },\n \"imageQuality\":{\n \"title\":\"Image quality that use lossy compression such as jpeg and webp\",\n \"type\":\"number\",\n \"default\": 0.92,\n \"min\": 0,\n \"max\": 1\n },\n \"maxWidth\": {\n \"title\": \"The maximal image width\",\n \"type\": \"number\",\n \"default\": 640\n }, \n \"maxHeight\": {\n \"title\": \"The maximal image heigth\",\n \"type\": \"number\",\n \"default\": 480\n }\n },\n \"required\": []\n },\n \"form\": [\n \"widgetTitle\",\n {\n \"key\": \"imageFormat\",\n \"type\": \"rc-select\",\n \"multiple\": false,\n \"items\": [\n {\n \"value\": \"image/jpeg\",\n \"label\": \"JPEG\"\n },\n {\n \"value\": \"image/png\",\n \"label\": \"PNG\"\n },\n {\n \"value\": \"image/webp\",\n \"label\": \"WEBP\"\n }\n ]\n },\n \"imageQuality\",\n \"maxWidth\",\n \"maxHeight\"\n ]\n}",
"dataKeySettingsSchema": "{}\n",
"defaultConfig": "{\"datasources\":[{\"type\":\"function\",\"name\":\"function\",\"dataKeys\":[{\"name\":\"f(x)\",\"type\":\"function\",\"label\":\"Random\",\"color\":\"#2196f3\",\"settings\":{},\"_hash\":0.15479322438769105,\"funcBody\":\"var value = prevValue + Math.random() * 100 - 50;\\nvar multiplier = Math.pow(10, 2 || 0);\\nvar value = Math.round(value * multiplier) / multiplier;\\nif (value < -1000) {\\n\\tvalue = -1000;\\n} else if (value > 1000) {\\n\\tvalue = 1000;\\n}\\nreturn value;\"}]}],\"timewindow\":{\"realtime\":{\"timewindowMs\":60000}},\"showTitle\":true,\"backgroundColor\":\"#fff\",\"color\":\"rgba(0, 0, 0, 0.87)\",\"padding\":\"8px\",\"settings\":{},\"title\":\"Web Camera Input\",\"showTitleIcon\":false,\"titleIcon\":\"more_horiz\",\"iconColor\":\"rgba(0, 0, 0, 0.87)\",\"iconSize\":\"24px\",\"titleTooltip\":\"\",\"dropShadow\":true,\"enableFullscreen\":false,\"widgetStyle\":{},\"titleStyle\":{\"fontSize\":\"16px\",\"fontWeight\":400},\"useDashboardTimewindow\":true,\"displayTimewindow\":true,\"showLegend\":false,\"actions\":{}}"
}
diff --git a/application/src/main/data/upgrade/3.1.1/schema_update_before.sql b/application/src/main/data/upgrade/3.1.1/schema_update_before.sql
index 216940b8f0..5f30f135e1 100644
--- a/application/src/main/data/upgrade/3.1.1/schema_update_before.sql
+++ b/application/src/main/data/upgrade/3.1.1/schema_update_before.sql
@@ -96,6 +96,7 @@ CREATE TABLE IF NOT EXISTS device_profile (
is_default boolean,
tenant_id uuid,
default_rule_chain_id uuid,
+ default_queue_name varchar(255),
provision_device_key varchar,
CONSTRAINT device_profile_name_unq_key UNIQUE (tenant_id, name),
CONSTRAINT device_provision_key_unq_key UNIQUE (provision_device_key),
diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
index 7f028932f2..52dbb38c9c 100644
--- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
+++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
@@ -32,7 +32,6 @@ import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.thingsboard.rule.engine.api.MailService;
-import org.thingsboard.rule.engine.api.RuleEngineDeviceProfileCache;
import org.thingsboard.server.actors.service.ActorService;
import org.thingsboard.server.actors.tenant.DebugTbRateLimits;
import org.thingsboard.server.common.data.DataConstants;
@@ -65,6 +64,8 @@ import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.user.UserService;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
+import org.thingsboard.server.queue.usagestats.TbApiUsageClient;
+import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.component.ComponentDiscoveryService;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.service.executors.DbCallbackExecutorService;
@@ -72,7 +73,7 @@ import org.thingsboard.server.service.executors.ExternalCallExecutorService;
import org.thingsboard.server.service.executors.SharedEventLoopGroupService;
import org.thingsboard.server.service.mail.MailExecutorService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
-import org.thingsboard.server.service.profile.TbTenantProfileCache;
+import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.queue.TbClusterService;
import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService;
import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService;
@@ -106,6 +107,14 @@ public class ActorSystemContext {
return debugPerTenantLimits;
}
+ @Autowired
+ @Getter
+ private TbApiUsageStateService apiUsageStateService;
+
+ @Autowired
+ @Getter
+ private TbApiUsageClient apiUsageClient;
+
@Autowired
@Getter
@Setter
diff --git a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
index 826596db56..ba8654854c 100644
--- a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
+++ b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
@@ -30,7 +30,6 @@ import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
-import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.page.PageDataIterable;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.msg.MsgType;
@@ -40,10 +39,8 @@ import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg;
import org.thingsboard.server.common.msg.queue.RuleEngineException;
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.service.profile.TbTenantProfileCache;
+import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper;
import java.util.HashSet;
@@ -150,15 +147,12 @@ public class AppActor extends ContextAwareActor {
private void onComponentLifecycleMsg(ComponentLifecycleMsg msg) {
TbActorRef target = null;
if (TenantId.SYS_TENANT_ID.equals(msg.getTenantId())) {
- if (msg.getEntityId().getEntityType() == EntityType.TENANT_PROFILE) {
- tenantProfileCache.evict(new TenantProfileId(msg.getEntityId().getId()));
- } else {
+ if (!EntityType.TENANT_PROFILE.equals(msg.getEntityId().getEntityType())) {
log.warn("Message has system tenant id: {}", msg);
}
} else {
- if (msg.getEntityId().getEntityType() == EntityType.TENANT) {
+ if (EntityType.TENANT.equals(msg.getEntityId().getEntityType())) {
TenantId tenantId = new TenantId(msg.getEntityId().getId());
- tenantProfileCache.evict(tenantId);
if (msg.getEvent() == ComponentLifecycleEvent.DELETED) {
log.info("[{}] Handling tenant deleted notification: {}", msg.getTenantId(), msg);
deletedTenants.add(tenantId);
diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
index c8bb083577..6749ed69ed 100644
--- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
+++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
@@ -38,6 +38,7 @@ import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.asset.Asset;
+import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.RuleNodeId;
@@ -73,6 +74,7 @@ import org.thingsboard.server.service.script.RuleNodeJsScriptEngine;
import java.util.Collections;
import java.util.Set;
+import java.util.function.BiConsumer;
import java.util.function.Consumer;
/**
@@ -257,26 +259,26 @@ class DefaultTbContext implements TbContext {
}
public TbMsg customerCreatedMsg(Customer customer, RuleNodeId ruleNodeId) {
- return entityCreatedMsg(customer, customer.getId(), ruleNodeId);
+ return entityActionMsg(customer, customer.getId(), ruleNodeId, DataConstants.ENTITY_CREATED);
}
public TbMsg deviceCreatedMsg(Device device, RuleNodeId ruleNodeId) {
- return entityCreatedMsg(device, device.getId(), ruleNodeId);
+ return entityActionMsg(device, device.getId(), ruleNodeId, DataConstants.ENTITY_CREATED);
}
public TbMsg assetCreatedMsg(Asset asset, RuleNodeId ruleNodeId) {
- return entityCreatedMsg(asset, asset.getId(), ruleNodeId);
+ return entityActionMsg(asset, asset.getId(), ruleNodeId, DataConstants.ENTITY_CREATED);
}
- public TbMsg alarmCreatedMsg(Alarm alarm, RuleNodeId ruleNodeId) {
- return entityCreatedMsg(alarm, alarm.getId(), ruleNodeId);
+ public TbMsg alarmActionMsg(Alarm alarm, RuleNodeId ruleNodeId, String action) {
+ return entityActionMsg(alarm, alarm.getId(), ruleNodeId, action);
}
- public TbMsg entityCreatedMsg(E entity, I id, RuleNodeId ruleNodeId) {
+ public TbMsg entityActionMsg(E entity, I id, RuleNodeId ruleNodeId, String action) {
try {
- return TbMsg.newMsg(DataConstants.ENTITY_CREATED, id, getActionMetaData(ruleNodeId), mapper.writeValueAsString(mapper.valueToTree(entity)));
+ return TbMsg.newMsg(action, id, getActionMetaData(ruleNodeId), mapper.writeValueAsString(mapper.valueToTree(entity)));
} catch (JsonProcessingException | IllegalArgumentException e) {
- throw new RuntimeException("Failed to process " + id.getEntityType().name().toLowerCase() + " created msg: " + e);
+ throw new RuntimeException("Failed to process " + id.getEntityType().name().toLowerCase() + " " + action + " msg: " + e);
}
}
@@ -312,7 +314,7 @@ class DefaultTbContext implements TbContext {
@Override
public ScriptEngine createJsScriptEngine(String script, String... argNames) {
- return new RuleNodeJsScriptEngine(mainCtx.getJsSandbox(), nodeCtx.getSelf().getId(), script, argNames);
+ return new RuleNodeJsScriptEngine(getTenantId(), mainCtx.getJsSandbox(), nodeCtx.getSelf().getId(), script, argNames);
}
@Override
@@ -479,8 +481,16 @@ class DefaultTbContext implements TbContext {
}
@Override
- public void addProfileListener(Consumer listener) {
- mainCtx.getDeviceProfileCache().addListener(getTenantId(), getSelfId(), listener);
+ public void removeRuleNodeStateForEntity(EntityId entityId) {
+ if (log.isDebugEnabled()) {
+ log.debug("[{}][{}][{}] Remove Rule Node State for entity.", getTenantId(), getSelfId(), entityId);
+ }
+ mainCtx.getRuleNodeStateService().removeByRuleNodeIdAndEntityId(getTenantId(), getSelfId(), entityId);
+ }
+
+ @Override
+ public void addDeviceProfileListeners(Consumer profileListener, BiConsumer deviceListener) {
+ mainCtx.getDeviceProfileCache().addListener(getTenantId(), getSelfId(), profileListener, deviceListener);
}
@Override
diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java
index 5ca394be35..ad5f97ab48 100644
--- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java
+++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java
@@ -23,6 +23,7 @@ import org.thingsboard.server.actors.TbActorRef;
import org.thingsboard.server.actors.TbEntityActorId;
import org.thingsboard.server.actors.service.DefaultActorService;
import org.thingsboard.server.actors.shared.ComponentMsgProcessor;
+import org.thingsboard.server.common.data.ApiUsageRecordKey;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.RuleChainId;
@@ -46,6 +47,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.queue.TbQueueCallback;
import org.thingsboard.server.queue.common.MultipleTbQueueTbMsgCallbackWrapper;
import org.thingsboard.server.queue.common.TbQueueTbMsgCallbackWrapper;
+import org.thingsboard.server.queue.usagestats.TbApiUsageClient;
import org.thingsboard.server.service.queue.TbClusterService;
import java.util.ArrayList;
@@ -68,15 +70,16 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor> nodeRoutes;
private final RuleChainService service;
private final TbClusterService clusterService;
+ private final TbApiUsageClient apiUsageClient;
private String ruleChainName;
private RuleNodeId firstId;
private RuleNodeCtx firstNode;
private boolean started;
- RuleChainActorMessageProcessor(TenantId tenantId, RuleChain ruleChain, ActorSystemContext systemContext
- , TbActorRef parent, TbActorRef self) {
+ RuleChainActorMessageProcessor(TenantId tenantId, RuleChain ruleChain, ActorSystemContext systemContext, TbActorRef parent, TbActorRef self) {
super(systemContext, tenantId, ruleChain.getId());
+ this.apiUsageClient = systemContext.getApiUsageClient();
this.ruleChainName = ruleChain.getName();
this.parent = parent;
this.self = self;
diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainManagerActor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainManagerActor.java
index eb46d7c098..0f828dfeb4 100644
--- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainManagerActor.java
+++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainManagerActor.java
@@ -64,6 +64,12 @@ public abstract class RuleChainManagerActor extends ContextAwareActor {
}
}
+ protected void destroyRuleChains() {
+ for (RuleChain ruleChain : new PageDataIterable<>(link -> ruleChainService.findTenantRuleChains(tenantId, link), ContextAwareActor.ENTITY_PACK_LIMIT)) {
+ ctx.stop(new TbEntityActorId(ruleChain.getId()));
+ }
+ }
+
protected void visit(RuleChain entity, TbActorRef actorRef) {
if (entity != null && entity.isRoot()) {
rootChain = entity;
diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java
index db55ff8edf..95d21090bd 100644
--- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java
+++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleNodeActorMessageProcessor.java
@@ -21,13 +21,18 @@ import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.TbActorCtx;
import org.thingsboard.server.actors.TbActorRef;
import org.thingsboard.server.actors.shared.ComponentMsgProcessor;
+import org.thingsboard.server.common.data.ApiUsageRecordKey;
+import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.RuleNodeId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleState;
import org.thingsboard.server.common.data.rule.RuleNode;
+import org.thingsboard.server.common.data.tenant.profile.TenantProfileConfiguration;
+import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.queue.PartitionChangeMsg;
import org.thingsboard.server.common.msg.queue.RuleNodeException;
import org.thingsboard.server.common.msg.queue.RuleNodeInfo;
+import org.thingsboard.server.queue.usagestats.TbApiUsageClient;
/**
* @author Andrew Shvayka
@@ -36,6 +41,7 @@ public class RuleNodeActorMessageProcessor extends ComponentMsgProcessor extends Abstract
this.entityId = id;
}
+ protected TenantProfileConfiguration getTenantProfileConfiguration() {
+ return systemContext.getTenantProfileCache().get(tenantId).getProfileData().getConfiguration();
+ }
+
public abstract String getComponentName();
public abstract void start(TbActorCtx context) throws Exception;
diff --git a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java
index 161239be7e..72351a480a 100644
--- a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java
+++ b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java
@@ -29,6 +29,7 @@ import org.thingsboard.server.actors.device.DeviceActorCreator;
import org.thingsboard.server.actors.ruleChain.RuleChainManagerActor;
import org.thingsboard.server.actors.service.ContextBasedCreator;
import org.thingsboard.server.actors.service.DefaultActorService;
+import org.thingsboard.server.common.data.ApiUsageState;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile;
@@ -57,6 +58,7 @@ public class TenantActor extends RuleChainManagerActor {
private boolean isRuleEngineForCurrentTenant;
private boolean isCore;
+ private ApiUsageState apiUsageState;
private TenantActor(ActorSystemContext systemContext, TenantId tenantId) {
super(systemContext, tenantId);
@@ -74,19 +76,24 @@ public class TenantActor extends RuleChainManagerActor {
cantFindTenant = true;
log.info("[{}] Started tenant actor for missing tenant.", tenantId);
} else {
+ apiUsageState = new ApiUsageState(systemContext.getApiUsageStateService().getApiUsageState(tenant.getId()));
+
// This Service may be started for specific tenant only.
Optional isolatedTenantId = systemContext.getServiceInfoProvider().getIsolatedTenant();
TenantProfile tenantProfile = systemContext.getTenantProfileCache().get(tenant.getTenantProfileId());
- isRuleEngineForCurrentTenant = systemContext.getServiceInfoProvider().isService(ServiceType.TB_RULE_ENGINE);
isCore = systemContext.getServiceInfoProvider().isService(ServiceType.TB_CORE);
-
+ isRuleEngineForCurrentTenant = systemContext.getServiceInfoProvider().isService(ServiceType.TB_RULE_ENGINE);
if (isRuleEngineForCurrentTenant) {
try {
if (isolatedTenantId.map(id -> id.equals(tenantId)).orElseGet(() -> !tenantProfile.isIsolatedTbRuleEngine())) {
- log.info("[{}] Going to init rule chains", tenantId);
- initRuleChains();
+ if (apiUsageState.isReExecEnabled()) {
+ log.info("[{}] Going to init rule chains", tenantId);
+ initRuleChains();
+ } else {
+ log.info("[{}] Skip init of the rule chains due to API limits", tenantId);
+ }
} else {
isRuleEngineForCurrentTenant = false;
}
@@ -98,8 +105,6 @@ public class TenantActor extends RuleChainManagerActor {
}
} catch (Exception e) {
log.warn("[{}] Unknown failure", tenantId, e);
-// TODO: throw this in 3.1?
-// throw new TbActorException("Failed to init actor", e);
}
}
@@ -115,7 +120,7 @@ public class TenantActor extends RuleChainManagerActor {
if (msg.getMsgType().equals(MsgType.QUEUE_TO_RULE_ENGINE_MSG)) {
QueueToRuleEngineMsg queueMsg = (QueueToRuleEngineMsg) msg;
queueMsg.getTbMsg().getCallback().onSuccess();
- } else if (msg.getMsgType().equals(MsgType.TRANSPORT_TO_DEVICE_ACTOR_MSG)){
+ } else if (msg.getMsgType().equals(MsgType.TRANSPORT_TO_DEVICE_ACTOR_MSG)) {
TransportToDeviceActorMsgWrapper transportMsg = (TransportToDeviceActorMsgWrapper) msg;
transportMsg.getCallback().onSuccess();
}
@@ -173,26 +178,33 @@ public class TenantActor extends RuleChainManagerActor {
return;
}
TbMsg tbMsg = msg.getTbMsg();
- if (tbMsg.getRuleChainId() == null) {
- if (getRootChainActor() != null) {
- getRootChainActor().tell(msg);
+ if (apiUsageState.isReExecEnabled()) {
+ if (tbMsg.getRuleChainId() == null) {
+ if (getRootChainActor() != null) {
+ getRootChainActor().tell(msg);
+ } else {
+ tbMsg.getCallback().onFailure(new RuleEngineException("No Root Rule Chain available!"));
+ log.info("[{}] No Root Chain: {}", tenantId, msg);
+ }
} else {
- tbMsg.getCallback().onFailure(new RuleEngineException("No Root Rule Chain available!"));
- log.info("[{}] No Root Chain: {}", tenantId, msg);
+ try {
+ ctx.tell(new TbEntityActorId(tbMsg.getRuleChainId()), msg);
+ } catch (TbActorNotRegisteredException ex) {
+ log.trace("Received message for non-existing rule chain: [{}]", tbMsg.getRuleChainId());
+ //TODO: 3.1 Log it to dead letters queue;
+ tbMsg.getCallback().onSuccess();
+ }
}
} else {
- try {
- ctx.tell(new TbEntityActorId(tbMsg.getRuleChainId()), msg);
- } catch (TbActorNotRegisteredException ex) {
- log.trace("Received message for non-existing rule chain: [{}]", tbMsg.getRuleChainId());
- //TODO: 3.1 Log it to dead letters queue;
- tbMsg.getCallback().onSuccess();
- }
+ log.trace("[{}] Ack message because Rule Engine is disabled", tenantId);
+ tbMsg.getCallback().onSuccess();
}
}
private void onRuleChainMsg(RuleChainAwareMsg msg) {
- getOrCreateActor(msg.getRuleChainId()).tell(msg);
+ if (apiUsageState.isReExecEnabled()) {
+ getOrCreateActor(msg.getRuleChainId()).tell(msg);
+ }
}
private void onToDeviceActorMsg(DeviceAwareMsg msg, boolean priority) {
@@ -208,6 +220,17 @@ public class TenantActor extends RuleChainManagerActor {
}
private void onComponentLifecycleMsg(ComponentLifecycleMsg msg) {
+ if (msg.getEntityId().getEntityType().equals(EntityType.API_USAGE_STATE)) {
+ ApiUsageState old = apiUsageState;
+ apiUsageState = new ApiUsageState(systemContext.getApiUsageStateService().getApiUsageState(tenantId));
+ if (old.isReExecEnabled() && !apiUsageState.isReExecEnabled()) {
+ log.info("[{}] Received API state update. Going to DISABLE Rule Engine execution.", tenantId);
+ destroyRuleChains();
+ } else if (!old.isReExecEnabled() && apiUsageState.isReExecEnabled()) {
+ log.info("[{}] Received API state update. Going to ENABLE Rule Engine execution.", tenantId);
+ initRuleChains();
+ }
+ }
if (isRuleEngineForCurrentTenant) {
TbActorRef target = getEntityActorRef(msg.getEntityId());
if (target != null) {
diff --git a/application/src/main/java/org/thingsboard/server/controller/AlarmController.java b/application/src/main/java/org/thingsboard/server/controller/AlarmController.java
index 4d063d07f1..47aa4aa550 100644
--- a/application/src/main/java/org/thingsboard/server/controller/AlarmController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/AlarmController.java
@@ -126,6 +126,7 @@ public class AlarmController extends BaseController {
long ackTs = System.currentTimeMillis();
alarmService.ackAlarm(getCurrentUser().getTenantId(), alarmId, ackTs).get();
alarm.setAckTs(ackTs);
+ alarm.setStatus(alarm.getStatus().isCleared() ? AlarmStatus.CLEARED_ACK : AlarmStatus.ACTIVE_ACK);
logEntityAction(alarm.getOriginator(), alarm, getCurrentUser().getCustomerId(), ActionType.ALARM_ACK, null);
} catch (Exception e) {
throw handleException(e);
@@ -143,6 +144,7 @@ public class AlarmController extends BaseController {
long clearTs = System.currentTimeMillis();
alarmService.clearAlarm(getCurrentUser().getTenantId(), alarmId, null, clearTs).get();
alarm.setClearTs(clearTs);
+ alarm.setStatus(alarm.getStatus().isAck() ? AlarmStatus.CLEARED_ACK : AlarmStatus.CLEARED_UNACK);
logEntityAction(alarm.getOriginator(), alarm, getCurrentUser().getCustomerId(), ActionType.ALARM_CLEAR, null);
} catch (Exception e) {
throw handleException(e);
diff --git a/application/src/main/java/org/thingsboard/server/controller/AuthController.java b/application/src/main/java/org/thingsboard/server/controller/AuthController.java
index 798b71e114..47119c4bb9 100644
--- a/application/src/main/java/org/thingsboard/server/controller/AuthController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/AuthController.java
@@ -166,7 +166,8 @@ public class AuthController extends BaseController {
try {
String email = resetPasswordByEmailRequest.get("email").asText();
UserCredentials userCredentials = userService.requestPasswordReset(TenantId.SYS_TENANT_ID, email);
- String baseUrl = MiscUtils.constructBaseUrl(request);
+ User user = userService.findUserById(TenantId.SYS_TENANT_ID, userCredentials.getUserId());
+ String baseUrl = systemSecurityService.getBaseUrl(user.getTenantId(), user.getCustomerId(), request);
String resetUrl = String.format("%s/api/noauth/resetPassword?resetToken=%s", baseUrl,
userCredentials.getResetToken());
@@ -214,7 +215,7 @@ public class AuthController extends BaseController {
User user = userService.findUserById(TenantId.SYS_TENANT_ID, credentials.getUserId());
UserPrincipal principal = new UserPrincipal(UserPrincipal.Type.USER_NAME, user.getEmail());
SecurityUser securityUser = new SecurityUser(user, credentials.isEnabled(), principal);
- String baseUrl = MiscUtils.constructBaseUrl(request);
+ String baseUrl = systemSecurityService.getBaseUrl(user.getTenantId(), user.getCustomerId(), request);
String loginUrl = String.format("%s/login", baseUrl);
String email = user.getEmail();
@@ -261,7 +262,7 @@ public class AuthController extends BaseController {
User user = userService.findUserById(TenantId.SYS_TENANT_ID, userCredentials.getUserId());
UserPrincipal principal = new UserPrincipal(UserPrincipal.Type.USER_NAME, user.getEmail());
SecurityUser securityUser = new SecurityUser(user, userCredentials.isEnabled(), principal);
- String baseUrl = MiscUtils.constructBaseUrl(request);
+ String baseUrl = systemSecurityService.getBaseUrl(user.getTenantId(), user.getCustomerId(), request);
String loginUrl = String.format("%s/login", baseUrl);
String email = user.getEmail();
mailService.sendPasswordWasResetEmail(loginUrl, email);
diff --git a/application/src/main/java/org/thingsboard/server/controller/BaseController.java b/application/src/main/java/org/thingsboard/server/controller/BaseController.java
index b8fe956217..869cbea127 100644
--- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java
@@ -35,7 +35,6 @@ import org.thingsboard.server.common.data.asset.AssetInfo;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode;
import org.thingsboard.server.common.data.exception.ThingsboardException;
-import org.thingsboard.server.common.data.id.*;
import org.thingsboard.server.common.data.id.AlarmId;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
@@ -94,7 +93,7 @@ import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.component.ComponentDiscoveryService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
-import org.thingsboard.server.service.profile.TbTenantProfileCache;
+import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.queue.TbClusterService;
import org.thingsboard.server.service.security.model.SecurityUser;
import org.thingsboard.server.service.security.permission.AccessControlService;
diff --git a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java
index c8fde63c4d..03ccc9fb93 100644
--- a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java
@@ -51,6 +51,7 @@ import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
+import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType;
@@ -119,6 +120,8 @@ public class DeviceController extends BaseController {
tbClusterService.pushMsgToCore(new DeviceNameOrTypeUpdateMsg(savedDevice.getTenantId(),
savedDevice.getId(), savedDevice.getName(), savedDevice.getType()), null);
+ tbClusterService.onEntityStateChange(savedDevice.getTenantId(), savedDevice.getId(),
+ device.getId() == null ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED);
logEntityAction(savedDevice.getId(), savedDevice,
savedDevice.getCustomerId(),
diff --git a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java
index 6d80c05c05..7392f74231 100644
--- a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java
@@ -160,7 +160,7 @@ public class RuleChainController extends BaseController {
public RuleChain saveRuleChain(@RequestBody DefaultRuleChainCreateRequest request) throws ThingsboardException {
try {
checkNotNull(request);
- checkNotNull(request.getName());
+ checkParameter(request.getName(), "name");
RuleChain savedRuleChain = installScripts.createDefaultRuleChain(getCurrentUser().getTenantId(), request.getName());
@@ -347,7 +347,7 @@ public class RuleChainController extends BaseController {
String errorText = "";
ScriptEngine engine = null;
try {
- engine = new RuleNodeJsScriptEngine(jsInvokeService, getCurrentUser().getId(), script, argNames);
+ engine = new RuleNodeJsScriptEngine(getTenantId(), jsInvokeService, getCurrentUser().getId(), script, argNames);
TbMsg inMsg = TbMsg.newMsg(msgType, null, new TbMsgMetaData(metadata), TbMsgDataType.JSON, data);
switch (scriptType) {
case "update":
diff --git a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java
index c8c41ea479..fb5fade521 100644
--- a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java
@@ -392,6 +392,11 @@ public class TelemetryController extends BaseController {
if (attributes.isEmpty()) {
return getImmediateDeferredResult("No attributes data found in request body!", HttpStatus.BAD_REQUEST);
}
+ for (AttributeKvEntry attributeKvEntry: attributes) {
+ if (attributeKvEntry.getKey().isEmpty() || attributeKvEntry.getKey().trim().length() == 0) {
+ return getImmediateDeferredResult("Key cannot be empty or contains only spaces", HttpStatus.BAD_REQUEST);
+ }
+ }
SecurityUser user = getCurrentUser();
return accessValidator.validateEntityAndCallback(getCurrentUser(), Operation.WRITE_ATTRIBUTES, entityIdSrc, (result, tenantId, entityId) -> {
tsSubService.saveAndNotify(tenantId, entityId, scope, attributes, new FutureCallback() {
diff --git a/application/src/main/java/org/thingsboard/server/controller/TenantController.java b/application/src/main/java/org/thingsboard/server/controller/TenantController.java
index 83e3d98b1b..a42cbb0f91 100644
--- a/application/src/main/java/org/thingsboard/server/controller/TenantController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/TenantController.java
@@ -93,6 +93,8 @@ public class TenantController extends BaseController {
}
tenantProfileCache.evict(tenant.getId());
tbClusterService.onTenantChange(tenant, null);
+ tbClusterService.onEntityStateChange(tenant.getId(), tenant.getId(),
+ newTenant ? ComponentLifecycleEvent.CREATED : ComponentLifecycleEvent.UPDATED);
return tenant;
} catch (Exception e) {
throw handleException(e);
diff --git a/application/src/main/java/org/thingsboard/server/controller/UserController.java b/application/src/main/java/org/thingsboard/server/controller/UserController.java
index fce56b509e..cc00850b4a 100644
--- a/application/src/main/java/org/thingsboard/server/controller/UserController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/UserController.java
@@ -52,6 +52,7 @@ import org.thingsboard.server.service.security.model.token.JwtToken;
import org.thingsboard.server.service.security.model.token.JwtTokenFactory;
import org.thingsboard.server.service.security.permission.Operation;
import org.thingsboard.server.service.security.permission.Resource;
+import org.thingsboard.server.service.security.system.SystemSecurityService;
import org.thingsboard.server.utils.MiscUtils;
import javax.servlet.http.HttpServletRequest;
@@ -78,6 +79,9 @@ public class UserController extends BaseController {
@Autowired
private RefreshTokenRepository refreshTokenRepository;
+ @Autowired
+ private SystemSecurityService systemSecurityService;
+
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/user/{userId}", method = RequestMethod.GET)
@@ -146,7 +150,7 @@ public class UserController extends BaseController {
if (sendEmail) {
SecurityUser authUser = getCurrentUser();
UserCredentials userCredentials = userService.findUserCredentialsByUserId(authUser.getTenantId(), savedUser.getId());
- String baseUrl = MiscUtils.constructBaseUrl(request);
+ String baseUrl = systemSecurityService.getBaseUrl(getTenantId(), getCurrentUser().getCustomerId(), request);
String activateUrl = String.format(ACTIVATE_URL_PATTERN, baseUrl,
userCredentials.getActivateToken());
String email = savedUser.getEmail();
@@ -186,7 +190,7 @@ public class UserController extends BaseController {
UserCredentials userCredentials = userService.findUserCredentialsByUserId(getCurrentUser().getTenantId(), user.getId());
if (!userCredentials.isEnabled()) {
- String baseUrl = MiscUtils.constructBaseUrl(request);
+ String baseUrl = systemSecurityService.getBaseUrl(getTenantId(), getCurrentUser().getCustomerId(), request);
String activateUrl = String.format(ACTIVATE_URL_PATTERN, baseUrl,
userCredentials.getActivateToken());
mailService.sendActivationEmail(activateUrl, email);
@@ -211,7 +215,7 @@ public class UserController extends BaseController {
SecurityUser authUser = getCurrentUser();
UserCredentials userCredentials = userService.findUserCredentialsByUserId(authUser.getTenantId(), user.getId());
if (!userCredentials.isEnabled()) {
- String baseUrl = MiscUtils.constructBaseUrl(request);
+ String baseUrl = systemSecurityService.getBaseUrl(getTenantId(), getCurrentUser().getCustomerId(), request);
String activateUrl = String.format(ACTIVATE_URL_PATTERN, baseUrl,
userCredentials.getActivateToken());
return activateUrl;
diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
new file mode 100644
index 0000000000..c7206712c3
--- /dev/null
+++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
@@ -0,0 +1,370 @@
+/**
+ * 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.boot.context.event.ApplicationReadyEvent;
+import org.springframework.context.annotation.Lazy;
+import org.springframework.context.event.EventListener;
+import org.springframework.core.annotation.Order;
+import org.springframework.data.util.Pair;
+import org.springframework.stereotype.Service;
+import org.thingsboard.server.common.data.ApiFeature;
+import org.thingsboard.server.common.data.ApiUsageRecordKey;
+import org.thingsboard.server.common.data.ApiUsageState;
+import org.thingsboard.server.common.data.ApiUsageStateValue;
+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.LongDataEntry;
+import org.thingsboard.server.common.data.kv.StringDataEntry;
+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.dao.tenant.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.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.ExecutionException;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
+
+@Slf4j
+@Service
+public class DefaultTbApiUsageStateService implements TbApiUsageStateService {
+
+ public static final String HOURLY = "Hourly";
+ public static final FutureCallback VOID_CALLBACK = new FutureCallback() {
+ @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 myTenantStates = new ConcurrentHashMap<>();
+ // Tenants that should be processed on other servers
+ private final Map 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.");
+ scheduler.scheduleAtFixedRate(this::checkStartOfNextCycle, nextCycleCheckInterval, nextCycleCheckInterval, TimeUnit.MILLISECONDS);
+ log.info("Started api usage service.");
+ }
+ }
+
+ @Override
+ public void process(TbProtoQueueMsg msg, TbCallback callback) {
+ ToUsageStatsServiceMsg statsMsg = msg.getValue();
+ TenantId tenantId = new TenantId(new UUID(statsMsg.getTenantIdMSB(), statsMsg.getTenantIdLSB()));
+ TenantApiUsageState tenantState;
+ List updatedEntries;
+ Map result;
+ 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);
+ Set apiFeatures = new HashSet<>();
+ 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)));
+ apiFeatures.add(recordKey.getApiFeature());
+ }
+ result = tenantState.checkStateUpdatedDueToThreshold(apiFeatures);
+ } 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);
+ if (state != null) {
+ otherTenantStates.put(tenantId, state);
+ }
+ }
+ } finally {
+ updateLock.unlock();
+ }
+ return state;
+ }
+ }
+ }
+ }
+
+ @Override
+ public void onApiUsageStateUpdate(TenantId tenantId) {
+ otherTenantStates.remove(tenantId);
+ }
+
+ @Override
+ public void onTenantProfileUpdate(TenantProfileId tenantProfileId) {
+ log.info("[{}] On Tenant Profile Update", 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) {
+ log.info("[{}] On Tenant Update.", 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 profile) {
+ TenantProfileData oldProfileData = state.getTenantProfileData();
+ state.setTenantProfileId(profile.getId());
+ state.setTenantProfileData(profile.getProfileData());
+ Map result = state.checkStateUpdatedDueToThresholds();
+ if (!result.isEmpty()) {
+ persistAndNotify(state, result);
+ }
+ updateProfileThresholds(state.getTenantId(), state.getApiUsageState().getId(),
+ oldProfileData.getConfiguration(), profile.getProfileData().getConfiguration());
+ }
+
+ private void updateProfileThresholds(TenantId tenantId, ApiUsageStateId id,
+ TenantProfileConfiguration oldData, TenantProfileConfiguration newData) {
+ long ts = System.currentTimeMillis();
+ List 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 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 stateTelemetry = new ArrayList<>();
+ result.forEach(((apiFeature, aState) -> stateTelemetry.add(new BasicTsKvEntry(ts, new StringDataEntry(apiFeature.getApiStateKey(), aState.name())))));
+ tsWsService.saveAndNotifyInternal(state.getTenantId(), state.getApiUsageState().getId(), stateTelemetry, VOID_CALLBACK);
+ //TODO: notify tenant admin via email!
+ }
+
+ 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 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 {
+ log.info("Initializing tenant states.");
+ PageDataIterable 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()) {
+ log.debug("[{}] Initializing tenant state.", tenant.getId());
+ updateLock.lock();
+ try {
+ updateTenantState(getOrFetchState(tenant.getId()), tenantProfileCache.get(tenant.getTenantProfileId()));
+ log.debug("[{}] Initialized tenant state.", tenant.getId());
+ } catch (Exception e) {
+ log.warn("[{}] Failed to initialize tenant API state", tenant.getId(), e);
+ } finally {
+ updateLock.unlock();
+ }
+ }
+ }
+ log.info("Initialized tenant states.");
+ } catch (Exception e) {
+ log.warn("Unknown failure", e);
+ }
+ }
+
+}
diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java
new file mode 100644
index 0000000000..e19c0a4154
--- /dev/null
+++ b/application/src/main/java/org/thingsboard/server/service/apiusage/TbApiUsageStateService.java
@@ -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 {
+
+ void process(TbProtoQueueMsg msg, TbCallback callback);
+
+ ApiUsageState getApiUsageState(TenantId tenantId);
+
+ void onTenantProfileUpdate(TenantProfileId tenantProfileId);
+
+ void onTenantUpdate(TenantId tenantId);
+
+ void onApiUsageStateUpdate(TenantId tenantId);
+}
diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java b/application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java
new file mode 100644
index 0000000000..a43afcd2c5
--- /dev/null
+++ b/application/src/main/java/org/thingsboard/server/service/apiusage/TenantApiUsageState.java
@@ -0,0 +1,186 @@
+/**
+ * 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.ApiFeature;
+import org.thingsboard.server.common.data.ApiUsageRecordKey;
+import org.thingsboard.server.common.data.ApiUsageState;
+import org.thingsboard.server.common.data.ApiUsageStateValue;
+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.TenantProfileData;
+import org.thingsboard.server.common.msg.tools.SchedulerUtils;
+
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+
+public class TenantApiUsageState {
+
+ private final Map currentCycleValues = new ConcurrentHashMap<>();
+ private final Map 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 long getProfileWarnThreshold(ApiUsageRecordKey key) {
+ return tenantProfileData.getConfiguration().getWarnThreshold(key);
+ }
+
+ public TenantId getTenantId() {
+ return apiUsageState.getTenantId();
+ }
+
+ public ApiUsageStateValue getFeatureValue(ApiFeature feature) {
+ switch (feature) {
+ case TRANSPORT:
+ return apiUsageState.getTransportState();
+ case RE:
+ return apiUsageState.getReExecState();
+ case DB:
+ return apiUsageState.getDbStorageState();
+ case JS:
+ return apiUsageState.getJsExecState();
+ default:
+ return ApiUsageStateValue.ENABLED;
+ }
+ }
+
+ public boolean setFeatureValue(ApiFeature feature, ApiUsageStateValue value) {
+ ApiUsageStateValue currentValue = getFeatureValue(feature);
+ switch (feature) {
+ case TRANSPORT:
+ apiUsageState.setTransportState(value);
+ break;
+ case RE:
+ apiUsageState.setReExecState(value);
+ break;
+ case DB:
+ apiUsageState.setDbStorageState(value);
+ break;
+ case JS:
+ apiUsageState.setJsExecState(value);
+ break;
+ }
+ return !currentValue.equals(value);
+ }
+
+ public Map checkStateUpdatedDueToThresholds() {
+ return checkStateUpdatedDueToThreshold(new HashSet<>(Arrays.asList(ApiFeature.values())));
+ }
+
+ public Map checkStateUpdatedDueToThreshold(Set features) {
+ Map result = new HashMap<>();
+ for (ApiFeature feature : features) {
+ Pair tmp = checkStateUpdatedDueToThreshold(feature);
+ if (tmp != null) {
+ result.put(tmp.getFirst(), tmp.getSecond());
+ }
+ }
+ return result;
+ }
+
+ public Pair checkStateUpdatedDueToThreshold(ApiFeature feature) {
+ ApiUsageStateValue featureValue = ApiUsageStateValue.ENABLED;
+ for (ApiUsageRecordKey recordKey : ApiUsageRecordKey.getKeys(feature)) {
+ long value = get(recordKey);
+ long threshold = getProfileThreshold(recordKey);
+ long warnThreshold = getProfileWarnThreshold(recordKey);
+ ApiUsageStateValue tmpValue;
+ if (threshold == 0 || value < warnThreshold) {
+ tmpValue = ApiUsageStateValue.ENABLED;
+ } else if (value < threshold) {
+ tmpValue = ApiUsageStateValue.WARNING;
+ } else {
+ tmpValue = ApiUsageStateValue.DISABLED;
+ }
+ featureValue = ApiUsageStateValue.toMoreRestricted(featureValue, tmpValue);
+ }
+ return setFeatureValue(feature, featureValue) ? Pair.of(feature, featureValue) : null;
+ }
+
+}
diff --git a/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java b/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java
index 1a0f9eca65..6c7678f023 100644
--- a/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java
+++ b/application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java
@@ -30,7 +30,8 @@ import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile;
-import org.thingsboard.server.common.data.TenantProfileData;
+import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
+import org.thingsboard.server.common.data.tenant.profile.TenantProfileData;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.CustomerId;
@@ -127,13 +128,16 @@ public class DefaultSystemDataLoaderService implements SystemDataLoaderService {
public void createDefaultTenantProfiles() throws Exception {
tenantProfileService.findOrCreateDefaultTenantProfile(TenantId.SYS_TENANT_ID);
+ TenantProfileData tenantProfileData = new TenantProfileData();
+ tenantProfileData.setConfiguration(new DefaultTenantProfileConfiguration());
+
TenantProfile isolatedTbCoreProfile = new TenantProfile();
isolatedTbCoreProfile.setDefault(false);
isolatedTbCoreProfile.setName("Isolated TB Core");
- isolatedTbCoreProfile.setProfileData(new TenantProfileData());
isolatedTbCoreProfile.setDescription("Isolated TB Core tenant profile");
isolatedTbCoreProfile.setIsolatedTbCore(true);
isolatedTbCoreProfile.setIsolatedTbRuleEngine(false);
+ isolatedTbCoreProfile.setProfileData(tenantProfileData);
try {
tenantProfileService.saveTenantProfile(TenantId.SYS_TENANT_ID, isolatedTbCoreProfile);
} catch (DataValidationException e) {
@@ -143,10 +147,11 @@ public class DefaultSystemDataLoaderService implements SystemDataLoaderService {
TenantProfile isolatedTbRuleEngineProfile = new TenantProfile();
isolatedTbRuleEngineProfile.setDefault(false);
isolatedTbRuleEngineProfile.setName("Isolated TB Rule Engine");
- isolatedTbRuleEngineProfile.setProfileData(new TenantProfileData());
isolatedTbRuleEngineProfile.setDescription("Isolated TB Rule Engine tenant profile");
isolatedTbRuleEngineProfile.setIsolatedTbCore(false);
isolatedTbRuleEngineProfile.setIsolatedTbRuleEngine(true);
+ isolatedTbRuleEngineProfile.setProfileData(tenantProfileData);
+
try {
tenantProfileService.saveTenantProfile(TenantId.SYS_TENANT_ID, isolatedTbRuleEngineProfile);
} catch (DataValidationException e) {
@@ -156,10 +161,11 @@ public class DefaultSystemDataLoaderService implements SystemDataLoaderService {
TenantProfile isolatedTbCoreAndTbRuleEngineProfile = new TenantProfile();
isolatedTbCoreAndTbRuleEngineProfile.setDefault(false);
isolatedTbCoreAndTbRuleEngineProfile.setName("Isolated TB Core and TB Rule Engine");
- isolatedTbCoreAndTbRuleEngineProfile.setProfileData(new TenantProfileData());
isolatedTbCoreAndTbRuleEngineProfile.setDescription("Isolated TB Core and TB Rule Engine tenant profile");
isolatedTbCoreAndTbRuleEngineProfile.setIsolatedTbCore(true);
isolatedTbCoreAndTbRuleEngineProfile.setIsolatedTbRuleEngine(true);
+ isolatedTbCoreAndTbRuleEngineProfile.setProfileData(tenantProfileData);
+
try {
tenantProfileService.saveTenantProfile(TenantId.SYS_TENANT_ID, isolatedTbCoreAndTbRuleEngineProfile);
} catch (DataValidationException e) {
@@ -173,6 +179,7 @@ public class DefaultSystemDataLoaderService implements SystemDataLoaderService {
generalSettings.setKey("general");
ObjectNode node = objectMapper.createObjectNode();
node.put("baseUrl", "http://localhost:8080");
+ node.put("prohibitDifferentUrl", true);
generalSettings.setJsonValue(node);
adminSettingsService.saveAdminSettings(TenantId.SYS_TENANT_ID, generalSettings);
diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java
index d6068b9616..19727684a3 100644
--- a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java
+++ b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java
@@ -28,6 +28,7 @@ import org.thingsboard.server.dao.dashboard.DashboardService;
import org.thingsboard.server.dao.device.DeviceProfileService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.tenant.TenantService;
+import org.thingsboard.server.dao.usagerecord.ApiUsageStateService;
import org.thingsboard.server.service.install.sql.SqlDbHelper;
import java.nio.charset.Charset;
@@ -96,6 +97,9 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
@Autowired
private DeviceProfileService deviceProfileService;
+ @Autowired
+ private ApiUsageStateService apiUsageStateService;
+
@Override
public void upgradeDatabase(String fromVersion) throws Exception {
@@ -352,6 +356,22 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
} catch (Exception e) {
}
+ try {
+ conn.createStatement().execute("CREATE TABLE IF NOT EXISTS api_usage_state (" +
+ " id uuid NOT NULL CONSTRAINT usage_record_pkey PRIMARY KEY," +
+ " created_time bigint NOT NULL," +
+ " tenant_id uuid," +
+ " entity_type varchar(32)," +
+ " entity_id uuid," +
+ " transport varchar(32)," +
+ " db_storage varchar(32)," +
+ " re_exec varchar(32)," +
+ " js_exec varchar(32)," +
+ " CONSTRAINT api_usage_state_unq_key UNIQUE (tenant_id, entity_id)\n" +
+ ");");
+ } catch (Exception e) {
+ }
+
schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "3.1.1", "schema_update_before.sql");
loadSql(schemaUpdateFile, conn);
@@ -367,6 +387,10 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
do {
pageData = tenantService.findTenants(pageLink);
for (Tenant tenant : pageData.getData()) {
+ try {
+ apiUsageStateService.createDefaultApiUsageState(tenant.getId());
+ } catch (Exception e) {
+ }
List deviceTypes = deviceService.findDeviceTypesByTenantId(tenant.getId()).get();
try {
deviceProfileService.createDefaultDeviceProfile(tenant.getId());
diff --git a/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java b/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java
index 6784405adb..3db23c65e2 100644
--- a/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java
+++ b/application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java
@@ -26,11 +26,11 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.dao.device.DeviceProfileService;
import org.thingsboard.server.dao.device.DeviceService;
-import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
+import java.util.function.BiConsumer;
import java.util.function.Consumer;
@Service
@@ -43,7 +43,8 @@ public class DefaultTbDeviceProfileCache implements TbDeviceProfileCache {
private final ConcurrentMap deviceProfilesMap = new ConcurrentHashMap<>();
private final ConcurrentMap devicesMap = new ConcurrentHashMap<>();
- private final ConcurrentMap>> listeners = new ConcurrentHashMap<>();
+ private final ConcurrentMap>> profileListeners = new ConcurrentHashMap<>();
+ private final ConcurrentMap>> deviceProfileListeners = new ConcurrentHashMap<>();
public DefaultTbDeviceProfileCache(DeviceProfileService deviceProfileService, DeviceService deviceService) {
this.deviceProfileService = deviceProfileService;
@@ -87,33 +88,37 @@ public class DefaultTbDeviceProfileCache implements TbDeviceProfileCache {
return get(tenantId, profileId);
}
- @Override
- public void put(DeviceProfile profile) {
- if (profile.getId() != null) {
- deviceProfilesMap.put(profile.getId(), profile);
- log.debug("[{}] pushed device profile to cache: {}", profile.getId(), profile);
- notifyListeners(profile);
- }
- }
-
@Override
public void evict(TenantId tenantId, DeviceProfileId profileId) {
DeviceProfile oldProfile = deviceProfilesMap.remove(profileId);
log.debug("[{}] evict device profile from cache: {}", profileId, oldProfile);
DeviceProfile newProfile = get(tenantId, profileId);
if (newProfile != null) {
- notifyListeners(newProfile);
+ notifyProfileListeners(newProfile);
}
}
@Override
- public void evict(DeviceId deviceId) {
- devicesMap.remove(deviceId);
+ public void evict(TenantId tenantId, DeviceId deviceId) {
+ DeviceProfileId old = devicesMap.remove(deviceId);
+ if (old != null) {
+ DeviceProfile newProfile = get(tenantId, deviceId);
+ if (newProfile == null || !old.equals(newProfile.getId())) {
+ notifyDeviceListeners(tenantId, deviceId, newProfile);
+ }
+ }
}
@Override
- public void addListener(TenantId tenantId, EntityId listenerId, Consumer listener) {
- listeners.computeIfAbsent(tenantId, id -> new ConcurrentHashMap<>()).put(listenerId, listener);
+ public void addListener(TenantId tenantId, EntityId listenerId,
+ Consumer profileListener,
+ BiConsumer deviceListener) {
+ if (profileListener != null) {
+ profileListeners.computeIfAbsent(tenantId, id -> new ConcurrentHashMap<>()).put(listenerId, profileListener);
+ }
+ if (deviceListener != null) {
+ deviceProfileListeners.computeIfAbsent(tenantId, id -> new ConcurrentHashMap<>()).put(listenerId, deviceListener);
+ }
}
@Override
@@ -128,17 +133,30 @@ public class DefaultTbDeviceProfileCache implements TbDeviceProfileCache {
@Override
public void removeListener(TenantId tenantId, EntityId listenerId) {
- ConcurrentMap> tenantListeners = listeners.get(tenantId);
+ ConcurrentMap> tenantListeners = profileListeners.get(tenantId);
if (tenantListeners != null) {
tenantListeners.remove(listenerId);
}
+ ConcurrentMap> deviceListeners = deviceProfileListeners.get(tenantId);
+ if (deviceListeners != null) {
+ deviceListeners.remove(listenerId);
+ }
}
- private void notifyListeners(DeviceProfile profile) {
- ConcurrentMap> tenantListeners = listeners.get(profile.getTenantId());
+ private void notifyProfileListeners(DeviceProfile profile) {
+ ConcurrentMap> tenantListeners = profileListeners.get(profile.getTenantId());
if (tenantListeners != null) {
tenantListeners.forEach((id, listener) -> listener.accept(profile));
}
}
+ private void notifyDeviceListeners(TenantId tenantId, DeviceId deviceId, DeviceProfile profile) {
+ if (profile != null) {
+ ConcurrentMap> tenantListeners = deviceProfileListeners.get(tenantId);
+ if (tenantListeners != null) {
+ tenantListeners.forEach((id, listener) -> listener.accept(deviceId, profile));
+ }
+ }
+ }
+
}
diff --git a/application/src/main/java/org/thingsboard/server/service/profile/TbDeviceProfileCache.java b/application/src/main/java/org/thingsboard/server/service/profile/TbDeviceProfileCache.java
index e65f297d53..10e28bd299 100644
--- a/application/src/main/java/org/thingsboard/server/service/profile/TbDeviceProfileCache.java
+++ b/application/src/main/java/org/thingsboard/server/service/profile/TbDeviceProfileCache.java
@@ -23,11 +23,9 @@ import org.thingsboard.server.common.data.id.TenantId;
public interface TbDeviceProfileCache extends RuleEngineDeviceProfileCache {
- void put(DeviceProfile profile);
-
void evict(TenantId tenantId, DeviceProfileId id);
- void evict(DeviceId id);
+ void evict(TenantId tenantId, DeviceId id);
DeviceProfile find(DeviceProfileId deviceProfileId);
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
index 9451b58e9f..edd53c6727 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
@@ -21,6 +21,7 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
import org.thingsboard.rule.engine.api.msg.ToDeviceActorNotificationMsg;
+import org.thingsboard.server.common.data.ApiUsageState;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.HasName;
@@ -47,6 +48,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotifica
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
import org.thingsboard.server.queue.TbQueueCallback;
import org.thingsboard.server.queue.TbQueueProducer;
+import org.thingsboard.server.queue.common.MultipleTbQueueCallbackWrapper;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
@@ -156,8 +158,16 @@ public class DefaultTbClusterService implements TbClusterService {
private TbMsg transformMsg(TbMsg tbMsg, DeviceProfile deviceProfile) {
if (deviceProfile != null) {
RuleChainId targetRuleChainId = deviceProfile.getDefaultRuleChainId();
- if (targetRuleChainId != null && !targetRuleChainId.equals(tbMsg.getRuleChainId())) {
+ String targetQueueName = deviceProfile.getDefaultQueueName();
+ boolean isRuleChainTransform = targetRuleChainId != null && !targetRuleChainId.equals(tbMsg.getRuleChainId());
+ boolean isQueueTransform = targetQueueName != null && !targetQueueName.equals(tbMsg.getQueueName());
+
+ if (isRuleChainTransform && isQueueTransform) {
+ tbMsg = TbMsg.transformMsg(tbMsg, targetRuleChainId, targetQueueName);
+ } else if (isRuleChainTransform) {
tbMsg = TbMsg.transformMsg(tbMsg, targetRuleChainId);
+ } else if (isQueueTransform) {
+ tbMsg = TbMsg.transformMsg(tbMsg, targetQueueName);
}
}
return tbMsg;
@@ -206,6 +216,12 @@ public class DefaultTbClusterService implements TbClusterService {
onEntityChange(TenantId.SYS_TENANT_ID, tenant.getId(), tenant, callback);
}
+ @Override
+ public void onApiStateChange(ApiUsageState apiUsageState, TbQueueCallback callback) {
+ onEntityChange(apiUsageState.getTenantId(), apiUsageState.getId(), apiUsageState, callback);
+ broadcast(new ComponentLifecycleMsg(apiUsageState.getTenantId(), apiUsageState.getId(), ComponentLifecycleEvent.UPDATED));
+ }
+
@Override
public void onDeviceProfileDelete(DeviceProfile entity, TbQueueCallback callback) {
onEntityDelete(entity.getTenantId(), entity.getId(), entity.getName(), callback);
@@ -221,13 +237,14 @@ public class DefaultTbClusterService implements TbClusterService {
onEntityDelete(TenantId.SYS_TENANT_ID, entity.getId(), entity.getName(), callback);
}
- public void onEntityChange(TenantId tenantId, EntityId entityid, T entity, TbQueueCallback callback) {
- log.trace("[{}][{}][{}] Processing [{}] change event", tenantId, entityid.getEntityType(), entityid.getId(), entity.getName());
+ public void onEntityChange(TenantId tenantId, EntityId entityid, T entity, TbQueueCallback callback) {
+ String entityName = (entity instanceof HasName) ? ((HasName) entity).getName() : entity.getClass().getName();
+ log.trace("[{}][{}][{}] Processing [{}] change event", tenantId, entityid.getEntityType(), entityid.getId(), entityName);
TransportProtos.EntityUpdateMsg entityUpdateMsg = TransportProtos.EntityUpdateMsg.newBuilder()
.setEntityType(entityid.getEntityType().name())
.setData(ByteString.copyFrom(encodingService.encode(entity))).build();
ToTransportMsg transportMsg = ToTransportMsg.newBuilder().setEntityUpdateMsg(entityUpdateMsg).build();
- broadcast(transportMsg);
+ broadcast(transportMsg, callback);
}
private void onEntityDelete(TenantId tenantId, EntityId entityId, String name, TbQueueCallback callback) {
@@ -238,15 +255,16 @@ public class DefaultTbClusterService implements TbClusterService {
.setEntityIdLSB(entityId.getId().getLeastSignificantBits())
.build();
ToTransportMsg transportMsg = ToTransportMsg.newBuilder().setEntityDeleteMsg(entityDeleteMsg).build();
- broadcast(transportMsg);
+ broadcast(transportMsg, callback);
}
- private void broadcast(ToTransportMsg transportMsg) {
+ private void broadcast(ToTransportMsg transportMsg, TbQueueCallback callback) {
TbQueueProducer> toTransportNfProducer = producerProvider.getTransportNotificationsMsgProducer();
Set tbTransportServices = partitionService.getAllServiceIds(ServiceType.TB_TRANSPORT);
+ TbQueueCallback proxyCallback = callback != null ? new MultipleTbQueueCallbackWrapper(tbTransportServices.size(), callback) : null;
for (String transportServiceId : tbTransportServices) {
TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_TRANSPORT, transportServiceId);
- toTransportNfProducer.send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), transportMsg), null);
+ toTransportNfProducer.send(tpi, new TbProtoQueueMsg<>(UUID.randomUUID(), transportMsg), proxyCallback);
toTransportNfs.incrementAndGet();
}
}
@@ -256,7 +274,9 @@ public class DefaultTbClusterService implements TbClusterService {
TbQueueProducer> toRuleEngineProducer = producerProvider.getRuleEngineNotificationsMsgProducer();
Set tbRuleEngineServices = new HashSet<>(partitionService.getAllServiceIds(ServiceType.TB_RULE_ENGINE));
if (msg.getEntityId().getEntityType().equals(EntityType.TENANT)
- || msg.getEntityId().getEntityType().equals(EntityType.DEVICE_PROFILE)) {
+ || msg.getEntityId().getEntityType().equals(EntityType.TENANT_PROFILE)
+ || msg.getEntityId().getEntityType().equals(EntityType.DEVICE_PROFILE)
+ || msg.getEntityId().getEntityType().equals(EntityType.API_USAGE_STATE)) {
TbQueueProducer> toCoreNfProducer = producerProvider.getTbCoreNotificationsMsgProducer();
Set tbCoreServices = partitionService.getAllServiceIds(ServiceType.TB_CORE);
for (String serviceId : tbCoreServices) {
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
index 5d240790f2..f56e08b10f 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
@@ -17,41 +17,45 @@ package org.thingsboard.server.service.queue;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
+import org.springframework.boot.context.event.ApplicationReadyEvent;
+import org.springframework.context.event.EventListener;
+import org.springframework.core.annotation.Order;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
+import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.rule.engine.api.RpcError;
import org.thingsboard.server.actors.ActorSystemContext;
-import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.alarm.Alarm;
-import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg;
-import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback;
+import org.thingsboard.server.common.stats.StatsFactory;
+import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.dao.util.mapping.JacksonUtil;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceStateServiceMsgProto;
import org.thingsboard.server.gen.transport.TransportProtos.FromDeviceRPCResponseProto;
import org.thingsboard.server.gen.transport.TransportProtos.LocalSubscriptionServiceMsgProto;
import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionMgrMsgProto;
-import org.thingsboard.server.gen.transport.TransportProtos.TbAttributeUpdateProto;
-import org.thingsboard.server.gen.transport.TransportProtos.TbAttributeDeleteProto;
-import org.thingsboard.server.gen.transport.TransportProtos.TbAlarmUpdateProto;
import org.thingsboard.server.gen.transport.TransportProtos.TbAlarmDeleteProto;
+import org.thingsboard.server.gen.transport.TransportProtos.TbAlarmUpdateProto;
+import org.thingsboard.server.gen.transport.TransportProtos.TbAttributeDeleteProto;
+import org.thingsboard.server.gen.transport.TransportProtos.TbAttributeUpdateProto;
import org.thingsboard.server.gen.transport.TransportProtos.TbSubscriptionCloseProto;
import org.thingsboard.server.gen.transport.TransportProtos.TbTimeSeriesUpdateProto;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
+import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg;
import org.thingsboard.server.gen.transport.TransportProtos.TransportToDeviceActorMsg;
import org.thingsboard.server.queue.TbQueueConsumer;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
-import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.queue.util.TbCoreComponent;
-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.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.queue.processing.AbstractConsumerService;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponse;
import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService;
@@ -70,6 +74,8 @@ import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.function.Function;
import java.util.stream.Collectors;
@@ -88,32 +94,57 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService> mainConsumer;
private final DeviceStateService stateService;
+ private final TbApiUsageStateService statsService;
private final TbLocalSubscriptionService localSubscriptionService;
private final SubscriptionManagerService subscriptionManagerService;
private final TbCoreDeviceRpcService tbCoreDeviceRpcService;
private final TbCoreConsumerStats stats;
+ protected final TbQueueConsumer> usageStatsConsumer;
+
+ protected volatile ExecutorService usageStatsExecutor;
- public DefaultTbCoreConsumerService(TbCoreQueueFactory tbCoreQueueFactory, ActorSystemContext actorContext,
- DeviceStateService stateService, TbLocalSubscriptionService localSubscriptionService,
- SubscriptionManagerService subscriptionManagerService, DataDecodingEncodingService encodingService,
- TbCoreDeviceRpcService tbCoreDeviceRpcService, StatsFactory statsFactory, TbDeviceProfileCache deviceProfileCache) {
- super(actorContext, encodingService, deviceProfileCache, tbCoreQueueFactory.createToCoreNotificationsMsgConsumer());
+ public DefaultTbCoreConsumerService(TbCoreQueueFactory tbCoreQueueFactory,
+ ActorSystemContext actorContext,
+ DeviceStateService stateService,
+ TbLocalSubscriptionService localSubscriptionService,
+ SubscriptionManagerService subscriptionManagerService,
+ 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.usageStatsConsumer = tbCoreQueueFactory.createToUsageStatsServiceMsgConsumer();
this.stateService = stateService;
this.localSubscriptionService = localSubscriptionService;
this.subscriptionManagerService = subscriptionManagerService;
this.tbCoreDeviceRpcService = tbCoreDeviceRpcService;
this.stats = new TbCoreConsumerStats(statsFactory);
+ this.statsService = statsService;
}
@PostConstruct
public void init() {
super.init("tb-core-consumer", "tb-core-notifications-consumer");
+ this.usageStatsExecutor = Executors.newCachedThreadPool(ThingsBoardThreadFactory.forName("tb-core-usage-stats-consumer"));
}
@PreDestroy
public void destroy() {
super.destroy();
+ if (usageStatsExecutor != null) {
+ usageStatsExecutor.shutdownNow();
+ }
+ }
+
+ @EventListener(ApplicationReadyEvent.class)
+ @Order(value = 2)
+ public void onApplicationEvent(ApplicationReadyEvent event) {
+ super.onApplicationEvent(event);
+ launchUsageStatsConsumer();
}
@Override
@@ -121,6 +152,12 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService tpi.newByTopic(usageStatsConsumer.getTopic()))
+ .collect(Collectors.toSet()));
}
}
@@ -223,6 +260,53 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService {
+ while (!stopped) {
+ try {
+ List> msgs = usageStatsConsumer.poll(getNotificationPollDuration());
+ if (msgs.isEmpty()) {
+ continue;
+ }
+ ConcurrentMap> pendingMap = msgs.stream().collect(
+ Collectors.toConcurrentMap(s -> UUID.randomUUID(), Function.identity()));
+ CountDownLatch processingTimeoutLatch = new CountDownLatch(1);
+ TbPackProcessingContext> ctx = new TbPackProcessingContext<>(
+ processingTimeoutLatch, pendingMap, new ConcurrentHashMap<>());
+ pendingMap.forEach((id, msg) -> {
+ log.trace("[{}] Creating usage stats callback for message: {}", id, msg.getValue());
+ TbCallback callback = new TbPackCallback<>(id, ctx);
+ try {
+ handleUsageStats(msg, callback);
+ } catch (Throwable e) {
+ log.warn("[{}] Failed to process usge stats: {}", id, msg, e);
+ callback.onFailure(e);
+ }
+ });
+ if (!processingTimeoutLatch.await(getNotificationPackProcessingTimeout(), TimeUnit.MILLISECONDS)) {
+ ctx.getAckMap().forEach((id, msg) -> log.warn("[{}] Timeout to process usage stats: {}", id, msg.getValue()));
+ ctx.getFailedMap().forEach((id, msg) -> log.warn("[{}] Failed to process usage stats: {}", id, msg.getValue()));
+ }
+ usageStatsConsumer.commit();
+ } catch (Exception e) {
+ if (!stopped) {
+ log.warn("Failed to obtain usage stats from queue.", e);
+ try {
+ Thread.sleep(getNotificationPollDuration());
+ } catch (InterruptedException e2) {
+ log.trace("Failed to wait until the server has capacity to handle new usage stats", e2);
+ }
+ }
+ }
+ }
+ log.info("TB Usage Stats Consumer stopped.");
+ });
+ }
+
+ private void handleUsageStats(TbProtoQueueMsg msg, TbCallback callback) {
+ statsService.process(msg, callback);
+ }
+
private void forwardToCoreRpcService(FromDeviceRPCResponseProto proto, TbCallback callback) {
RpcError error = proto.getError() > 0 ? RpcError.values()[proto.getError()] : null;
FromDeviceRpcResponse response = new FromDeviceRpcResponse(new UUID(proto.getRequestIdMSB(), proto.getRequestIdLSB())
@@ -321,6 +405,9 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService> nfConsumer;
public AbstractConsumerService(ActorSystemContext actorContext, DataDecodingEncodingService encodingService,
- TbDeviceProfileCache deviceProfileCache, TbQueueConsumer> nfConsumer) {
+ TbTenantProfileCache tenantProfileCache, TbDeviceProfileCache deviceProfileCache, TbApiUsageStateService apiUsageStateService, TbQueueConsumer> nfConsumer) {
this.actorContext = actorContext;
this.encodingService = encodingService;
+ this.tenantProfileCache = tenantProfileCache;
this.deviceProfileCache = deviceProfileCache;
+ this.apiUsageStateService = apiUsageStateService;
this.nfConsumer = nfConsumer;
}
@@ -77,6 +86,7 @@ public abstract class AbstractConsumerService scriptIdToNameMap = new ConcurrentHashMap<>();
- protected Map blackListedFunctions = new ConcurrentHashMap<>();
+ protected Map disabledFunctions = new ConcurrentHashMap<>();
+
+ protected AbstractJsInvokeService(TbApiUsageStateService apiUsageStateService, TbApiUsageClient apiUsageClient) {
+ this.apiUsageStateService = apiUsageStateService;
+ this.apiUsageClient = apiUsageClient;
+ }
public void init(long maxRequestsTimeout) {
if (maxRequestsTimeout > 0) {
@@ -50,24 +61,33 @@ public abstract class AbstractJsInvokeService implements JsInvokeService {
}
@Override
- public ListenableFuture eval(JsScriptType scriptType, String scriptBody, String... argNames) {
- UUID scriptId = UUID.randomUUID();
- String functionName = "invokeInternal_" + scriptId.toString().replace('-', '_');
- String jsScript = generateJsScript(scriptType, functionName, scriptBody, argNames);
- return doEval(scriptId, functionName, jsScript);
+ public ListenableFuture eval(TenantId tenantId, JsScriptType scriptType, String scriptBody, String... argNames) {
+ if (apiUsageStateService.getApiUsageState(tenantId).isJsExecEnabled()) {
+ UUID scriptId = UUID.randomUUID();
+ String functionName = "invokeInternal_" + scriptId.toString().replace('-', '_');
+ String jsScript = generateJsScript(scriptType, functionName, scriptBody, argNames);
+ return doEval(scriptId, functionName, jsScript);
+ } else {
+ return Futures.immediateFailedFuture(new RuntimeException("JS Execution is disabled due to API limits!"));
+ }
}
@Override
- public ListenableFuture