Browse Source

Merge remote-tracking branch 'ce/master' into feature/http_header_as_array_if_multivalue

pull/7587/head
Sergey Matvienko 4 years ago
parent
commit
86127d662b
  1. 62
      application/src/main/data/upgrade/3.4.1/schema_update.sql
  2. 9
      application/src/main/data/upgrade/3.4.1/schema_update_before.sql
  3. 3
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  4. 5
      application/src/main/java/org/thingsboard/server/controller/AssetController.java
  5. 8
      application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java
  6. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  7. 64
      application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java
  8. 58
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java
  9. 14
      application/src/main/java/org/thingsboard/server/service/entitiy/asset/DefaultTbAssetService.java
  10. 37
      application/src/main/java/org/thingsboard/server/service/install/CassandraKeyspaceService.java
  11. 26
      application/src/main/java/org/thingsboard/server/service/install/DbUpgradeExecutorService.java
  12. 19
      application/src/main/java/org/thingsboard/server/service/install/NoSqlKeyspaceService.java
  13. 68
      application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java
  14. 5
      application/src/main/java/org/thingsboard/server/service/install/update/DefaultCacheCleanupService.java
  15. 13
      application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java
  16. 24
      application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
  17. 22
      application/src/main/java/org/thingsboard/server/service/ttl/EdgeEventsCleanUpService.java
  18. 11
      application/src/main/resources/thingsboard.yml
  19. 69
      application/src/test/java/org/thingsboard/server/controller/BaseEdgeEventControllerTest.java
  20. 7
      application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java
  21. 98
      application/src/test/java/org/thingsboard/server/service/script/MvelInvokeServiceTest.java
  22. 48
      application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java
  23. 1
      application/src/test/java/org/thingsboard/server/transport/TransportNoSqlTestSuite.java
  24. 5
      application/src/test/resources/application-test.properties
  25. 2
      application/src/test/resources/logback-test.xml
  26. 16
      common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/AbstractCassandraCluster.java
  27. 22
      common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaDriverContext.java
  28. 14
      common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaSessionBuilder.java
  29. 27
      common/dao-api/src/main/java/org/thingsboard/server/dao/util/NoSqlAnyDaoNonCloud.java
  30. 2
      common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java
  31. 26
      common/data/src/main/java/org/thingsboard/server/common/data/util/TbPair.java
  32. 5
      common/edge-api/src/main/proto/edge.proto
  33. 4
      common/script/script-api/pom.xml
  34. 78
      common/script/script-api/src/main/java/org/thingsboard/script/api/mvel/DefaultMvelInvokeService.java
  35. 1
      common/script/script-api/src/main/java/org/thingsboard/script/api/mvel/MvelScript.java
  36. 3
      dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java
  37. 9
      dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java
  38. 2
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java
  39. 2
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java
  40. 18
      dao/src/main/java/org/thingsboard/server/dao/sql/asset/AssetRepository.java
  41. 9
      dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java
  42. 83
      dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java
  43. 17
      dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java
  44. 1
      dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java
  45. 5
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java
  46. 4
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java
  47. 15
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TsKvTimescaleRepository.java
  48. 11
      dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvRepository.java
  49. 21
      dao/src/main/resources/cassandra/schema-keyspace.cql
  50. 6
      dao/src/main/resources/cassandra/schema-ts-latest.cql
  51. 6
      dao/src/main/resources/cassandra/schema-ts.cql
  52. 2
      dao/src/main/resources/sql/schema-entities-idx.sql
  53. 4
      dao/src/main/resources/sql/schema-entities.sql
  54. 1
      dao/src/test/java/org/thingsboard/server/dao/NoSqlDaoServiceTestSuite.java
  55. 3
      docker/tb-js-executor.env
  56. 2
      msa/js-executor/api/jsExecutor.models.ts
  57. 11
      msa/js-executor/api/jsInvokeMessageProcessor.ts
  58. 1
      msa/js-executor/config/custom-environment-variables.yml
  59. 1
      msa/js-executor/config/default.yml
  60. 2
      msa/js-executor/docker/start-js-executor.sh
  61. 2
      pom.xml
  62. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java
  63. 55
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java
  64. 119
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNodeTest.java
  65. 2
      tools/src/main/java/org/thingsboard/client/tools/migrator/README.md
  66. 4
      ui-ngx/src/app/core/http/attribute.service.ts
  67. 22
      ui-ngx/src/app/modules/home/components/dashboard-page/layout/manage-dashboard-layouts-dialog.component.html
  68. 4
      ui-ngx/src/app/modules/home/components/dashboard-page/layout/manage-dashboard-layouts-dialog.component.scss
  69. 31
      ui-ngx/src/app/modules/home/components/dashboard-page/layout/manage-dashboard-layouts-dialog.component.ts

62
application/src/main/data/upgrade/3.4.1/schema_update.sql

@ -14,6 +14,7 @@
-- limitations under the License. -- limitations under the License.
-- --
-- AUDIT LOGS MIGRATION START
DO DO
$$ $$
DECLARE table_partition RECORD; DECLARE table_partition RECORD;
@ -73,3 +74,64 @@ BEGIN
WHERE created_time >= start_time_ms AND created_time < end_time_ms; WHERE created_time >= start_time_ms AND created_time < end_time_ms;
END; END;
$$; $$;
-- AUDIT LOGS MIGRATION END
-- EDGE EVENTS MIGRATION START
DO
$$
DECLARE table_partition RECORD;
BEGIN
-- in case of running the upgrade script a second time:
IF NOT (SELECT exists(SELECT FROM pg_tables WHERE tablename = 'old_edge_event')) THEN
ALTER TABLE edge_event RENAME TO old_edge_event;
ALTER INDEX IF EXISTS idx_edge_event_tenant_id_and_created_time RENAME TO idx_old_edge_event_tenant_id_and_created_time;
FOR table_partition IN SELECT tablename AS name, split_part(tablename, '_', 3) AS partition_ts
FROM pg_tables WHERE tablename LIKE 'edge_event_%'
LOOP
EXECUTE format('ALTER TABLE %s RENAME TO old_edge_event_%s', table_partition.name, table_partition.partition_ts);
END LOOP;
ELSE
RAISE NOTICE 'Table old_edge_event already exists, leaving as is';
END IF;
END;
$$;
CREATE TABLE IF NOT EXISTS edge_event (
id uuid NOT NULL,
created_time bigint NOT NULL,
edge_id uuid,
edge_event_type varchar(255),
edge_event_uid varchar(255),
entity_id uuid,
edge_event_action varchar(255),
body varchar(10000000),
tenant_id uuid,
ts bigint NOT NULL
) PARTITION BY RANGE (created_time);
CREATE INDEX IF NOT EXISTS idx_edge_event_tenant_id_and_created_time ON edge_event(tenant_id, created_time DESC);
CREATE OR REPLACE PROCEDURE migrate_edge_event(IN start_time_ms BIGINT, IN end_time_ms BIGINT, IN partition_size_ms BIGINT)
LANGUAGE plpgsql AS
$$
DECLARE
p RECORD;
partition_end_ts BIGINT;
BEGIN
FOR p IN SELECT DISTINCT (created_time - created_time % partition_size_ms) AS partition_ts FROM old_edge_event
WHERE created_time >= start_time_ms AND created_time < end_time_ms
LOOP
partition_end_ts = p.partition_ts + partition_size_ms;
RAISE NOTICE '[edge_event] Partition to create : [%-%]', p.partition_ts, partition_end_ts;
EXECUTE format('CREATE TABLE IF NOT EXISTS edge_event_%s PARTITION OF edge_event ' ||
'FOR VALUES FROM ( %s ) TO ( %s )', p.partition_ts, p.partition_ts, partition_end_ts);
END LOOP;
INSERT INTO edge_event
SELECT id, created_time, edge_id, edge_event_type, edge_event_uid, entity_id, edge_event_action, body, tenant_id, ts
FROM old_edge_event
WHERE created_time >= start_time_ms AND created_time < end_time_ms;
END;
$$;
-- EDGE EVENTS MIGRATION END

9
application/src/main/data/upgrade/3.4.1/schema_update_before.sql

@ -37,9 +37,10 @@ CREATE OR REPLACE PROCEDURE update_asset_profiles()
LANGUAGE plpgsql AS LANGUAGE plpgsql AS
$$ $$
BEGIN BEGIN
UPDATE asset as a SET asset_profile_id = p.id UPDATE asset a SET asset_profile_id = COALESCE(
FROM (SELECT id from asset_profile p WHERE p.tenant_id = a.tenant_id AND a.type = p.name),
(SELECT id, tenant_id, name from asset_profile) as p (SELECT id from asset_profile p WHERE p.tenant_id = a.tenant_id AND p.name = 'default')
WHERE a.asset_profile_id IS NULL AND p.tenant_id = a.tenant_id AND a.type = p.name; )
WHERE a.asset_profile_id IS NULL;
END; END;
$$; $$;

3
application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java

@ -823,6 +823,9 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
body.put("expirationTime", msg.getExpirationTime()); body.put("expirationTime", msg.getExpirationTime());
body.put("method", msg.getBody().getMethod()); body.put("method", msg.getBody().getMethod());
body.put("params", msg.getBody().getParams()); body.put("params", msg.getBody().getParams());
body.put("persisted", msg.isPersisted());
body.put("retries", msg.getRetries());
body.put("additionalInfo", msg.getAdditionalInfo());
EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(tenantId, edgeId, EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL, deviceId, body); EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(tenantId, edgeId, EdgeEventType.DEVICE, EdgeEventActionType.RPC_CALL, deviceId, body);

5
application/src/main/java/org/thingsboard/server/controller/AssetController.java

@ -38,7 +38,6 @@ import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.asset.AssetInfo; import org.thingsboard.server.common.data.asset.AssetInfo;
import org.thingsboard.server.common.data.asset.AssetSearchQuery; import org.thingsboard.server.common.data.asset.AssetSearchQuery;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode;
import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.AssetProfileId; import org.thingsboard.server.common.data.id.AssetProfileId;
@ -86,7 +85,6 @@ import static org.thingsboard.server.controller.ControllerConstants.TENANT_AUTHO
import static org.thingsboard.server.controller.ControllerConstants.TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH; import static org.thingsboard.server.controller.ControllerConstants.TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH;
import static org.thingsboard.server.controller.ControllerConstants.UUID_WIKI_LINK; import static org.thingsboard.server.controller.ControllerConstants.UUID_WIKI_LINK;
import static org.thingsboard.server.controller.EdgeController.EDGE_ID; import static org.thingsboard.server.controller.EdgeController.EDGE_ID;
import static org.thingsboard.server.dao.asset.BaseAssetService.TB_SERVICE_QUEUE;
@RestController @RestController
@TbCoreComponent @TbCoreComponent
@ -148,9 +146,6 @@ public class AssetController extends BaseController {
@RequestMapping(value = "/asset", method = RequestMethod.POST) @RequestMapping(value = "/asset", method = RequestMethod.POST)
@ResponseBody @ResponseBody
public Asset saveAsset(@ApiParam(value = "A JSON value representing the asset.") @RequestBody Asset asset) throws Exception { public Asset saveAsset(@ApiParam(value = "A JSON value representing the asset.") @RequestBody Asset asset) throws Exception {
if (TB_SERVICE_QUEUE.equals(asset.getType())) {
throw new ThingsboardException("Unable to save asset with type " + TB_SERVICE_QUEUE, ThingsboardErrorCode.BAD_REQUEST_PARAMS);
}
asset.setTenantId(getTenantId()); asset.setTenantId(getTenantId());
checkEntity(asset.getId(), asset, Resource.ASSET); checkEntity(asset.getId(), asset, Resource.ASSET);
return tbAssetService.save(asset, getCurrentUser()); return tbAssetService.save(asset, getCurrentUser());

8
application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java

@ -26,6 +26,7 @@ import org.thingsboard.server.service.component.ComponentDiscoveryService;
import org.thingsboard.server.service.install.DatabaseEntitiesUpgradeService; import org.thingsboard.server.service.install.DatabaseEntitiesUpgradeService;
import org.thingsboard.server.service.install.DatabaseTsUpgradeService; import org.thingsboard.server.service.install.DatabaseTsUpgradeService;
import org.thingsboard.server.service.install.EntityDatabaseSchemaService; import org.thingsboard.server.service.install.EntityDatabaseSchemaService;
import org.thingsboard.server.service.install.NoSqlKeyspaceService;
import org.thingsboard.server.service.install.SystemDataLoaderService; import org.thingsboard.server.service.install.SystemDataLoaderService;
import org.thingsboard.server.service.install.TsDatabaseSchemaService; import org.thingsboard.server.service.install.TsDatabaseSchemaService;
import org.thingsboard.server.service.install.TsLatestDatabaseSchemaService; import org.thingsboard.server.service.install.TsLatestDatabaseSchemaService;
@ -51,6 +52,9 @@ public class ThingsboardInstallService {
@Autowired @Autowired
private EntityDatabaseSchemaService entityDatabaseSchemaService; private EntityDatabaseSchemaService entityDatabaseSchemaService;
@Autowired(required = false)
private NoSqlKeyspaceService noSqlKeyspaceService;
@Autowired @Autowired
private TsDatabaseSchemaService tsDatabaseSchemaService; private TsDatabaseSchemaService tsDatabaseSchemaService;
@ -252,6 +256,10 @@ public class ThingsboardInstallService {
log.info("Installing DataBase schema for timeseries..."); log.info("Installing DataBase schema for timeseries...");
if (noSqlKeyspaceService != null) {
noSqlKeyspaceService.createDatabaseSchema();
}
tsDatabaseSchemaService.createDatabaseSchema(); tsDatabaseSchemaService.createDatabaseSchema();
if (tsLatestDatabaseSchemaService != null) { if (tsLatestDatabaseSchemaService != null) {

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java

@ -611,7 +611,7 @@ public final class EdgeGrpcSession implements Closeable {
} }
if (uplinkMsg.getDeviceRpcCallMsgCount() > 0) { if (uplinkMsg.getDeviceRpcCallMsgCount() > 0) {
for (DeviceRpcCallMsg deviceRpcCallMsg : uplinkMsg.getDeviceRpcCallMsgList()) { for (DeviceRpcCallMsg deviceRpcCallMsg : uplinkMsg.getDeviceRpcCallMsgList()) {
result.add(ctx.getDeviceProcessor().processDeviceRpcCallResponseFromEdge(edge.getTenantId(), deviceRpcCallMsg)); result.add(ctx.getDeviceProcessor().processDeviceRpcCallFromEdge(edge.getTenantId(), edge, deviceRpcCallMsg));
} }
} }
if (uplinkMsg.getWidgetBundleTypesRequestMsgCount() > 0) { if (uplinkMsg.getWidgetBundleTypesRequestMsgCount() > 0) {

64
application/src/main/java/org/thingsboard/server/service/edge/rpc/constructor/DeviceMsgConstructor.java

@ -21,13 +21,13 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg;
import org.thingsboard.server.gen.edge.v1.DeviceRpcCallMsg; import org.thingsboard.server.gen.edge.v1.DeviceRpcCallMsg;
import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg; import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg;
import org.thingsboard.server.gen.edge.v1.RpcRequestMsg; import org.thingsboard.server.gen.edge.v1.RpcRequestMsg;
import org.thingsboard.server.gen.edge.v1.RpcResponseMsg;
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; import org.thingsboard.server.gen.edge.v1.UpdateMsgType;
import org.thingsboard.server.queue.util.DataDecodingEncodingService; import org.thingsboard.server.queue.util.DataDecodingEncodingService;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
@ -97,25 +97,55 @@ public class DeviceMsgConstructor {
} }
public DeviceRpcCallMsg constructDeviceRpcCallMsg(UUID deviceId, JsonNode body) { public DeviceRpcCallMsg constructDeviceRpcCallMsg(UUID deviceId, JsonNode body) {
int requestId = body.get("requestId").asInt(); DeviceRpcCallMsg.Builder builder = constructDeviceRpcMsg(deviceId, body);
boolean oneway = body.get("oneway").asBoolean(); if (body.has("error") || body.has("response")) {
UUID requestUUID = UUID.fromString(body.get("requestUUID").asText()); RpcResponseMsg.Builder responseBuilder = RpcResponseMsg.newBuilder();
long expirationTime = body.get("expirationTime").asLong(); if (body.has("error")) {
String method = body.get("method").asText(); responseBuilder.setError(body.get("error").asText());
String params = body.get("params").asText(); } else {
responseBuilder.setResponse(body.get("response").asText());
}
builder.setResponseMsg(responseBuilder.build());
} else {
RpcRequestMsg.Builder requestBuilder = RpcRequestMsg.newBuilder();
requestBuilder.setMethod(body.get("method").asText());
requestBuilder.setParams(body.get("params").asText());
builder.setRequestMsg(requestBuilder.build());
}
return builder.build();
}
RpcRequestMsg.Builder requestBuilder = RpcRequestMsg.newBuilder(); private DeviceRpcCallMsg.Builder constructDeviceRpcMsg(UUID deviceId, JsonNode body) {
requestBuilder.setMethod(method);
requestBuilder.setParams(params);
DeviceRpcCallMsg.Builder builder = DeviceRpcCallMsg.newBuilder() DeviceRpcCallMsg.Builder builder = DeviceRpcCallMsg.newBuilder()
.setDeviceIdMSB(deviceId.getMostSignificantBits()) .setDeviceIdMSB(deviceId.getMostSignificantBits())
.setDeviceIdLSB(deviceId.getLeastSignificantBits()) .setDeviceIdLSB(deviceId.getLeastSignificantBits())
.setRequestUuidMSB(requestUUID.getMostSignificantBits()) .setRequestId(body.get("requestId").asInt());
.setRequestUuidLSB(requestUUID.getLeastSignificantBits()) if (body.get("oneway") != null) {
.setRequestId(requestId) builder.setOneway(body.get("oneway").asBoolean());
.setExpirationTime(expirationTime) }
.setOneway(oneway) if (body.get("requestUUID") != null) {
.setRequestMsg(requestBuilder.build()); UUID requestUUID = UUID.fromString(body.get("requestUUID").asText());
return builder.build(); builder.setRequestUuidMSB(requestUUID.getMostSignificantBits())
.setRequestUuidLSB(requestUUID.getLeastSignificantBits());
}
if (body.get("expirationTime") != null) {
builder.setExpirationTime(body.get("expirationTime").asLong());
}
if (body.get("persisted") != null) {
builder.setPersisted(body.get("persisted").asBoolean());
}
if (body.get("retries") != null) {
builder.setRetries(body.get("retries").asInt());
}
if (body.get("additionalInfo") != null) {
builder.setAdditionalInfo(JacksonUtil.toString(body.get("additionalInfo")));
}
if (body.get("serviceId") != null) {
builder.setServiceId(body.get("serviceId").asText());
}
if (body.get("sessionId") != null) {
builder.setSessionId(body.get("sessionId").asText());
}
return builder;
} }
} }

58
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/DeviceEdgeProcessor.java

@ -52,6 +52,7 @@ import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType; import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse; import org.thingsboard.server.common.msg.rpc.FromDeviceRpcResponse;
import org.thingsboard.server.common.msg.session.SessionMsgType;
import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.gen.edge.v1.DeviceCredentialsRequestMsg; import org.thingsboard.server.gen.edge.v1.DeviceCredentialsRequestMsg;
import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg; import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg;
@ -325,8 +326,17 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor {
return metaData; return metaData;
} }
public ListenableFuture<Void> processDeviceRpcCallResponseFromEdge(TenantId tenantId, DeviceRpcCallMsg deviceRpcCallMsg) { public ListenableFuture<Void> processDeviceRpcCallFromEdge(TenantId tenantId, Edge edge, DeviceRpcCallMsg deviceRpcCallMsg) {
log.trace("[{}] processDeviceRpcCallResponseMsg [{}]", tenantId, deviceRpcCallMsg); log.trace("[{}] processDeviceRpcCallFromEdge [{}]", tenantId, deviceRpcCallMsg);
if (deviceRpcCallMsg.hasResponseMsg()) {
return processDeviceRpcResponseFromEdge(tenantId, deviceRpcCallMsg);
} else if (deviceRpcCallMsg.hasRequestMsg()) {
return processDeviceRpcRequestFromEdge(tenantId, edge, deviceRpcCallMsg);
}
return Futures.immediateFuture(null);
}
private ListenableFuture<Void> processDeviceRpcResponseFromEdge(TenantId tenantId, DeviceRpcCallMsg deviceRpcCallMsg) {
SettableFuture<Void> futureToSet = SettableFuture.create(); SettableFuture<Void> futureToSet = SettableFuture.create();
UUID requestUuid = new UUID(deviceRpcCallMsg.getRequestUuidMSB(), deviceRpcCallMsg.getRequestUuidLSB()); UUID requestUuid = new UUID(deviceRpcCallMsg.getRequestUuidMSB(), deviceRpcCallMsg.getRequestUuidLSB());
DeviceId deviceId = new DeviceId(new UUID(deviceRpcCallMsg.getDeviceIdMSB(), deviceRpcCallMsg.getDeviceIdLSB())); DeviceId deviceId = new DeviceId(new UUID(deviceRpcCallMsg.getDeviceIdMSB(), deviceRpcCallMsg.getDeviceIdLSB()));
@ -357,6 +367,46 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor {
return futureToSet; return futureToSet;
} }
private ListenableFuture<Void> processDeviceRpcRequestFromEdge(TenantId tenantId, Edge edge, DeviceRpcCallMsg deviceRpcCallMsg) {
DeviceId deviceId = new DeviceId(new UUID(deviceRpcCallMsg.getDeviceIdMSB(), deviceRpcCallMsg.getDeviceIdLSB()));
try {
TbMsgMetaData metaData = new TbMsgMetaData();
String requestId = Integer.toString(deviceRpcCallMsg.getRequestId());
metaData.putValue("requestId", requestId);
metaData.putValue("serviceId", deviceRpcCallMsg.getServiceId());
metaData.putValue("sessionId", deviceRpcCallMsg.getSessionId());
metaData.putValue(DataConstants.EDGE_ID, edge.getId().toString());
Device device = deviceService.findDeviceById(tenantId, deviceId);
if (device != null) {
metaData.putValue("deviceName", device.getName());
metaData.putValue("deviceType", device.getType());
metaData.putValue(DataConstants.DEVICE_ID, deviceId.getId().toString());
}
ObjectNode data = JacksonUtil.OBJECT_MAPPER.createObjectNode();
data.put("method", deviceRpcCallMsg.getRequestMsg().getMethod());
data.put("params", deviceRpcCallMsg.getRequestMsg().getParams());
TbMsg tbMsg = TbMsg.newMsg(SessionMsgType.TO_SERVER_RPC_REQUEST.name(), deviceId, null, metaData,
TbMsgDataType.JSON, JacksonUtil.OBJECT_MAPPER.writeValueAsString(data));
tbClusterService.pushMsgToRuleEngine(tenantId, deviceId, tbMsg, new TbQueueCallback() {
@Override
public void onSuccess(TbQueueMsgMetadata metadata) {
log.debug("Successfully send TO_SERVER_RPC_REQUEST to rule engine [{}], deviceRpcCallMsg {}",
device, deviceRpcCallMsg);
}
@Override
public void onFailure(Throwable t) {
log.debug("Failed to send TO_SERVER_RPC_REQUEST to rule engine [{}], deviceRpcCallMsg {}",
device, deviceRpcCallMsg, t);
}
});
} catch (JsonProcessingException | IllegalArgumentException e) {
log.warn("[{}] Failed to push TO_SERVER_RPC_REQUEST to rule engine. deviceRpcCallMsg {}", deviceId, deviceRpcCallMsg, e);
}
return Futures.immediateFuture(null);
}
public DownlinkMsg convertDeviceEventToDownlink(EdgeEvent edgeEvent) { public DownlinkMsg convertDeviceEventToDownlink(EdgeEvent edgeEvent) {
DeviceId deviceId = new DeviceId(edgeEvent.getEntityId()); DeviceId deviceId = new DeviceId(edgeEvent.getEntityId());
DownlinkMsg downlinkMsg = null; DownlinkMsg downlinkMsg = null;
@ -413,11 +463,9 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor {
private DownlinkMsg convertRpcCallEventToDownlink(EdgeEvent edgeEvent) { private DownlinkMsg convertRpcCallEventToDownlink(EdgeEvent edgeEvent) {
log.trace("Executing convertRpcCallEventToDownlink, edgeEvent [{}]", edgeEvent); log.trace("Executing convertRpcCallEventToDownlink, edgeEvent [{}]", edgeEvent);
DeviceRpcCallMsg deviceRpcCallMsg =
deviceMsgConstructor.constructDeviceRpcCallMsg(edgeEvent.getEntityId(), edgeEvent.getBody());
return DownlinkMsg.newBuilder() return DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) .setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addDeviceRpcCallMsg(deviceRpcCallMsg) .addDeviceRpcCallMsg(deviceMsgConstructor.constructDeviceRpcCallMsg(edgeEvent.getEntityId(), edgeEvent.getBody()))
.build(); .build();
} }

14
application/src/main/java/org/thingsboard/server/service/entitiy/asset/DefaultTbAssetService.java

@ -22,8 +22,10 @@ import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.asset.AssetProfile;
import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode;
import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
@ -31,20 +33,32 @@ import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.dao.asset.AssetService; import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.service.entitiy.AbstractTbEntityService; import org.thingsboard.server.service.entitiy.AbstractTbEntityService;
import org.thingsboard.server.service.profile.TbAssetProfileCache;
import java.util.List; import java.util.List;
import static org.thingsboard.server.dao.asset.BaseAssetService.TB_SERVICE_QUEUE;
@Service @Service
@AllArgsConstructor @AllArgsConstructor
public class DefaultTbAssetService extends AbstractTbEntityService implements TbAssetService { public class DefaultTbAssetService extends AbstractTbEntityService implements TbAssetService {
private final AssetService assetService; private final AssetService assetService;
private final TbAssetProfileCache assetProfileCache;
@Override @Override
public Asset save(Asset asset, User user) throws Exception { public Asset save(Asset asset, User user) throws Exception {
ActionType actionType = asset.getId() == null ? ActionType.ADDED : ActionType.UPDATED; ActionType actionType = asset.getId() == null ? ActionType.ADDED : ActionType.UPDATED;
TenantId tenantId = asset.getTenantId(); TenantId tenantId = asset.getTenantId();
try { try {
if (TB_SERVICE_QUEUE.equals(asset.getType())) {
throw new ThingsboardException("Unable to save asset with type " + TB_SERVICE_QUEUE, ThingsboardErrorCode.BAD_REQUEST_PARAMS);
} else if (asset.getAssetProfileId() != null) {
AssetProfile assetProfile = assetProfileCache.get(tenantId, asset.getAssetProfileId());
if (assetProfile != null && TB_SERVICE_QUEUE.equals(assetProfile.getName())) {
throw new ThingsboardException("Unable to save asset with profile " + TB_SERVICE_QUEUE, ThingsboardErrorCode.BAD_REQUEST_PARAMS);
}
}
Asset savedAsset = checkNotNull(assetService.saveAsset(asset)); Asset savedAsset = checkNotNull(assetService.saveAsset(asset));
autoCommit(user, savedAsset.getId()); autoCommit(user, savedAsset.getId());
notificationEntityService.notifyCreateOrUpdateEntity(tenantId, savedAsset.getId(), savedAsset, notificationEntityService.notifyCreateOrUpdateEntity(tenantId, savedAsset.getId(), savedAsset,

37
application/src/main/java/org/thingsboard/server/service/install/CassandraKeyspaceService.java

@ -0,0 +1,37 @@
/**
* Copyright © 2016-2022 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.install;
import org.springframework.context.annotation.Profile;
import org.springframework.stereotype.Service;
import org.thingsboard.server.dao.util.NoSqlAnyDaoNonCloud;
/*
* Create keyspace for Cassandra NoSQL database for non-cloud deployment.
* For cloud service like Astra DBaas admin have to create keyspace manually on cloud UI.
* Then create tokens with database admin role and put it on Thingsboard parameters.
* Without this service cloud DB will end up with exception like
* UnauthorizedException: Missing correct permission on thingsboard
* */
@Service
@NoSqlAnyDaoNonCloud
@Profile("install")
public class CassandraKeyspaceService extends CassandraAbstractDatabaseSchemaService
implements NoSqlKeyspaceService {
public CassandraKeyspaceService() {
super("schema-keyspace.cql");
}
}

26
application/src/main/java/org/thingsboard/server/service/install/DbUpgradeExecutorService.java

@ -0,0 +1,26 @@
/**
* Copyright © 2016-2022 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.install;
import org.springframework.context.annotation.Profile;
import org.springframework.stereotype.Component;
import org.thingsboard.server.service.executors.DbCallbackExecutorService;
@Component
@Profile("install")
public class DbUpgradeExecutorService extends DbCallbackExecutorService {
}

19
application/src/main/java/org/thingsboard/server/service/install/NoSqlKeyspaceService.java

@ -0,0 +1,19 @@
/**
* Copyright © 2016-2022 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.install;
public interface NoSqlKeyspaceService extends DatabaseSchemaService {
}

68
application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java

@ -15,6 +15,8 @@
*/ */
package org.thingsboard.server.service.install; package org.thingsboard.server.service.install;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils; import org.apache.commons.collections.CollectionUtils;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
@ -32,14 +34,14 @@ import org.thingsboard.server.common.data.queue.ProcessingStrategyType;
import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.data.queue.Queue;
import org.thingsboard.server.common.data.queue.SubmitStrategy; import org.thingsboard.server.common.data.queue.SubmitStrategy;
import org.thingsboard.server.common.data.queue.SubmitStrategyType; import org.thingsboard.server.common.data.queue.SubmitStrategyType;
import org.thingsboard.server.dao.asset.AssetDao;
import org.thingsboard.server.dao.asset.AssetProfileService; import org.thingsboard.server.dao.asset.AssetProfileService;
import org.thingsboard.server.dao.asset.AssetService;
import org.thingsboard.server.dao.dashboard.DashboardService; import org.thingsboard.server.dao.dashboard.DashboardService;
import org.thingsboard.server.dao.device.DeviceProfileService; import org.thingsboard.server.dao.device.DeviceProfileService;
import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.common.data.util.TbPair;
import org.thingsboard.server.dao.queue.QueueService; import org.thingsboard.server.dao.queue.QueueService;
import org.thingsboard.server.dao.rule.RuleChainService; import org.thingsboard.server.dao.sql.tenant.TenantRepository;
import org.thingsboard.server.dao.tenant.TenantProfileService;
import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.dao.usagerecord.ApiUsageStateService; import org.thingsboard.server.dao.usagerecord.ApiUsageStateService;
import org.thingsboard.server.queue.settings.TbRuleEngineQueueConfiguration; import org.thingsboard.server.queue.settings.TbRuleEngineQueueConfiguration;
@ -56,7 +58,9 @@ import java.sql.SQLException;
import java.sql.SQLSyntaxErrorException; import java.sql.SQLSyntaxErrorException;
import java.sql.SQLWarning; import java.sql.SQLWarning;
import java.sql.Statement; import java.sql.Statement;
import java.util.ArrayList;
import java.util.List; import java.util.List;
import java.util.UUID;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import static org.thingsboard.server.service.install.DatabaseHelper.ADDITIONAL_INFO; import static org.thingsboard.server.service.install.DatabaseHelper.ADDITIONAL_INFO;
@ -106,11 +110,14 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
@Autowired @Autowired
private TenantService tenantService; private TenantService tenantService;
@Autowired
private TenantRepository tenantRepository;
@Autowired @Autowired
private DeviceService deviceService; private DeviceService deviceService;
@Autowired @Autowired
private AssetService assetService; private AssetDao assetDao;
@Autowired @Autowired
private DeviceProfileService deviceProfileService; private DeviceProfileService deviceProfileService;
@ -129,10 +136,7 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
private TbRuleEngineQueueConfigService queueConfig; private TbRuleEngineQueueConfigService queueConfig;
@Autowired @Autowired
private RuleChainService ruleChainService; private DbUpgradeExecutorService dbUpgradeExecutor;
@Autowired
private TenantProfileService tenantProfileService;
@Override @Override
public void upgradeDatabase(String fromVersion) throws Exception { public void upgradeDatabase(String fromVersion) throws Exception {
@ -620,26 +624,45 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "3.4.1", "schema_update_before.sql"); schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "3.4.1", "schema_update_before.sql");
loadSql(schemaUpdateFile, conn); loadSql(schemaUpdateFile, conn);
conn.createStatement().execute("DELETE FROM asset a WHERE NOT exists(SELECT id FROM tenant WHERE id = a.tenant_id);");
log.info("Creating default asset profiles..."); log.info("Creating default asset profiles...");
PageLink pageLink = new PageLink(100);
PageData<Tenant> pageData; PageLink pageLink = new PageLink(1000);
PageData<TenantId> tenantIds;
do { do {
pageData = tenantService.findTenants(pageLink); List<ListenableFuture<?>> futures = new ArrayList<>();
for (Tenant tenant : pageData.getData()) { tenantIds = tenantService.findTenantsIds(pageLink);
List<EntitySubtype> assetTypes = assetService.findAssetTypesByTenantId(tenant.getId()).get(); for (TenantId tenantId : tenantIds.getData()) {
try { futures.add(dbUpgradeExecutor.submit(() -> {
assetProfileService.createDefaultAssetProfile(tenant.getId());
} catch (Exception e) {
}
for (EntitySubtype assetType : assetTypes) {
try { try {
assetProfileService.findOrCreateAssetProfile(tenant.getId(), assetType.getType()); assetProfileService.createDefaultAssetProfile(tenantId);
} catch (Exception e) { } catch (Exception e) {}
} }));
}
Futures.allAsList(futures).get();
pageLink = pageLink.nextPageLink();
} while (tenantIds.hasNext());
pageLink = new PageLink(1000);
PageData<TbPair<UUID, String>> pairs;
do {
List<ListenableFuture<?>> futures = new ArrayList<>();
pairs = assetDao.getAllAssetTypes(pageLink);
for (TbPair<UUID, String> pair : pairs.getData()) {
TenantId tenantId = new TenantId(pair.getFirst());
String assetType = pair.getSecond();
if (!"default".equals(assetType)) {
futures.add(dbUpgradeExecutor.submit(() -> {
try {
assetProfileService.findOrCreateAssetProfile(tenantId, assetType);
} catch (Exception e) {}
}));
} }
} }
Futures.allAsList(futures).get();
pageLink = pageLink.nextPageLink(); pageLink = pageLink.nextPageLink();
} while (pageData.hasNext()); } while (pairs.hasNext());
log.info("Updating asset profiles..."); log.info("Updating asset profiles...");
conn.createStatement().execute("call update_asset_profiles()"); conn.createStatement().execute("call update_asset_profiles()");
@ -728,5 +751,4 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
return queue; return queue;
} }
} }

5
application/src/main/java/org/thingsboard/server/service/install/update/DefaultCacheCleanupService.java

@ -73,6 +73,11 @@ public class DefaultCacheCleanupService implements CacheCleanupService {
log.info("Clear cache to upgrade from version 3.3.4 to 3.4.0 ..."); log.info("Clear cache to upgrade from version 3.3.4 to 3.4.0 ...");
clearAll(); clearAll();
break; break;
case "3.4.1":
log.info("Clear cache to upgrade from version 3.4.1 to 3.4.2 ...");
clearCacheByName("assets");
clearCacheByName("repositorySettings");
break;
default: default:
//Do nothing, since cache cleanup is optional. //Do nothing, since cache cleanup is optional.
} }

13
application/src/main/java/org/thingsboard/server/service/install/update/DefaultDataUpdateService.java

@ -64,6 +64,7 @@ import org.thingsboard.server.common.data.tenant.profile.TenantProfileQueueConfi
import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.alarm.AlarmDao; import org.thingsboard.server.dao.alarm.AlarmDao;
import org.thingsboard.server.dao.audit.AuditLogDao; import org.thingsboard.server.dao.audit.AuditLogDao;
import org.thingsboard.server.dao.edge.EdgeEventDao;
import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.entityview.EntityViewService;
import org.thingsboard.server.dao.event.EventService; import org.thingsboard.server.dao.event.EventService;
@ -142,6 +143,9 @@ public class DefaultDataUpdateService implements DataUpdateService {
@Autowired @Autowired
private AuditLogDao auditLogDao; private AuditLogDao auditLogDao;
@Autowired
private EdgeEventDao edgeEventDao;
@Override @Override
public void updateData(String fromVersion) throws Exception { public void updateData(String fromVersion) throws Exception {
switch (fromVersion) { switch (fromVersion) {
@ -181,14 +185,21 @@ public class DefaultDataUpdateService implements DataUpdateService {
} }
break; break;
case "3.4.1": case "3.4.1":
log.info("Updating data from version 3.4.1 to 3.4.2 ...");
boolean skipAuditLogsMigration = getEnv("TB_SKIP_AUDIT_LOGS_MIGRATION", false); boolean skipAuditLogsMigration = getEnv("TB_SKIP_AUDIT_LOGS_MIGRATION", false);
if (!skipAuditLogsMigration) { if (!skipAuditLogsMigration) {
log.info("Updating data from version 3.4.1 to 3.4.2 ...");
log.info("Starting audit logs migration. Can be skipped with TB_SKIP_AUDIT_LOGS_MIGRATION env variable set to true"); log.info("Starting audit logs migration. Can be skipped with TB_SKIP_AUDIT_LOGS_MIGRATION env variable set to true");
auditLogDao.migrateAuditLogs(); auditLogDao.migrateAuditLogs();
} else { } else {
log.info("Skipping audit logs migration"); log.info("Skipping audit logs migration");
} }
boolean skipEdgeEventsMigration = getEnv("TB_SKIP_EDGE_EVENTS_MIGRATION", false);
if (!skipEdgeEventsMigration) {
log.info("Starting edge events migration. Can be skipped with TB_SKIP_EDGE_EVENTS_MIGRATION env variable set to true");
edgeEventDao.migrateEdgeEvents();
} else {
log.info("Skipping edge events migration");
}
break; break;
default: default:
throw new RuntimeException("Unable to update data, unsupported fromVersion: " + fromVersion); throw new RuntimeException("Unable to update data, unsupported fromVersion: " + fromVersion);

24
application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java

@ -24,6 +24,7 @@ import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.ListeningExecutorService; import com.google.common.util.concurrent.ListeningExecutorService;
import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.MoreExecutors;
import lombok.Getter; import lombok.Getter;
import lombok.Setter;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
@ -121,7 +122,8 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
new EntityKey(EntityKeyType.TIME_SERIES, INACTIVITY_TIMEOUT), new EntityKey(EntityKeyType.TIME_SERIES, INACTIVITY_TIMEOUT),
new EntityKey(EntityKeyType.TIME_SERIES, ACTIVITY_STATE), new EntityKey(EntityKeyType.TIME_SERIES, ACTIVITY_STATE),
new EntityKey(EntityKeyType.TIME_SERIES, LAST_CONNECT_TIME), new EntityKey(EntityKeyType.TIME_SERIES, LAST_CONNECT_TIME),
new EntityKey(EntityKeyType.TIME_SERIES, LAST_DISCONNECT_TIME)); new EntityKey(EntityKeyType.TIME_SERIES, LAST_DISCONNECT_TIME),
new EntityKey(EntityKeyType.SERVER_ATTRIBUTE, INACTIVITY_TIMEOUT));
private static final List<EntityKey> PERSISTENT_ATTRIBUTE_KEYS = Arrays.asList( private static final List<EntityKey> PERSISTENT_ATTRIBUTE_KEYS = Arrays.asList(
new EntityKey(EntityKeyType.SERVER_ATTRIBUTE, LAST_ACTIVITY_TIME), new EntityKey(EntityKeyType.SERVER_ATTRIBUTE, LAST_ACTIVITY_TIME),
@ -152,14 +154,21 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
@Value("${state.defaultInactivityTimeoutInSec}") @Value("${state.defaultInactivityTimeoutInSec}")
@Getter @Getter
@Setter
private long defaultInactivityTimeoutInSec; private long defaultInactivityTimeoutInSec;
@Value("#{${state.defaultInactivityTimeoutInSec} * 1000}")
@Getter
@Setter
private long defaultInactivityTimeoutMs;
@Value("${state.defaultStateCheckIntervalInSec}") @Value("${state.defaultStateCheckIntervalInSec}")
@Getter @Getter
private int defaultStateCheckIntervalInSec; private int defaultStateCheckIntervalInSec;
@Value("${state.persistToTelemetry:false}") @Value("${state.persistToTelemetry:false}")
@Getter @Getter
@Setter
private boolean persistToTelemetry; private boolean persistToTelemetry;
@Value("${state.initFetchPackSize:50000}") @Value("${state.initFetchPackSize:50000}")
@ -540,7 +549,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
private ListenableFuture<DeviceStateData> transformInactivityTimeout(ListenableFuture<DeviceStateData> future) { private ListenableFuture<DeviceStateData> transformInactivityTimeout(ListenableFuture<DeviceStateData> future) {
return Futures.transformAsync(future, deviceStateData -> { return Futures.transformAsync(future, deviceStateData -> {
if (!persistToTelemetry || deviceStateData.getState().getInactivityTimeout() != TimeUnit.SECONDS.toMillis(defaultInactivityTimeoutInSec)) { if (!persistToTelemetry || deviceStateData.getState().getInactivityTimeout() != defaultInactivityTimeoutMs) {
return future; //fail fast return future; //fail fast
} }
var attributesFuture = attributesService.find(TenantId.SYS_TENANT_ID, deviceStateData.getDeviceId(), SERVER_SCOPE, INACTIVITY_TIMEOUT); var attributesFuture = attributesService.find(TenantId.SYS_TENANT_ID, deviceStateData.getDeviceId(), SERVER_SCOPE, INACTIVITY_TIMEOUT);
@ -563,7 +572,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
try { try {
long lastActivityTime = getEntryValue(data, LAST_ACTIVITY_TIME, 0L); long lastActivityTime = getEntryValue(data, LAST_ACTIVITY_TIME, 0L);
long inactivityAlarmTime = getEntryValue(data, INACTIVITY_ALARM_TIME, 0L); long inactivityAlarmTime = getEntryValue(data, INACTIVITY_ALARM_TIME, 0L);
long inactivityTimeout = getEntryValue(data, INACTIVITY_TIMEOUT, TimeUnit.SECONDS.toMillis(defaultInactivityTimeoutInSec)); long inactivityTimeout = getEntryValue(data, INACTIVITY_TIMEOUT, defaultInactivityTimeoutMs);
//Actual active state by wall-clock will updated outside this method. This method is only for fetch persistent state //Actual active state by wall-clock will updated outside this method. This method is only for fetch persistent state
final boolean active = getEntryValue(data, ACTIVITY_STATE, false); final boolean active = getEntryValue(data, ACTIVITY_STATE, false);
DeviceState deviceState = DeviceState.builder() DeviceState deviceState = DeviceState.builder()
@ -634,10 +643,15 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
} }
private DeviceStateData toDeviceStateData(EntityData ed, DeviceIdInfo deviceIdInfo) { DeviceStateData toDeviceStateData(EntityData ed, DeviceIdInfo deviceIdInfo) {
long lastActivityTime = getEntryValue(ed, getKeyType(), LAST_ACTIVITY_TIME, 0L); long lastActivityTime = getEntryValue(ed, getKeyType(), LAST_ACTIVITY_TIME, 0L);
long inactivityAlarmTime = getEntryValue(ed, getKeyType(), INACTIVITY_ALARM_TIME, 0L); long inactivityAlarmTime = getEntryValue(ed, getKeyType(), INACTIVITY_ALARM_TIME, 0L);
long inactivityTimeout = getEntryValue(ed, getKeyType(), INACTIVITY_TIMEOUT, TimeUnit.SECONDS.toMillis(defaultInactivityTimeoutInSec)); long inactivityTimeout = getEntryValue(ed, getKeyType(), INACTIVITY_TIMEOUT, defaultInactivityTimeoutMs);
if (persistToTelemetry && inactivityTimeout == defaultInactivityTimeoutMs) {
log.trace("[{}] default value for inactivity timeout fetched {}, going to fetch inactivity timeout from attributes",
deviceIdInfo.getDeviceId(), inactivityTimeout);
inactivityTimeout = getEntryValue(ed, EntityKeyType.SERVER_ATTRIBUTE, INACTIVITY_TIMEOUT, defaultInactivityTimeoutMs);
}
//Actual active state by wall-clock will updated outside this method. This method is only for fetch persistent state //Actual active state by wall-clock will updated outside this method. This method is only for fetch persistent state
final boolean active = getEntryValue(ed, getKeyType(), ACTIVITY_STATE, false); final boolean active = getEntryValue(ed, getKeyType(), ACTIVITY_STATE, false);
DeviceState deviceState = DeviceState.builder() DeviceState deviceState = DeviceState.builder()

22
application/src/main/java/org/thingsboard/server/service/ttl/EdgeEventsCleanUpService.java

@ -17,15 +17,22 @@ package org.thingsboard.server.service.ttl;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.scheduling.annotation.Scheduled; import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.dao.edge.EdgeEventService; import org.thingsboard.server.dao.edge.EdgeEventService;
import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository;
import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.queue.util.TbCoreComponent;
import java.util.concurrent.TimeUnit;
import static org.thingsboard.server.dao.model.ModelConstants.EDGE_EVENT_COLUMN_FAMILY_NAME;
@TbCoreComponent @TbCoreComponent
@Slf4j @Slf4j
@Service @Service
@ConditionalOnExpression("${sql.ttl.edge_events.enabled:true} && ${sql.ttl.edge_events.edge_event_ttl:0} > 0")
public class EdgeEventsCleanUpService extends AbstractCleanUpService { public class EdgeEventsCleanUpService extends AbstractCleanUpService {
public static final String RANDOM_DELAY_INTERVAL_MS_EXPRESSION = public static final String RANDOM_DELAY_INTERVAL_MS_EXPRESSION =
@ -34,20 +41,29 @@ public class EdgeEventsCleanUpService extends AbstractCleanUpService {
@Value("${sql.ttl.edge_events.edge_events_ttl}") @Value("${sql.ttl.edge_events.edge_events_ttl}")
private long ttl; private long ttl;
@Value("${sql.ttl.edge_events.enabled}") @Value("${sql.edge_events.partition_size:168}")
private int partitionSizeInHours;
@Value("${sql.ttl.edge_events.enabled:true}")
private boolean ttlTaskExecutionEnabled; private boolean ttlTaskExecutionEnabled;
private final EdgeEventService edgeEventService; private final EdgeEventService edgeEventService;
public EdgeEventsCleanUpService(PartitionService partitionService, EdgeEventService edgeEventService) { private final SqlPartitioningRepository partitioningRepository;
public EdgeEventsCleanUpService(PartitionService partitionService, EdgeEventService edgeEventService, SqlPartitioningRepository partitioningRepository) {
super(partitionService); super(partitionService);
this.edgeEventService = edgeEventService; this.edgeEventService = edgeEventService;
this.partitioningRepository = partitioningRepository;
} }
@Scheduled(initialDelayString = RANDOM_DELAY_INTERVAL_MS_EXPRESSION, fixedDelayString = "${sql.ttl.edge_events.execution_interval_ms}") @Scheduled(initialDelayString = RANDOM_DELAY_INTERVAL_MS_EXPRESSION, fixedDelayString = "${sql.ttl.edge_events.execution_interval_ms}")
public void cleanUp() { public void cleanUp() {
long edgeEventsExpTime = System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(ttl);
if (ttlTaskExecutionEnabled && isSystemTenantPartitionMine()) { if (ttlTaskExecutionEnabled && isSystemTenantPartitionMine()) {
edgeEventService.cleanupEvents(ttl); edgeEventService.cleanupEvents(edgeEventsExpTime);
} else {
partitioningRepository.cleanupPartitionsCache(EDGE_EVENT_COLUMN_FAMILY_NAME, edgeEventsExpTime, TimeUnit.HOURS.toMillis(partitionSizeInHours));
} }
} }

11
application/src/main/resources/thingsboard.yml

@ -200,6 +200,14 @@ cassandra:
username: "${CASSANDRA_USERNAME:}" username: "${CASSANDRA_USERNAME:}"
# Specify your password # Specify your password
password: "${CASSANDRA_PASSWORD:}" password: "${CASSANDRA_PASSWORD:}"
# Astra DB connect https://astra.datastax.com/
cloud:
# /etc/thingsboard/astra/secure-connect-thingsboard.zip
secure_connect_bundle_path: "${CASSANDRA_CLOUD_SECURE_BUNDLE_PATH:}"
# DucitQPHMzPCBOZqFYexAfKk
client_id: "${CASSANDRA_CLOUD_CLIENT_ID:}"
# ZnF7FpuHp43FP5BzM+KY8wGmSb4Ql6BhT4Z7sOU13ze+gXQ-n7OkFpNuB,oACUIQObQnK0g4bSPoZhK5ejkcF9F.j6f64j71Sr.tiRe0Fsq2hPS1ZCGSfAaIgg63IydG
client_secret: "${CASSANDRA_CLOUD_CLIENT_SECRET:}"
# Cassandra cluster connection socket parameters # # Cassandra cluster connection socket parameters #
socket: socket:
@ -265,6 +273,7 @@ sql:
batch_size: "${SQL_EDGE_EVENTS_BATCH_SIZE:1000}" batch_size: "${SQL_EDGE_EVENTS_BATCH_SIZE:1000}"
batch_max_delay: "${SQL_EDGE_EVENTS_BATCH_MAX_DELAY_MS:100}" batch_max_delay: "${SQL_EDGE_EVENTS_BATCH_MAX_DELAY_MS:100}"
stats_print_interval_ms: "${SQL_EDGE_EVENTS_BATCH_STATS_PRINT_MS:10000}" stats_print_interval_ms: "${SQL_EDGE_EVENTS_BATCH_STATS_PRINT_MS:10000}"
partition_size: "${SQL_EDGE_EVENTS_PARTITION_SIZE_HOURS:168}" # Number of hours to partition the events. The current value corresponds to one week.
audit_logs: audit_logs:
partition_size: "${SQL_AUDIT_LOGS_PARTITION_SIZE_HOURS:168}" # Default value - 1 week partition_size: "${SQL_AUDIT_LOGS_PARTITION_SIZE_HOURS:168}" # Default value - 1 week
# Specify whether to sort entities before batch update. Should be enabled for cluster mode to avoid deadlocks # Specify whether to sort entities before batch update. Should be enabled for cluster mode to avoid deadlocks
@ -534,6 +543,7 @@ spring.servlet.multipart.max-file-size: "50MB"
spring.servlet.multipart.max-request-size: "50MB" spring.servlet.multipart.max-request-size: "50MB"
spring.jpa.properties.hibernate.jdbc.lob.non_contextual_creation: "true" spring.jpa.properties.hibernate.jdbc.lob.non_contextual_creation: "true"
# Note: as for current Spring JPA version, custom NullHandling for the Sort.Order is ignored and this parameter is used
spring.jpa.properties.hibernate.order_by.default_null_ordering: "${SPRING_JPA_PROPERTIES_HIBERNATE_ORDER_BY_DEFAULT_NULL_ORDERING:last}" spring.jpa.properties.hibernate.order_by.default_null_ordering: "${SPRING_JPA_PROPERTIES_HIBERNATE_ORDER_BY_DEFAULT_NULL_ORDERING:last}"
# SQL DAO Configuration # SQL DAO Configuration
@ -615,6 +625,7 @@ mvel:
max_black_list_duration_sec: "${MVEL_MAX_BLACKLIST_DURATION_SEC:60}" max_black_list_duration_sec: "${MVEL_MAX_BLACKLIST_DURATION_SEC:60}"
# Specify thread pool size for javascript executor service # Specify thread pool size for javascript executor service
thread_pool_size: "${MVEL_THREAD_POOL_SIZE:50}" thread_pool_size: "${MVEL_THREAD_POOL_SIZE:50}"
compiled_scripts_cache_size: "${MVEL_COMPILED_SCRIPTS_CACHE_SIZE:1000}"
stats: stats:
enabled: "${TB_MVEL_STATS_ENABLED:false}" enabled: "${TB_MVEL_STATS_ENABLED:false}"
print_interval_ms: "${TB_MVEL_STATS_PRINT_INTERVAL_MS:10000}" print_interval_ms: "${TB_MVEL_STATS_PRINT_INTERVAL_MS:10000}"

69
application/src/test/java/org/thingsboard/server/controller/BaseEdgeEventControllerTest.java

@ -22,6 +22,9 @@ import org.junit.After;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Before; import org.junit.Before;
import org.junit.Test; import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.TestPropertySource; import org.springframework.test.context.TestPropertySource;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.Tenant;
@ -29,16 +32,27 @@ import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.asset.Asset; import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType; import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.dao.edge.EdgeEventDao;
import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository;
import org.thingsboard.server.service.ttl.EdgeEventsCleanUpService;
import java.time.LocalDate;
import java.time.ZoneOffset;
import java.util.List; import java.util.List;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.reset;
import static org.mockito.Mockito.verify;
import static org.assertj.core.api.Assertions.assertThat;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@TestPropertySource(properties = { @TestPropertySource(properties = {
@ -50,6 +64,18 @@ public abstract class BaseEdgeEventControllerTest extends AbstractControllerTest
private Tenant savedTenant; private Tenant savedTenant;
private User tenantAdmin; private User tenantAdmin;
@Autowired
private EdgeEventDao edgeEventDao;
@SpyBean
private SqlPartitioningRepository partitioningRepository;
@Autowired
private EdgeEventsCleanUpService edgeEventsCleanUpService;
@Value("#{${sql.edge_events.partition_size} * 60 * 60 * 1000}")
private long partitionDurationInMs;
@Value("${sql.ttl.edge_events.edge_event_ttl}")
private long edgeEventTtlInSec;
@Before @Before
public void beforeTest() throws Exception { public void beforeTest() throws Exception {
loginSysAdmin(); loginSysAdmin();
@ -114,6 +140,34 @@ public abstract class BaseEdgeEventControllerTest extends AbstractControllerTest
Assert.assertTrue(edgeEvents.stream().anyMatch(ee -> EdgeEventType.RELATION.equals(ee.getType()))); Assert.assertTrue(edgeEvents.stream().anyMatch(ee -> EdgeEventType.RELATION.equals(ee.getType())));
} }
@Test
public void saveEdgeEvent_thenCreatePartitionIfNotExist() {
reset(partitioningRepository);
EdgeEvent edgeEvent = createEdgeEvent();
verify(partitioningRepository).createPartitionIfNotExists(eq("edge_event"), eq(edgeEvent.getCreatedTime()), eq(partitionDurationInMs));
List<Long> partitions = partitioningRepository.fetchPartitions("edge_event");
assertThat(partitions).singleElement().satisfies(partitionStartTs -> {
assertThat(partitionStartTs).isEqualTo(partitioningRepository.calculatePartitionStartTime(edgeEvent.getCreatedTime(), partitionDurationInMs));
});
}
@Test
public void cleanUpEdgeEventByTtl_dropOldPartitions() {
long oldEdgeEventTs = LocalDate.of(2020, 10, 1).atStartOfDay().toInstant(ZoneOffset.UTC).toEpochMilli();
long partitionStartTs = partitioningRepository.calculatePartitionStartTime(oldEdgeEventTs, partitionDurationInMs);
partitioningRepository.createPartitionIfNotExists("edge_event", oldEdgeEventTs, partitionDurationInMs);
List<Long> partitions = partitioningRepository.fetchPartitions("edge_event");
assertThat(partitions).contains(partitionStartTs);
edgeEventsCleanUpService.cleanUp();
partitions = partitioningRepository.fetchPartitions("edge_event");
assertThat(partitions).doesNotContain(partitionStartTs);
assertThat(partitions).allSatisfy(partitionsStart -> {
long partitionEndTs = partitionsStart + partitionDurationInMs;
assertThat(partitionEndTs).isGreaterThan(System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(edgeEventTtlInSec));
});
}
private List<EdgeEvent> findEdgeEvents(EdgeId edgeId) throws Exception { private List<EdgeEvent> findEdgeEvents(EdgeId edgeId) throws Exception {
return doGetTypedWithTimePageLink("/api/edge/" + edgeId.toString() + "/events?", return doGetTypedWithTimePageLink("/api/edge/" + edgeId.toString() + "/events?",
new TypeReference<PageData<EdgeEvent>>() { new TypeReference<PageData<EdgeEvent>>() {
@ -134,4 +188,19 @@ public abstract class BaseEdgeEventControllerTest extends AbstractControllerTest
return asset; return asset;
} }
private EdgeEvent createEdgeEvent() {
EdgeEvent edgeEvent = new EdgeEvent();
edgeEvent.setCreatedTime(System.currentTimeMillis());
edgeEvent.setTenantId(tenantId);
edgeEvent.setAction(EdgeEventActionType.ADDED);
edgeEvent.setEntityId(tenantAdmin.getUuidId());
edgeEvent.setType(EdgeEventType.ALARM);
try {
edgeEventDao.saveAsync(edgeEvent).get();
} catch (InterruptedException | ExecutionException e) {
throw new RuntimeException(e);
}
return edgeEvent;
}
} }

7
application/src/test/java/org/thingsboard/server/edge/BaseDeviceEdgeTest.java

@ -511,8 +511,11 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest {
body.put("expirationTime", System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10)); body.put("expirationTime", System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(10));
body.put("method", "test_method"); body.put("method", "test_method");
body.put("params", "{\"param1\":\"value1\"}"); body.put("params", "{\"param1\":\"value1\"}");
body.put("persisted", true);
body.put("retries", 2);
EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL, device.getId().getId(), EdgeEventType.DEVICE, body); EdgeEvent edgeEvent = constructEdgeEvent(tenantId, edge.getId(), EdgeEventActionType.RPC_CALL,
device.getId().getId(), EdgeEventType.DEVICE, body);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent).get(); edgeEventService.saveAsync(edgeEvent).get();
clusterService.onEdgeEventUpdate(tenantId, edge.getId()); clusterService.onEdgeEventUpdate(tenantId, edge.getId());
@ -522,6 +525,8 @@ abstract public class BaseDeviceEdgeTest extends AbstractEdgeTest {
Assert.assertTrue(latestMessage instanceof DeviceRpcCallMsg); Assert.assertTrue(latestMessage instanceof DeviceRpcCallMsg);
DeviceRpcCallMsg latestDeviceRpcCallMsg = (DeviceRpcCallMsg) latestMessage; DeviceRpcCallMsg latestDeviceRpcCallMsg = (DeviceRpcCallMsg) latestMessage;
Assert.assertEquals("test_method", latestDeviceRpcCallMsg.getRequestMsg().getMethod()); Assert.assertEquals("test_method", latestDeviceRpcCallMsg.getRequestMsg().getMethod());
Assert.assertTrue(latestDeviceRpcCallMsg.getPersisted());
Assert.assertEquals(2, latestDeviceRpcCallMsg.getRetries());
} }
private void sendAttributesRequestAndVerify(Device device, String scope, String attributesDataStr, String expectedKey, private void sendAttributesRequestAndVerify(Device device, String scope, String attributesDataStr, String expectedKey,

98
application/src/test/java/org/thingsboard/server/service/script/MvelInvokeServiceTest.java

@ -16,6 +16,7 @@
package org.thingsboard.server.service.script; package org.thingsboard.server.service.script;
import com.fasterxml.jackson.databind.node.ObjectNode; import com.fasterxml.jackson.databind.node.ObjectNode;
import com.github.benmanes.caffeine.cache.Cache;
import org.junit.Assert; import org.junit.Assert;
import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
@ -24,15 +25,22 @@ import org.springframework.test.context.TestPropertySource;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.script.api.ScriptType; import org.thingsboard.script.api.ScriptType;
import org.thingsboard.script.api.mvel.MvelInvokeService; import org.thingsboard.script.api.mvel.MvelInvokeService;
import org.thingsboard.script.api.mvel.MvelScript;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.controller.AbstractControllerTest; import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.dao.service.DaoSqlTest;
import java.io.Serializable;
import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.assertj.core.api.Assertions.assertThatThrownBy;
@DaoSqlTest @DaoSqlTest
@ -41,6 +49,7 @@ import static org.assertj.core.api.Assertions.assertThatThrownBy;
"mvel.max_total_args_size=50", "mvel.max_total_args_size=50",
"mvel.max_result_size=50", "mvel.max_result_size=50",
"mvel.max_errors=2", "mvel.max_errors=2",
"mvel.compiled_scripts_cache_size=100"
}) })
class MvelInvokeServiceTest extends AbstractControllerTest { class MvelInvokeServiceTest extends AbstractControllerTest {
@ -110,6 +119,89 @@ class MvelInvokeServiceTest extends AbstractControllerTest {
assertThatScriptIsBlocked(scriptId); assertThatScriptIsBlocked(scriptId);
} }
@Test
void givenScriptsWithSameBody_thenCompileAndCacheOnlyOnce() throws Exception {
String script = "return msg.temperature > 20;";
List<UUID> scriptsIds = new ArrayList<>();
for (int i = 0; i < 100; i++) {
UUID scriptId = evalScript(script);
scriptsIds.add(scriptId);
}
Map<UUID, String> scriptIdToHash = getFieldValue(invokeService, "scriptIdToHash");
Map<String, MvelScript> scriptMap = getFieldValue(invokeService, "scriptMap");
Cache<String, Serializable> compiledScriptsCache = getFieldValue(invokeService, "compiledScriptsCache");
String scriptHash = scriptIdToHash.get(scriptsIds.get(0));
assertThat(scriptsIds.stream().map(scriptIdToHash::get)).containsOnly(scriptHash);
assertThat(scriptMap).containsKey(scriptHash);
assertThat(compiledScriptsCache.getIfPresent(scriptHash)).isNotNull();
}
@Test
public void whenReleasingScript_thenCheckForScriptHashUsages() throws Exception {
String script = "return msg.temperature > 20;";
List<UUID> scriptsIds = new ArrayList<>();
for (int i = 0; i < 10; i++) {
UUID scriptId = evalScript(script);
scriptsIds.add(scriptId);
}
Map<UUID, String> scriptIdToHash = getFieldValue(invokeService, "scriptIdToHash");
Map<String, MvelScript> scriptMap = getFieldValue(invokeService, "scriptMap");
Cache<String, Serializable> compiledScriptsCache = getFieldValue(invokeService, "compiledScriptsCache");
String scriptHash = scriptIdToHash.get(scriptsIds.get(0));
for (int i = 0; i < 9; i++) {
UUID scriptId = scriptsIds.get(i);
assertThat(scriptIdToHash).containsKey(scriptId);
invokeService.release(scriptId);
assertThat(scriptIdToHash).doesNotContainKey(scriptId);
}
assertThat(scriptMap).containsKey(scriptHash);
assertThat(compiledScriptsCache.getIfPresent(scriptHash)).isNotNull();
invokeService.release(scriptsIds.get(9));
assertThat(scriptMap).doesNotContainKey(scriptHash);
assertThat(compiledScriptsCache.getIfPresent(scriptHash)).isNull();
}
@Test
public void whenCompiledScriptsCacheIsTooBig_thenRemoveRarelyUsedScripts() throws Exception {
Map<UUID, String> scriptIdToHash = getFieldValue(invokeService, "scriptIdToHash");
Cache<String, Serializable> compiledScriptsCache = getFieldValue(invokeService, "compiledScriptsCache");
List<UUID> scriptsIds = new ArrayList<>();
for (int i = 0; i < 110; i++) { // mvel.compiled_scripts_cache_size = 100
String script = "return msg.temperature > " + i;
UUID scriptId = evalScript(script);
scriptsIds.add(scriptId);
for (int j = 0; j < i; j++) {
invokeScript(scriptId, "{ \"temperature\": 12 }"); // so that scriptsIds is ordered by number of invocations
}
}
ConcurrentMap<String, Serializable> cache = compiledScriptsCache.asMap();
for (int i = 0; i < 10; i++) { // iterating rarely used scripts
UUID scriptId = scriptsIds.get(i);
String scriptHash = scriptIdToHash.get(scriptId);
assertThat(cache).doesNotContainKey(scriptHash);
}
for (int i = 10; i < 110; i++) {
UUID scriptId = scriptsIds.get(i);
String scriptHash = scriptIdToHash.get(scriptId);
assertThat(cache).containsKey(scriptHash);
}
UUID scriptRemovedFromCache = scriptsIds.get(0);
assertThat(compiledScriptsCache.getIfPresent(scriptIdToHash.get(scriptRemovedFromCache))).isNull();
invokeScript(scriptRemovedFromCache, "{ \"temperature\": 12 }");
assertThat(compiledScriptsCache.getIfPresent(scriptIdToHash.get(scriptRemovedFromCache))).isNotNull();
}
private void assertThatScriptIsBlocked(UUID scriptId) { private void assertThatScriptIsBlocked(UUID scriptId) {
assertThatThrownBy(() -> { assertThatThrownBy(() -> {
invokeScript(scriptId, "{}"); invokeScript(scriptId, "{}");
@ -125,4 +217,10 @@ class MvelInvokeServiceTest extends AbstractControllerTest {
return invokeService.invokeScript(TenantId.SYS_TENANT_ID, null, scriptId, msg, "{}", "POST_TELEMETRY_REQUEST").get().toString(); return invokeService.invokeScript(TenantId.SYS_TENANT_ID, null, scriptId, msg, "{}", "POST_TELEMETRY_REQUEST").get().toString();
} }
private <T> T getFieldValue(Object target, String fieldName) throws Exception {
Field field = target.getClass().getDeclaredField(fieldName);
field.setAccessible(true);
return (T) field.get(target);
}
} }

48
application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java

@ -15,27 +15,37 @@
*/ */
package org.thingsboard.server.service.state; package org.thingsboard.server.service.state;
import org.junit.Assert;
import org.junit.Before; import org.junit.Before;
import org.junit.Test; import org.junit.Test;
import org.junit.runner.RunWith; import org.junit.runner.RunWith;
import org.mockito.Mock; import org.mockito.Mock;
import org.mockito.Mockito; import org.mockito.Mockito;
import org.mockito.junit.MockitoJUnitRunner; import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.DeviceIdInfo;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.query.EntityData;
import org.thingsboard.server.common.data.query.EntityKeyType;
import org.thingsboard.server.common.data.query.TsValue;
import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import java.util.Map;
import java.util.UUID;
import static org.hamcrest.CoreMatchers.is; import static org.hamcrest.CoreMatchers.is;
import static org.hamcrest.MatcherAssert.assertThat; import static org.hamcrest.MatcherAssert.assertThat;
import static org.mockito.BDDMockito.willReturn; import static org.mockito.BDDMockito.willReturn;
import static org.mockito.Mockito.never; import static org.mockito.Mockito.never;
import static org.mockito.Mockito.spy; import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.times; import static org.mockito.Mockito.times;
import static org.thingsboard.server.service.state.DefaultDeviceStateService.INACTIVITY_TIMEOUT;
@RunWith(MockitoJUnitRunner.class) @RunWith(MockitoJUnitRunner.class)
public class DefaultDeviceStateServiceTest { public class DefaultDeviceStateServiceTest {
@ -83,4 +93,40 @@ public class DefaultDeviceStateServiceTest {
Mockito.verify(service, times(1)).fetchDeviceStateDataUsingEntityDataQuery(deviceId); Mockito.verify(service, times(1)).fetchDeviceStateDataUsingEntityDataQuery(deviceId);
} }
@Test
public void givenPersistToTelemetryAndDefaultInactivityTimeoutFetched_whenTransformingToDeviceStateData_thenTryGetInactivityFromAttribute() {
var defaultInactivityTimeoutInSec = 60L;
var latest =
Map.of(
EntityKeyType.TIME_SERIES, Map.of(INACTIVITY_TIMEOUT, new TsValue(0, Long.toString(defaultInactivityTimeoutInSec * 1000))),
EntityKeyType.SERVER_ATTRIBUTE, Map.of(INACTIVITY_TIMEOUT, new TsValue(0, Long.toString(5000L)))
);
process(latest, defaultInactivityTimeoutInSec);
}
@Test
public void givenPersistToTelemetryAndNoInactivityTimeoutFetchedFromTimeSeries_whenTransformingToDeviceStateData_thenTryGetInactivityFromAttribute() {
var defaultInactivityTimeoutInSec = 60L;
var latest =
Map.of(
EntityKeyType.SERVER_ATTRIBUTE, Map.of(INACTIVITY_TIMEOUT, new TsValue(0, Long.toString(5000L)))
);
process(latest, defaultInactivityTimeoutInSec);
}
private void process(Map<EntityKeyType, Map<String, TsValue>> latest, long defaultInactivityTimeoutInSec) {
service.setDefaultInactivityTimeoutInSec(defaultInactivityTimeoutInSec);
service.setDefaultInactivityTimeoutMs(defaultInactivityTimeoutInSec * 1000);
service.setPersistToTelemetry(true);
var deviceUuid = UUID.randomUUID();
var deviceId = new DeviceId(deviceUuid);
DeviceStateData deviceStateData = service.toDeviceStateData(new EntityData(deviceId, latest, Map.of()), new DeviceIdInfo(TenantId.SYS_TENANT_ID.getId(), UUID.randomUUID(), deviceUuid));
Assert.assertEquals(5000L, deviceStateData.getState().getInactivityTimeout());
}
} }

1
application/src/test/java/org/thingsboard/server/transport/TransportNoSqlTestSuite.java

@ -35,6 +35,7 @@ public class TransportNoSqlTestSuite {
public static CustomCassandraCQLUnit cassandraUnit = public static CustomCassandraCQLUnit cassandraUnit =
new CustomCassandraCQLUnit( new CustomCassandraCQLUnit(
Arrays.asList( Arrays.asList(
new ClassPathCQLDataSet("cassandra/schema-keyspace.cql", false, false),
new ClassPathCQLDataSet("cassandra/schema-ts.cql", false, false), new ClassPathCQLDataSet("cassandra/schema-ts.cql", false, false),
new ClassPathCQLDataSet("cassandra/schema-ts-latest.cql", false, false) new ClassPathCQLDataSet("cassandra/schema-ts-latest.cql", false, false)
), ),

5
application/src/test/resources/application-test.properties

@ -59,4 +59,7 @@ queue.rule-engine.queues[2].processing-strategy.max-pause-between-retries=0
usage.stats.report.enabled=false usage.stats.report.enabled=false
sql.audit_logs.partition_size=24 sql.audit_logs.partition_size=24
sql.ttl.audit_logs.ttl=2592000 sql.ttl.audit_logs.ttl=2592000
sql.edge_events.partition_size=168
sql.ttl.edge_events.edge_event_ttl=2592000

2
application/src/test/resources/logback-test.xml

@ -16,6 +16,8 @@
<logger name="org.cassandraunit" level="INFO"/> <logger name="org.cassandraunit" level="INFO"/>
<logger name="org.eclipse.leshan" level="INFO"/> <logger name="org.eclipse.leshan" level="INFO"/>
<!-- mute TelemetryEdgeSqlTest that causes a lot of randomly generated errors -->
<logger name="org.thingsboard.server.service.edge.rpc.EdgeGrpcSession" level="OFF"/>
<root level="WARN"> <root level="WARN">
<appender-ref ref="console"/> <appender-ref ref="console"/>

16
common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/AbstractCassandraCluster.java

@ -24,11 +24,13 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.core.env.Environment; import org.springframework.core.env.Environment;
import org.springframework.core.env.Profiles; import org.springframework.core.env.Profiles;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.dao.cassandra.guava.GuavaSession; import org.thingsboard.server.dao.cassandra.guava.GuavaSession;
import org.thingsboard.server.dao.cassandra.guava.GuavaSessionBuilder; import org.thingsboard.server.dao.cassandra.guava.GuavaSessionBuilder;
import org.thingsboard.server.dao.cassandra.guava.GuavaSessionUtils; import org.thingsboard.server.dao.cassandra.guava.GuavaSessionUtils;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import java.nio.file.Paths;
@Slf4j @Slf4j
public abstract class AbstractCassandraCluster { public abstract class AbstractCassandraCluster {
@ -40,6 +42,13 @@ public abstract class AbstractCassandraCluster {
@Value("${cassandra.local_datacenter:datacenter1}") @Value("${cassandra.local_datacenter:datacenter1}")
private String localDatacenter; private String localDatacenter;
@Value("${cassandra.cloud.secure_connect_bundle_path:}")
private String cloudSecureConnectBundlePath;
@Value("${cassandra.cloud.client_id:}")
private String cloudClientId;
@Value("${cassandra.cloud.client_secret:}")
private String cloudClientSecret;
@Autowired @Autowired
private CassandraDriverOptions driverOptions; private CassandraDriverOptions driverOptions;
@ -86,7 +95,14 @@ public abstract class AbstractCassandraCluster {
this.sessionBuilder.withKeyspace(this.keyspaceName); this.sessionBuilder.withKeyspace(this.keyspaceName);
} }
this.sessionBuilder.withLocalDatacenter(localDatacenter); this.sessionBuilder.withLocalDatacenter(localDatacenter);
if (StringUtils.isNotBlank(cloudSecureConnectBundlePath)) {
this.sessionBuilder.withCloudSecureConnectBundle(Paths.get(cloudSecureConnectBundlePath));
this.sessionBuilder.withAuthCredentials(cloudClientId, cloudClientSecret);
}
session = sessionBuilder.build(); session = sessionBuilder.build();
if (this.metrics && this.jmx) { if (this.metrics && this.jmx) {
MetricRegistry registry = MetricRegistry registry =
session.getMetrics().orElseThrow( session.getMetrics().orElseThrow(

22
common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaDriverContext.java

@ -40,26 +40,8 @@ import java.util.function.Predicate;
*/ */
public class GuavaDriverContext extends DefaultDriverContext { public class GuavaDriverContext extends DefaultDriverContext {
public GuavaDriverContext( public GuavaDriverContext(DriverConfigLoader configLoader, ProgrammaticArguments programmaticArguments) {
DriverConfigLoader configLoader, super(configLoader, programmaticArguments);
List<TypeCodec<?>> typeCodecs,
NodeStateListener nodeStateListener,
SchemaChangeListener schemaChangeListener,
RequestTracker requestTracker,
Map<String, String> localDatacenters,
Map<String, Predicate<Node>> nodeFilters,
ClassLoader classLoader) {
super(
configLoader,
ProgrammaticArguments.builder()
.addTypeCodecs(typeCodecs.toArray(new TypeCodec<?>[0]))
.withNodeStateListener(nodeStateListener)
.withSchemaChangeListener(schemaChangeListener)
.withRequestTracker(requestTracker)
.withLocalDatacenters(localDatacenters)
.withNodeFilters(nodeFilters)
.withClassLoader(classLoader)
.build());
} }
@Override @Override

14
common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaSessionBuilder.java

@ -25,18 +25,8 @@ import edu.umd.cs.findbugs.annotations.NonNull;
public class GuavaSessionBuilder extends SessionBuilder<GuavaSessionBuilder, GuavaSession> { public class GuavaSessionBuilder extends SessionBuilder<GuavaSessionBuilder, GuavaSession> {
@Override @Override
protected DriverContext buildContext( protected DriverContext buildContext(DriverConfigLoader configLoader, ProgrammaticArguments programmaticArguments) {
DriverConfigLoader configLoader, return new GuavaDriverContext(configLoader, programmaticArguments);
ProgrammaticArguments programmaticArguments) {
return new GuavaDriverContext(
configLoader,
programmaticArguments.getTypeCodecs(),
programmaticArguments.getNodeStateListener(),
programmaticArguments.getSchemaChangeListener(),
programmaticArguments.getRequestTracker(),
programmaticArguments.getLocalDatacenters(),
programmaticArguments.getNodeFilters(),
programmaticArguments.getClassLoader());
} }
@Override @Override

27
common/dao-api/src/main/java/org/thingsboard/server/dao/util/NoSqlAnyDaoNonCloud.java

@ -0,0 +1,27 @@
/**
* Copyright © 2016-2022 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.util;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
@Retention(RetentionPolicy.RUNTIME)
@ConditionalOnExpression("('${database.ts.type}'=='cassandra' || '${database.ts_latest.type}'=='cassandra') " +
"&& ('${cassandra.cloud.secure_connect_bundle_path}' == null || '${cassandra.cloud.secure_connect_bundle_path}'.isBlank() )")
public @interface NoSqlAnyDaoNonCloud {
}

2
common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java

@ -40,6 +40,8 @@ public class DataConstants {
public static final String EXPIRATION_TIME = "expirationTime"; public static final String EXPIRATION_TIME = "expirationTime";
public static final String ADDITIONAL_INFO = "additionalInfo"; public static final String ADDITIONAL_INFO = "additionalInfo";
public static final String RETRIES = "retries"; public static final String RETRIES = "retries";
public static final String EDGE_ID = "edgeId";
public static final String DEVICE_ID = "deviceId";
public static final String COAP_TRANSPORT_NAME = "COAP"; public static final String COAP_TRANSPORT_NAME = "COAP";
public static final String LWM2M_TRANSPORT_NAME = "LWM2M"; public static final String LWM2M_TRANSPORT_NAME = "LWM2M";
public static final String MQTT_TRANSPORT_NAME = "MQTT"; public static final String MQTT_TRANSPORT_NAME = "MQTT";

26
common/data/src/main/java/org/thingsboard/server/common/data/util/TbPair.java

@ -0,0 +1,26 @@
/**
* Copyright © 2016-2022 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.util;
import lombok.AllArgsConstructor;
import lombok.Data;
@Data
@AllArgsConstructor
public class TbPair<S, T> {
private S first;
private T second;
}

5
common/edge-api/src/main/proto/edge.proto

@ -430,6 +430,11 @@ message DeviceRpcCallMsg {
bool oneway = 7; bool oneway = 7;
RpcRequestMsg requestMsg = 8; RpcRequestMsg requestMsg = 8;
RpcResponseMsg responseMsg = 9; RpcResponseMsg responseMsg = 9;
optional bool persisted = 10;
optional int32 retries = 11;
optional string additionalInfo = 12;
optional string serviceId = 13;
optional string sessionId = 14;
} }
message RpcRequestMsg { message RpcRequestMsg {

4
common/script/script-api/pom.xml

@ -56,6 +56,10 @@
<groupId>com.google.code.gson</groupId> <groupId>com.google.code.gson</groupId>
<artifactId>gson</artifactId> <artifactId>gson</artifactId>
</dependency> </dependency>
<dependency>
<groupId>com.github.ben-manes.caffeine</groupId>
<artifactId>caffeine</artifactId>
</dependency>
<dependency> <dependency>
<groupId>org.slf4j</groupId> <groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId> <artifactId>slf4j-api</artifactId>

78
common/script/script-api/src/main/java/org/thingsboard/script/api/mvel/DefaultMvelInvokeService.java

@ -15,6 +15,10 @@
*/ */
package org.thingsboard.script.api.mvel; package org.thingsboard.script.api.mvel;
import com.github.benmanes.caffeine.cache.Cache;
import com.github.benmanes.caffeine.cache.Caffeine;
import com.google.common.hash.Hasher;
import com.google.common.hash.Hashing;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.ListeningExecutorService; import com.google.common.util.concurrent.ListeningExecutorService;
import com.google.common.util.concurrent.MoreExecutors; import com.google.common.util.concurrent.MoreExecutors;
@ -42,12 +46,15 @@ import org.thingsboard.server.common.stats.TbApiUsageStateClient;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import java.io.Serializable; import java.io.Serializable;
import java.nio.charset.StandardCharsets;
import java.util.Collections; import java.util.Collections;
import java.util.Map; import java.util.Map;
import java.util.Optional; import java.util.Optional;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Executor; import java.util.concurrent.Executor;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import java.util.regex.Pattern; import java.util.regex.Pattern;
@Slf4j @Slf4j
@ -55,7 +62,10 @@ import java.util.regex.Pattern;
@Service @Service
public class DefaultMvelInvokeService extends AbstractScriptInvokeService implements MvelInvokeService { public class DefaultMvelInvokeService extends AbstractScriptInvokeService implements MvelInvokeService {
protected Map<UUID, MvelScript> scriptMap = new ConcurrentHashMap<>(); protected final Map<UUID, String> scriptIdToHash = new ConcurrentHashMap<>();
protected final Map<String, MvelScript> scriptMap = new ConcurrentHashMap<>();
protected Cache<String, Serializable> compiledScriptsCache;
private SandboxedParserConfiguration parserConfig; private SandboxedParserConfiguration parserConfig;
private static final Pattern NEW_KEYWORD_PATTERN = Pattern.compile("new\\s"); private static final Pattern NEW_KEYWORD_PATTERN = Pattern.compile("new\\s");
@ -92,8 +102,13 @@ public class DefaultMvelInvokeService extends AbstractScriptInvokeService implem
@Value("${mvel.max_memory_limit_mb:8}") @Value("${mvel.max_memory_limit_mb:8}")
private long maxMemoryLimitMb; private long maxMemoryLimitMb;
@Value("${mvel.compiled_scripts_cache_size:1000}")
private int compiledScriptsCacheSize;
private ListeningExecutorService executor; private ListeningExecutorService executor;
private final Lock lock = new ReentrantLock();
protected DefaultMvelInvokeService(Optional<TbApiUsageStateClient> apiUsageStateClient, Optional<TbApiUsageReportClient> apiUsageReportClient) { protected DefaultMvelInvokeService(Optional<TbApiUsageStateClient> apiUsageStateClient, Optional<TbApiUsageReportClient> apiUsageReportClient) {
super(apiUsageStateClient, apiUsageReportClient); super(apiUsageStateClient, apiUsageReportClient);
} }
@ -115,11 +130,14 @@ public class DefaultMvelInvokeService extends AbstractScriptInvokeService implem
executor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(threadPoolSize, "mvel-executor")); executor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(threadPoolSize, "mvel-executor"));
try { try {
// Special command to warm up MVEL engine // Special command to warm up MVEL engine
Serializable script = MVEL.compileExpression("var warmUp = {}; warmUp", new SandboxedParserContext(parserConfig)); Serializable script = compileScript("var warmUp = {}; warmUp");
MVEL.executeTbExpression(script, new ExecutionContext(parserConfig), Collections.emptyMap()); MVEL.executeTbExpression(script, new ExecutionContext(parserConfig), Collections.emptyMap());
} catch (Exception e) { } catch (Exception e) {
// do nothing // do nothing
} }
compiledScriptsCache = Caffeine.newBuilder()
.maximumSize(compiledScriptsCacheSize)
.build();
} }
@PreDestroy @PreDestroy
@ -141,16 +159,26 @@ public class DefaultMvelInvokeService extends AbstractScriptInvokeService implem
@Override @Override
protected boolean isScriptPresent(UUID scriptId) { protected boolean isScriptPresent(UUID scriptId) {
return scriptMap.containsKey(scriptId); return scriptIdToHash.containsKey(scriptId);
} }
@Override @Override
protected ListenableFuture<UUID> doEvalScript(TenantId tenantId, ScriptType scriptType, String scriptBody, UUID scriptId, String[] argNames) { protected ListenableFuture<UUID> doEvalScript(TenantId tenantId, ScriptType scriptType, String scriptBody, UUID scriptId, String[] argNames) {
return executor.submit(() -> { return executor.submit(() -> {
try { try {
Serializable compiledScript = MVEL.compileExpression(scriptBody, new SandboxedParserContext(parserConfig)); String scriptHash = hash(scriptBody, argNames);
MvelScript script = new MvelScript(compiledScript, scriptBody, argNames); compiledScriptsCache.get(scriptHash, k -> {
scriptMap.put(scriptId, script); return compileScript(scriptBody);
});
lock.lock();
try {
scriptIdToHash.put(scriptId, scriptHash);
scriptMap.computeIfAbsent(scriptHash, k -> {
return new MvelScript(scriptBody, argNames);
});
} finally {
lock.unlock();
}
return scriptId; return scriptId;
} catch (Exception e) { } catch (Exception e) {
throw new TbScriptException(scriptId, TbScriptException.ErrorCode.COMPILATION, scriptBody, e); throw new TbScriptException(scriptId, TbScriptException.ErrorCode.COMPILATION, scriptBody, e);
@ -162,12 +190,16 @@ public class DefaultMvelInvokeService extends AbstractScriptInvokeService implem
protected MvelScriptExecutionTask doInvokeFunction(UUID scriptId, Object[] args) { protected MvelScriptExecutionTask doInvokeFunction(UUID scriptId, Object[] args) {
ExecutionContext executionContext = new ExecutionContext(this.parserConfig, maxMemoryLimitMb * 1024 * 1024); ExecutionContext executionContext = new ExecutionContext(this.parserConfig, maxMemoryLimitMb * 1024 * 1024);
return new MvelScriptExecutionTask(executionContext, executor.submit(() -> { return new MvelScriptExecutionTask(executionContext, executor.submit(() -> {
MvelScript script = scriptMap.get(scriptId); String scriptHash = scriptIdToHash.get(scriptId);
if (script == null) { if (scriptHash == null) {
throw new TbScriptException(scriptId, TbScriptException.ErrorCode.OTHER, null, new RuntimeException("Script not found!")); throw new TbScriptException(scriptId, TbScriptException.ErrorCode.OTHER, null, new RuntimeException("Script not found!"));
} }
MvelScript script = scriptMap.get(scriptHash);
Serializable compiledScript = compiledScriptsCache.get(scriptHash, k -> {
return compileScript(script.getScriptBody());
});
try { try {
return MVEL.executeTbExpression(script.getCompiledScript(), executionContext, script.createVars(args)); return MVEL.executeTbExpression(compiledScript, executionContext, script.createVars(args));
} catch (ScriptMemoryOverflowException e) { } catch (ScriptMemoryOverflowException e) {
throw new TbScriptException(scriptId, TbScriptException.ErrorCode.OTHER, script.getScriptBody(), new RuntimeException("Script memory overflow!")); throw new TbScriptException(scriptId, TbScriptException.ErrorCode.OTHER, script.getScriptBody(), new RuntimeException("Script memory overflow!"));
} catch (Exception e) { } catch (Exception e) {
@ -178,6 +210,32 @@ public class DefaultMvelInvokeService extends AbstractScriptInvokeService implem
@Override @Override
protected void doRelease(UUID scriptId) throws Exception { protected void doRelease(UUID scriptId) throws Exception {
scriptMap.remove(scriptId); String scriptHash = scriptIdToHash.remove(scriptId);
if (scriptHash != null) {
lock.lock();
try {
if (!scriptIdToHash.containsValue(scriptHash)) {
scriptMap.remove(scriptHash);
compiledScriptsCache.invalidate(scriptHash);
}
} finally {
lock.unlock();
}
}
} }
private Serializable compileScript(String scriptBody) {
return MVEL.compileExpression(scriptBody, new SandboxedParserContext(parserConfig));
}
@SuppressWarnings("UnstableApiUsage")
protected String hash(String scriptBody, String[] argNames) {
Hasher hasher = Hashing.murmur3_128().newHasher();
hasher.putUnencodedChars(scriptBody);
for (String argName : argNames) {
hasher.putString(argName, StandardCharsets.UTF_8);
}
return hasher.hash().toString();
}
} }

1
common/script/script-api/src/main/java/org/thingsboard/script/api/mvel/MvelScript.java

@ -24,7 +24,6 @@ import java.util.Map;
@Data @Data
public class MvelScript { public class MvelScript {
private final Serializable compiledScript;
private final String scriptBody; private final String scriptBody;
private final String[] argNames; private final String[] argNames;

3
dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java

@ -26,6 +26,7 @@ import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.Dao; import org.thingsboard.server.dao.Dao;
import org.thingsboard.server.dao.ExportableEntityDao; import org.thingsboard.server.dao.ExportableEntityDao;
import org.thingsboard.server.dao.TenantEntityDao; import org.thingsboard.server.dao.TenantEntityDao;
import org.thingsboard.server.common.data.util.TbPair;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.Optional;
@ -222,4 +223,6 @@ public interface AssetDao extends Dao<Asset>, TenantEntityDao, ExportableEntityD
* @return the list of asset objects * @return the list of asset objects
*/ */
PageData<Asset> findAssetsByTenantIdAndEdgeIdAndType(UUID tenantId, UUID edgeId, String type, PageLink pageLink); PageData<Asset> findAssetsByTenantIdAndEdgeIdAndType(UUID tenantId, UUID edgeId, String type, PageLink pageLink);
PageData<TbPair<UUID, String>> getAllAssetTypes(PageLink pageLink);
} }

9
dao/src/main/java/org/thingsboard/server/dao/edge/BaseEdgeEventService.java

@ -16,8 +16,8 @@
package org.thingsboard.server.dao.edge; package org.thingsboard.server.dao.edge;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.AllArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EdgeId;
@ -28,13 +28,12 @@ import org.thingsboard.server.dao.service.DataValidator;
@Service @Service
@Slf4j @Slf4j
@AllArgsConstructor
public class BaseEdgeEventService implements EdgeEventService { public class BaseEdgeEventService implements EdgeEventService {
@Autowired private final EdgeEventDao edgeEventDao;
private EdgeEventDao edgeEventDao;
@Autowired private final DataValidator<EdgeEvent> edgeEventValidator;
private DataValidator<EdgeEvent> edgeEventValidator;
@Override @Override
public ListenableFuture<Void> saveAsync(EdgeEvent edgeEvent) { public ListenableFuture<Void> saveAsync(EdgeEvent edgeEvent) {

2
dao/src/main/java/org/thingsboard/server/dao/edge/EdgeEventDao.java

@ -54,4 +54,6 @@ public interface EdgeEventDao extends Dao<EdgeEvent> {
*/ */
void cleanupEvents(long ttl); void cleanupEvents(long ttl);
void migrateEdgeEvents();
} }

2
dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java

@ -731,6 +731,8 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
ConstraintViolationException e = extractConstraintViolationException(t).orElse(null); ConstraintViolationException e = extractConstraintViolationException(t).orElse(null);
if (e != null && e.getConstraintName() != null && e.getConstraintName().equalsIgnoreCase("fk_default_rule_chain_device_profile")) { if (e != null && e.getConstraintName() != null && e.getConstraintName().equalsIgnoreCase("fk_default_rule_chain_device_profile")) {
throw new DataValidationException("The rule chain referenced by the device profiles cannot be deleted!"); throw new DataValidationException("The rule chain referenced by the device profiles cannot be deleted!");
} else if (e != null && e.getConstraintName() != null && e.getConstraintName().equalsIgnoreCase("fk_default_rule_chain_asset_profile")) {
throw new DataValidationException("The rule chain referenced by the asset profiles cannot be deleted!");
} else { } else {
throw t; throw t;
} }

18
dao/src/main/java/org/thingsboard/server/dao/sql/asset/AssetRepository.java

@ -23,6 +23,7 @@ import org.springframework.data.repository.query.Param;
import org.thingsboard.server.dao.ExportableEntityRepository; import org.thingsboard.server.dao.ExportableEntityRepository;
import org.thingsboard.server.dao.model.sql.AssetEntity; import org.thingsboard.server.dao.model.sql.AssetEntity;
import org.thingsboard.server.dao.model.sql.AssetInfoEntity; import org.thingsboard.server.dao.model.sql.AssetInfoEntity;
import org.thingsboard.server.common.data.util.TbPair;
import java.util.List; import java.util.List;
import java.util.UUID; import java.util.UUID;
@ -70,9 +71,9 @@ public interface AssetRepository extends JpaRepository<AssetEntity, UUID>, Expor
"AND a.assetProfileId = :profileId " + "AND a.assetProfileId = :profileId " +
"AND LOWER(a.searchText) LIKE LOWER(CONCAT('%', :searchText, '%'))") "AND LOWER(a.searchText) LIKE LOWER(CONCAT('%', :searchText, '%'))")
Page<AssetEntity> findByTenantIdAndProfileId(@Param("tenantId") UUID tenantId, Page<AssetEntity> findByTenantIdAndProfileId(@Param("tenantId") UUID tenantId,
@Param("profileId") UUID profileId, @Param("profileId") UUID profileId,
@Param("searchText") String searchText, @Param("searchText") String searchText,
Pageable pageable); Pageable pageable);
@Query("SELECT new org.thingsboard.server.dao.model.sql.AssetInfoEntity(a, c.title, c.additionalInfo, p.name) " + @Query("SELECT new org.thingsboard.server.dao.model.sql.AssetInfoEntity(a, c.title, c.additionalInfo, p.name) " +
"FROM AssetEntity a " + "FROM AssetEntity a " +
@ -186,14 +187,17 @@ public interface AssetRepository extends JpaRepository<AssetEntity, UUID>, Expor
"AND a.type = :type " + "AND a.type = :type " +
"AND LOWER(a.searchText) LIKE LOWER(CONCAT('%', :searchText, '%'))") "AND LOWER(a.searchText) LIKE LOWER(CONCAT('%', :searchText, '%'))")
Page<AssetEntity> findByTenantIdAndEdgeIdAndType(@Param("tenantId") UUID tenantId, Page<AssetEntity> findByTenantIdAndEdgeIdAndType(@Param("tenantId") UUID tenantId,
@Param("edgeId") UUID edgeId, @Param("edgeId") UUID edgeId,
@Param("type") String type, @Param("type") String type,
@Param("searchText") String searchText, @Param("searchText") String searchText,
Pageable pageable); Pageable pageable);
Long countByTenantIdAndTypeIsNot(UUID tenantId, String type); Long countByTenantIdAndTypeIsNot(UUID tenantId, String type);
@Query("SELECT externalId FROM AssetEntity WHERE id = :id") @Query("SELECT externalId FROM AssetEntity WHERE id = :id")
UUID getExternalIdById(@Param("id") UUID id); UUID getExternalIdById(@Param("id") UUID id);
@Query(value = "SELECT DISTINCT new org.thingsboard.server.common.data.util.TbPair(a.tenantId , a.type) FROM AssetEntity a")
Page<TbPair<UUID, String>> getAllAssetTypes(Pageable pageable);
} }

9
dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java

@ -28,14 +28,17 @@ import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.page.SortOrder;
import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.asset.AssetDao; import org.thingsboard.server.dao.asset.AssetDao;
import org.thingsboard.server.dao.model.sql.AssetEntity; import org.thingsboard.server.dao.model.sql.AssetEntity;
import org.thingsboard.server.dao.model.sql.AssetInfoEntity; import org.thingsboard.server.dao.model.sql.AssetInfoEntity;
import org.thingsboard.server.common.data.util.TbPair;
import org.thingsboard.server.dao.sql.JpaAbstractSearchTextDao; import org.thingsboard.server.dao.sql.JpaAbstractSearchTextDao;
import org.thingsboard.server.dao.util.SqlDao; import org.thingsboard.server.dao.util.SqlDao;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections; import java.util.Collections;
import java.util.List; import java.util.List;
import java.util.Objects; import java.util.Objects;
@ -243,6 +246,12 @@ public class JpaAssetDao extends JpaAbstractSearchTextDao<AssetEntity, Asset> im
DaoUtil.toPageable(pageLink))); DaoUtil.toPageable(pageLink)));
} }
public PageData<TbPair<UUID, String>> getAllAssetTypes(PageLink pageLink) {
log.debug("Try to find all asset types and pageLink [{}]", pageLink);
return DaoUtil.pageToPageData(assetRepository.getAllAssetTypes(
DaoUtil.toPageable(pageLink, Arrays.asList(new SortOrder("tenantId"), new SortOrder("type")))));
}
@Override @Override
public Long countByTenantId(TenantId tenantId) { public Long countByTenantId(TenantId tenantId) {
return assetRepository.countByTenantIdAndTypeIsNot(tenantId.getId(), TB_SERVICE_QUEUE); return assetRepository.countByTenantIdAndTypeIsNot(tenantId.getId(), TB_SERVICE_QUEUE);

83
dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java

@ -17,10 +17,11 @@ package org.thingsboard.server.dao.sql.edge;
import com.datastax.oss.driver.api.core.uuid.Uuids; import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.edge.EdgeEvent; import org.thingsboard.server.common.data.edge.EdgeEvent;
@ -31,19 +32,17 @@ import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.edge.EdgeEventDao; import org.thingsboard.server.dao.edge.EdgeEventDao;
import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.model.sql.EdgeEventEntity; import org.thingsboard.server.dao.model.sql.EdgeEventEntity;
import org.thingsboard.server.dao.sql.JpaAbstractSearchTextDao; import org.thingsboard.server.dao.sql.JpaAbstractSearchTextDao;
import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper;
import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository;
import org.thingsboard.server.dao.util.SqlDao; import org.thingsboard.server.dao.util.SqlDao;
import javax.annotation.PostConstruct; import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy; import javax.annotation.PreDestroy;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.Comparator; import java.util.Comparator;
import java.util.Objects; import java.util.Objects;
import java.util.UUID; import java.util.UUID;
@ -52,18 +51,25 @@ import java.util.function.Function;
import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID; import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID;
@Slf4j
@Component @Component
@SqlDao @SqlDao
@RequiredArgsConstructor
@Slf4j
public class JpaBaseEdgeEventDao extends JpaAbstractSearchTextDao<EdgeEventEntity, EdgeEvent> implements EdgeEventDao { public class JpaBaseEdgeEventDao extends JpaAbstractSearchTextDao<EdgeEventEntity, EdgeEvent> implements EdgeEventDao {
private final UUID systemTenantId = NULL_UUID; private final UUID systemTenantId = NULL_UUID;
@Autowired private final ScheduledLogExecutorComponent logExecutor;
ScheduledLogExecutorComponent logExecutor;
private final StatsFactory statsFactory;
private final EdgeEventRepository edgeEventRepository;
private final EdgeEventInsertRepository edgeEventInsertRepository;
@Autowired private final SqlPartitioningRepository partitioningRepository;
private StatsFactory statsFactory;
private final JdbcTemplate jdbcTemplate;
@Value("${sql.edge_events.batch_size:1000}") @Value("${sql.edge_events.batch_size:1000}")
private int batchSize; private int batchSize;
@ -74,13 +80,15 @@ public class JpaBaseEdgeEventDao extends JpaAbstractSearchTextDao<EdgeEventEntit
@Value("${sql.edge_events.stats_print_interval_ms:10000}") @Value("${sql.edge_events.stats_print_interval_ms:10000}")
private long statsPrintIntervalMs; private long statsPrintIntervalMs;
private TbSqlBlockingQueueWrapper<EdgeEventEntity> queue; @Value("${sql.edge_events.partitions_size:168}")
private int partitionSizeInHours;
@Autowired @Value("${sql.ttl.edge_events.edge_events_ttl:2628000}")
private EdgeEventRepository edgeEventRepository; private long edge_events_ttl;
@Autowired private static final String TABLE_NAME = ModelConstants.EDGE_EVENT_COLUMN_FAMILY_NAME;
private EdgeEventInsertRepository edgeEventInsertRepository;
private TbSqlBlockingQueueWrapper<EdgeEventEntity> queue;
@Override @Override
protected Class<EdgeEventEntity> getEntityClass() { protected Class<EdgeEventEntity> getEntityClass() {
@ -140,6 +148,7 @@ public class JpaBaseEdgeEventDao extends JpaAbstractSearchTextDao<EdgeEventEntit
if (StringUtils.isEmpty(edgeEvent.getUid())) { if (StringUtils.isEmpty(edgeEvent.getUid())) {
edgeEvent.setUid(edgeEvent.getId().toString()); edgeEvent.setUid(edgeEvent.getId().toString());
} }
partitioningRepository.createPartitionIfNotExists(TABLE_NAME, edgeEvent.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours));
return save(new EdgeEventEntity(edgeEvent)); return save(new EdgeEventEntity(edgeEvent));
} }
@ -189,20 +198,36 @@ public class JpaBaseEdgeEventDao extends JpaAbstractSearchTextDao<EdgeEventEntit
@Override @Override
public void cleanupEvents(long ttl) { public void cleanupEvents(long ttl) {
log.info("Going to cleanup old edge events using ttl: {}s", ttl); partitioningRepository.dropPartitionsBefore(TABLE_NAME, ttl, TimeUnit.HOURS.toMillis(partitionSizeInHours));
try (Connection connection = dataSource.getConnection(); }
PreparedStatement stmt = connection.prepareStatement("call cleanup_edge_events_by_ttl(?,?)")) {
stmt.setLong(1, ttl); @Override
stmt.setLong(2, 0); public void migrateEdgeEvents() {
stmt.setQueryTimeout((int) TimeUnit.HOURS.toSeconds(1)); long startTime = edge_events_ttl > 0 ? System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(edge_events_ttl) : 1629158400000L;
stmt.execute();
printWarnings(stmt); long currentTime = System.currentTimeMillis();
try (ResultSet resultSet = stmt.getResultSet()) { var partitionStepInMs = TimeUnit.HOURS.toMillis(partitionSizeInHours);
resultSet.next(); long numberOfPartitions = (currentTime - startTime) / partitionStepInMs;
log.info("Total edge events removed by TTL: [{}]", resultSet.getLong(1));
} if (numberOfPartitions > 1000) {
} catch (SQLException e) { String error = "Please adjust your edge event partitioning configuration. Configuration with partition size " +
log.error("SQLException occurred during edge events TTL task execution ", e); "of " + partitionSizeInHours + " hours and corresponding TTL will use " + numberOfPartitions + " " +
"(> 1000) partitions which is not recommended!";
log.error(error);
throw new RuntimeException(error);
} }
while (startTime < currentTime) {
var endTime = startTime + partitionStepInMs;
log.info("Migrating edge event for time period: {} - {}", startTime, endTime);
callMigrationFunction(startTime, endTime, partitionStepInMs);
startTime = endTime;
}
log.info("Event edge migration finished");
jdbcTemplate.execute("DROP TABLE IF EXISTS old_edge_event");
}
private void callMigrationFunction(long startTime, long endTime, long partitionSIzeInMs) {
jdbcTemplate.update("CALL migrate_edge_event(?, ?, ?)", startTime, endTime, partitionSIzeInMs);
} }
} }

17
dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultAlarmQueryRepository.java

@ -140,6 +140,7 @@ public class DefaultAlarmQueryRepository implements AlarmQueryRepository {
selectPart.append(" a.originator_id as entity_id "); selectPart.append(" a.originator_id as entity_id ");
} }
EntityDataSortOrder sortOrder = pageLink.getSortOrder(); EntityDataSortOrder sortOrder = pageLink.getSortOrder();
String textSearchQuery = buildTextSearchQuery(ctx, query.getAlarmFields(), pageLink.getTextSearch());
if (sortOrder != null && sortOrder.getKey().getType().equals(EntityKeyType.ALARM_FIELD)) { if (sortOrder != null && sortOrder.getKey().getType().equals(EntityKeyType.ALARM_FIELD)) {
String sortOrderKey = sortOrder.getKey().getKey(); String sortOrderKey = sortOrder.getKey().getKey();
sortPart.append(alarmFieldColumnMap.getOrDefault(sortOrderKey, sortOrderKey)) sortPart.append(alarmFieldColumnMap.getOrDefault(sortOrderKey, sortOrderKey))
@ -166,7 +167,11 @@ public class DefaultAlarmQueryRepository implements AlarmQueryRepository {
} }
joinPart.append(" as e(id, priority)) e "); joinPart.append(" as e(id, priority)) e ");
if (pageLink.isSearchPropagatedAlarms()) { if (pageLink.isSearchPropagatedAlarms()) {
joinPart.append("on ea.entity_id = e.id"); if (textSearchQuery.isEmpty()) {
joinPart.append("on ea.entity_id = e.id");
} else {
joinPart.append("on a.entity_id = e.id");
}
} else { } else {
joinPart.append("on a.originator_id = e.id"); joinPart.append("on a.originator_id = e.id");
} }
@ -230,13 +235,11 @@ public class DefaultAlarmQueryRepository implements AlarmQueryRepository {
} }
} }
String textSearchQuery = buildTextSearchQuery(ctx, query.getAlarmFields(), pageLink.getTextSearch()); String mainQuery = String.format("%s%s", selectPart, fromPart);
String mainQuery; if (textSearchQuery.isEmpty()) {
if (!textSearchQuery.isEmpty()) { mainQuery = String.format("%s%s%s", mainQuery, joinPart, wherePart);
mainQuery = selectPart.toString() + fromPart.toString() + wherePart.toString();
mainQuery = String.format("select * from (%s) a %s WHERE %s", mainQuery, joinPart, textSearchQuery);
} else { } else {
mainQuery = selectPart.toString() + fromPart.toString() + joinPart.toString() + wherePart.toString(); mainQuery = String.format("select * from (%s%s) a %s WHERE %s", mainQuery, wherePart, joinPart, textSearchQuery);
} }
String countQuery = String.format("select count(*) from (%s) result", mainQuery); String countQuery = String.format("select count(*) from (%s) result", mainQuery);
long queryTs = System.currentTimeMillis(); long queryTs = System.currentTimeMillis();

1
dao/src/main/java/org/thingsboard/server/dao/sql/tenant/TenantRepository.java

@ -53,5 +53,4 @@ public interface TenantRepository extends JpaRepository<TenantEntity, UUID> {
@Query("SELECT t.id FROM TenantEntity t where t.tenantProfileId = :tenantProfileId") @Query("SELECT t.id FROM TenantEntity t where t.tenantProfileId = :tenantProfileId")
List<UUID> findTenantIdsByTenantProfileId(@Param("tenantProfileId") UUID tenantProfileId); List<UUID> findTenantIdsByTenantProfileId(@Param("tenantProfileId") UUID tenantProfileId);
} }

5
dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java

@ -20,7 +20,7 @@ import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.domain.PageRequest; import org.springframework.data.domain.PageRequest;
import org.springframework.data.domain.Sort; import org.springframework.data.domain.Sort.Direction;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.kv.Aggregation;
@ -143,8 +143,7 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq
keyId, keyId,
query.getStartTs(), query.getStartTs(),
query.getEndTs(), query.getEndTs(),
PageRequest.of(0, query.getLimit(), PageRequest.ofSize(query.getLimit()).withSort(Direction.fromString(query.getOrder()), "ts"));
Sort.by(new Sort.Order(Sort.Direction.fromString(query.getOrder()), "ts").nullsNative())));
tsKvEntities.forEach(tsKvEntity -> tsKvEntity.setStrKey(query.getKey())); tsKvEntities.forEach(tsKvEntity -> tsKvEntity.setStrKey(query.getKey()));
List<TsKvEntry> tsKvEntries = DaoUtil.convertDataList(tsKvEntities); List<TsKvEntry> tsKvEntries = DaoUtil.convertDataList(tsKvEntities);
long lastTs = tsKvEntries.stream().map(TsKvEntry::getTs).max(Long::compare).orElse(query.getStartTs()); long lastTs = tsKvEntries.stream().map(TsKvEntry::getTs).max(Long::compare).orElse(query.getStartTs());

4
dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java

@ -175,9 +175,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
keyId, keyId,
query.getStartTs(), query.getStartTs(),
query.getEndTs(), query.getEndTs(),
PageRequest.of(0, query.getLimit(), PageRequest.ofSize(query.getLimit()).withSort(Sort.Direction.fromString(query.getOrder()), "ts"));
Sort.by(new Sort.Order(Sort.Direction.fromString(query.getOrder()), "ts").nullsNative())));
;
timescaleTsKvEntities.forEach(tsKvEntity -> tsKvEntity.setStrKey(strKey)); timescaleTsKvEntities.forEach(tsKvEntity -> tsKvEntity.setStrKey(strKey));
var tsKvEntries = DaoUtil.convertDataList(timescaleTsKvEntities); var tsKvEntries = DaoUtil.convertDataList(timescaleTsKvEntities);
long lastTs = tsKvEntries.stream().map(TsKvEntry::getTs).max(Long::compare).orElse(query.getStartTs()); long lastTs = tsKvEntries.stream().map(TsKvEntry::getTs).max(Long::compare).orElse(query.getStartTs());

15
dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TsKvTimescaleRepository.java

@ -31,14 +31,13 @@ import java.util.UUID;
@TimescaleDBTsOrTsLatestDao @TimescaleDBTsOrTsLatestDao
public interface TsKvTimescaleRepository extends JpaRepository<TimescaleTsKvEntity, TimescaleTsKvCompositeKey> { public interface TsKvTimescaleRepository extends JpaRepository<TimescaleTsKvEntity, TimescaleTsKvCompositeKey> {
@Query("SELECT tskv FROM TimescaleTsKvEntity tskv WHERE tskv.entityId = :entityId " + @Query(value = "SELECT * FROM ts_kv WHERE entity_id = :entityId " +
"AND tskv.key = :entityKey " + "AND key = :entityKey AND ts >= :startTs AND ts < :endTs", nativeQuery = true)
"AND tskv.ts >= :startTs AND tskv.ts < :endTs") List<TimescaleTsKvEntity> findAllWithLimit(@Param("entityId") UUID entityId,
List<TimescaleTsKvEntity> findAllWithLimit( @Param("entityKey") int key,
@Param("entityId") UUID entityId, @Param("startTs") long startTs,
@Param("entityKey") int key, @Param("endTs") long endTs,
@Param("startTs") long startTs, Pageable pageable);
@Param("endTs") long endTs, Pageable pageable);
@Transactional @Transactional
@Modifying @Modifying

11
dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvRepository.java

@ -29,8 +29,15 @@ import java.util.UUID;
public interface TsKvRepository extends JpaRepository<TsKvEntity, TsKvCompositeKey> { public interface TsKvRepository extends JpaRepository<TsKvEntity, TsKvCompositeKey> {
@Query("SELECT tskv FROM TsKvEntity tskv WHERE tskv.entityId = :entityId " + /*
"AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs") * Using native query to avoid adding 'nulls first' or 'nulls last' (ignoring spring.jpa.properties.hibernate.order_by.default_null_ordering)
* to the order so that index scan is done instead of full scan.
*
* Note: even when setting custom NullHandling for the Sort.Order for non-native queries,
* it will be ignored and default_null_ordering will be used
* */
@Query(value = "SELECT * FROM ts_kv WHERE entity_id = :entityId " +
"AND key = :entityKey AND ts >= :startTs AND ts < :endTs ", nativeQuery = true)
List<TsKvEntity> findAllWithLimit(@Param("entityId") UUID entityId, List<TsKvEntity> findAllWithLimit(@Param("entityId") UUID entityId,
@Param("entityKey") int key, @Param("entityKey") int key,
@Param("startTs") long startTs, @Param("startTs") long startTs,

21
dao/src/main/resources/cassandra/schema-keyspace.cql

@ -0,0 +1,21 @@
--
-- Copyright © 2016-2022 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.
--
CREATE KEYSPACE IF NOT EXISTS thingsboard
WITH replication = {
'class' : 'SimpleStrategy',
'replication_factor' : 1
};

6
dao/src/main/resources/cassandra/schema-ts-latest.cql

@ -14,12 +14,6 @@
-- limitations under the License. -- limitations under the License.
-- --
CREATE KEYSPACE IF NOT EXISTS thingsboard
WITH replication = {
'class' : 'SimpleStrategy',
'replication_factor' : 1
};
CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_latest_cf ( CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_latest_cf (
entity_type text, // (DEVICE, CUSTOMER, TENANT) entity_type text, // (DEVICE, CUSTOMER, TENANT)
entity_id timeuuid, entity_id timeuuid,

6
dao/src/main/resources/cassandra/schema-ts.cql

@ -14,12 +14,6 @@
-- limitations under the License. -- limitations under the License.
-- --
CREATE KEYSPACE IF NOT EXISTS thingsboard
WITH replication = {
'class' : 'SimpleStrategy',
'replication_factor' : 1
};
CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_cf ( CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_cf (
entity_type text, // (DEVICE, CUSTOMER, TENANT) entity_type text, // (DEVICE, CUSTOMER, TENANT)
entity_id timeuuid, entity_id timeuuid,

2
dao/src/main/resources/sql/schema-entities-idx.sql

@ -50,6 +50,8 @@ CREATE INDEX IF NOT EXISTS idx_attribute_kv_by_key_and_last_update_ts ON attribu
CREATE INDEX IF NOT EXISTS idx_audit_log_tenant_id_and_created_time ON audit_log(tenant_id, created_time DESC); CREATE INDEX IF NOT EXISTS idx_audit_log_tenant_id_and_created_time ON audit_log(tenant_id, created_time DESC);
CREATE INDEX IF NOT EXISTS idx_edge_event_tenant_id_and_created_time ON edge_event(tenant_id, created_time DESC);
CREATE INDEX IF NOT EXISTS idx_rpc_tenant_id_device_id ON rpc(tenant_id, device_id); CREATE INDEX IF NOT EXISTS idx_rpc_tenant_id_device_id ON rpc(tenant_id, device_id);
CREATE INDEX IF NOT EXISTS idx_device_external_id ON device(tenant_id, external_id); CREATE INDEX IF NOT EXISTS idx_device_external_id ON device(tenant_id, external_id);

4
dao/src/main/resources/sql/schema-entities.sql

@ -719,7 +719,7 @@ CREATE TABLE IF NOT EXISTS edge (
); );
CREATE TABLE IF NOT EXISTS edge_event ( CREATE TABLE IF NOT EXISTS edge_event (
id uuid NOT NULL CONSTRAINT edge_event_pkey PRIMARY KEY, id uuid NOT NULL,
created_time bigint NOT NULL, created_time bigint NOT NULL,
edge_id uuid, edge_id uuid,
edge_event_type varchar(255), edge_event_type varchar(255),
@ -729,7 +729,7 @@ CREATE TABLE IF NOT EXISTS edge_event (
body varchar(10000000), body varchar(10000000),
tenant_id uuid, tenant_id uuid,
ts bigint NOT NULL ts bigint NOT NULL
); ) PARTITION BY RANGE(created_time);
CREATE TABLE IF NOT EXISTS rpc ( CREATE TABLE IF NOT EXISTS rpc (
id uuid NOT NULL CONSTRAINT rpc_pkey PRIMARY KEY, id uuid NOT NULL CONSTRAINT rpc_pkey PRIMARY KEY,

1
dao/src/test/java/org/thingsboard/server/dao/NoSqlDaoServiceTestSuite.java

@ -33,6 +33,7 @@ public class NoSqlDaoServiceTestSuite {
public static CustomCassandraCQLUnit cassandraUnit = public static CustomCassandraCQLUnit cassandraUnit =
new CustomCassandraCQLUnit( new CustomCassandraCQLUnit(
Arrays.asList( Arrays.asList(
new ClassPathCQLDataSet("cassandra/schema-keyspace.cql", false, false),
new ClassPathCQLDataSet("cassandra/schema-ts.cql", false, false), new ClassPathCQLDataSet("cassandra/schema-ts.cql", false, false),
new ClassPathCQLDataSet("cassandra/schema-ts-latest.cql", false, false) new ClassPathCQLDataSet("cassandra/schema-ts-latest.cql", false, false)
), ),

3
docker/tb-js-executor.env

@ -3,4 +3,5 @@ LOGGER_LEVEL=info
LOG_FOLDER=logs LOG_FOLDER=logs
LOGGER_FILENAME=tb-js-executor-%DATE%.log LOGGER_FILENAME=tb-js-executor-%DATE%.log
DOCKER_MODE=true DOCKER_MODE=true
SCRIPT_BODY_TRACE_FREQUENCY=1000 SCRIPT_BODY_TRACE_FREQUENCY=1000
NODE_OPTIONS="--max-old-space-size=200"

2
msa/js-executor/api/jsExecutor.models.ts

@ -56,7 +56,7 @@ export interface JsCompileResponse extends TbMessage {
export interface JsInvokeResponse { export interface JsInvokeResponse {
success: boolean; success: boolean;
result: string; result?: string;
errorCode?: number; errorCode?: number;
errorDetails?: string; errorDetails?: string;
} }

11
msa/js-executor/api/jsInvokeMessageProcessor.ts

@ -39,6 +39,7 @@ const TIMEOUT_ERROR = 2;
const NOT_FOUND_ERROR = 3; const NOT_FOUND_ERROR = 3;
const statFrequency = Number(config.get('script.stat_print_frequency')); const statFrequency = Number(config.get('script.stat_print_frequency'));
const memoryUsageTraceFrequency = Number(config.get('script.memory_usage_trace_frequency'));
const scriptBodyTraceFrequency = Number(config.get('script.script_body_trace_frequency')); const scriptBodyTraceFrequency = Number(config.get('script.script_body_trace_frequency'));
const useSandbox = config.get('script.use_sandbox') === 'true'; const useSandbox = config.get('script.use_sandbox') === 'true';
const maxActiveScripts = Number(config.get('script.max_active_scripts')); const maxActiveScripts = Number(config.get('script.max_active_scripts'));
@ -167,11 +168,15 @@ export class JsInvokeMessageProcessor {
if (this.executedScriptsCounter % scriptBodyTraceFrequency == 0) { if (this.executedScriptsCounter % scriptBodyTraceFrequency == 0) {
this.logger.info('[%s] Executing script body: [%s]', scriptId, invokeRequest.scriptBody); this.logger.info('[%s] Executing script body: [%s]', scriptId, invokeRequest.scriptBody);
} }
if (this.executedScriptsCounter % memoryUsageTraceFrequency == 0) {
this.logger.info('Current memory usage: [%s]', process.memoryUsage());
}
this.getOrCompileScript(scriptId, invokeRequest.scriptBody).then( this.getOrCompileScript(scriptId, invokeRequest.scriptBody).then(
(script) => { (script) => {
this.executor.executeScript(script, invokeRequest.args, invokeRequest.timeout).then( this.executor.executeScript(script, invokeRequest.args, invokeRequest.timeout).then(
(result) => { (result: string | undefined) => {
if (result.length <= maxResultSize) { if (!result || result.length <= maxResultSize) {
const invokeResponse = JsInvokeMessageProcessor.createInvokeResponse(result, true); const invokeResponse = JsInvokeMessageProcessor.createInvokeResponse(result, true);
this.logger.debug('[%s] Sending success invoke response, scriptId: [%s]', requestId, scriptId); this.logger.debug('[%s] Sending success invoke response, scriptId: [%s]', requestId, scriptId);
this.sendResponse(requestId, responseTopic, headers, scriptId, undefined, invokeResponse); this.sendResponse(requestId, responseTopic, headers, scriptId, undefined, invokeResponse);
@ -323,7 +328,7 @@ export class JsInvokeMessageProcessor {
} }
} }
private static createInvokeResponse(result: string, success: boolean, errorCode?: number, err?: any): JsInvokeResponse { private static createInvokeResponse(result: string | undefined, success: boolean, errorCode?: number, err?: any): JsInvokeResponse {
return { return {
errorCode: errorCode, errorCode: errorCode,
success: success, success: success,

1
msa/js-executor/config/custom-environment-variables.yml

@ -75,6 +75,7 @@ logger:
script: script:
use_sandbox: "SCRIPT_USE_SANDBOX" use_sandbox: "SCRIPT_USE_SANDBOX"
memory_usage_trace_frequency: "MEMORY_USAGE_TRACE_FREQUENCY"
stat_print_frequency: "SCRIPT_STAT_PRINT_FREQUENCY" stat_print_frequency: "SCRIPT_STAT_PRINT_FREQUENCY"
script_body_trace_frequency: "SCRIPT_BODY_TRACE_FREQUENCY" script_body_trace_frequency: "SCRIPT_BODY_TRACE_FREQUENCY"
max_active_scripts: "MAX_ACTIVE_SCRIPTS" max_active_scripts: "MAX_ACTIVE_SCRIPTS"

1
msa/js-executor/config/default.yml

@ -64,6 +64,7 @@ logger:
script: script:
use_sandbox: "true" use_sandbox: "true"
memory_usage_trace_frequency: "1000"
script_body_trace_frequency: "10000" script_body_trace_frequency: "10000"
stat_print_frequency: "10000" stat_print_frequency: "10000"
max_active_scripts: "1000" max_active_scripts: "1000"

2
msa/js-executor/docker/start-js-executor.sh

@ -27,4 +27,4 @@ source "${CONF_FOLDER}/${configfile}"
cd ${pkg.installFolder} cd ${pkg.installFolder}
# This will forward this PID 1 to the node.js and forward SIGTERM for graceful shutdown as well # This will forward this PID 1 to the node.js and forward SIGTERM for graceful shutdown as well
exec node server.js exec node --no-compilation-cache server.js

2
pom.xml

@ -77,7 +77,7 @@
<zookeeper.version>3.5.5</zookeeper.version> <zookeeper.version>3.5.5</zookeeper.version>
<protobuf.version>3.21.9</protobuf.version> <protobuf.version>3.21.9</protobuf.version>
<grpc.version>1.42.1</grpc.version> <grpc.version>1.42.1</grpc.version>
<mvel.version>2.4.23TB</mvel.version> <mvel.version>2.4.25TB</mvel.version>
<lombok.version>1.18.18</lombok.version> <lombok.version>1.18.18</lombok.version>
<paho.client.version>1.2.4</paho.client.version> <paho.client.version>1.2.4</paho.client.version>
<netty.version>4.1.75.Final</netty.version> <netty.version>4.1.75.Final</netty.version>

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/edge/AbstractTbMsgPushNode.java

@ -142,8 +142,11 @@ public abstract class AbstractTbMsgPushNode<T extends BaseTbMsgPushNodeConfigura
actionType = EdgeEventActionType.ATTRIBUTES_UPDATED; actionType = EdgeEventActionType.ATTRIBUTES_UPDATED;
} else if (SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType)) { } else if (SessionMsgType.POST_ATTRIBUTES_REQUEST.name().equals(msgType)) {
actionType = EdgeEventActionType.POST_ATTRIBUTES; actionType = EdgeEventActionType.POST_ATTRIBUTES;
} else { } else if (DataConstants.ATTRIBUTES_DELETED.equals(msgType)) {
actionType = EdgeEventActionType.ATTRIBUTES_DELETED; actionType = EdgeEventActionType.ATTRIBUTES_DELETED;
} else {
log.warn("Unsupported msg type [{}]", msgType);
throw new IllegalArgumentException("Unsupported msg type: " + msgType);
} }
return actionType; return actionType;
} }

55
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNode.java

@ -15,15 +15,27 @@
*/ */
package org.thingsboard.rule.engine.rpc; package org.thingsboard.rule.engine.rpc;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.RuleNode; import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode; import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.api.util.TbNodeUtils; import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EdgeUtils;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
@ -65,9 +77,46 @@ public class TbSendRPCReplyNode implements TbNode {
} else if (StringUtils.isEmpty(msg.getData())) { } else if (StringUtils.isEmpty(msg.getData())) {
ctx.tellFailure(msg, new RuntimeException("Request body is empty!")); ctx.tellFailure(msg, new RuntimeException("Request body is empty!"));
} else { } else {
ctx.getRpcService().sendRpcReplyToDevice(serviceIdStr, UUID.fromString(sessionIdStr), Integer.parseInt(requestIdStr), msg.getData()); if (StringUtils.isNotBlank(msg.getMetaData().getValue(DataConstants.EDGE_ID))) {
ctx.tellSuccess(msg); saveRpcResponseToEdgeQueue(ctx, msg, serviceIdStr, sessionIdStr, requestIdStr);
} else {
ctx.getRpcService().sendRpcReplyToDevice(serviceIdStr, UUID.fromString(sessionIdStr), Integer.parseInt(requestIdStr), msg.getData());
ctx.tellSuccess(msg);
}
} }
} }
private void saveRpcResponseToEdgeQueue(TbContext ctx, TbMsg msg, String serviceIdStr, String sessionIdStr, String requestIdStr) {
EdgeId edgeId;
DeviceId deviceId;
try {
edgeId = new EdgeId(UUID.fromString(msg.getMetaData().getValue(DataConstants.EDGE_ID)));
deviceId = new DeviceId(UUID.fromString(msg.getMetaData().getValue(DataConstants.DEVICE_ID)));
} catch (Exception e) {
String errMsg = String.format("[%s] Failed to parse edgeId or deviceId from metadata %s!", ctx.getTenantId(), msg.getMetaData());
ctx.tellFailure(msg, new RuntimeException(errMsg));
return;
}
ObjectNode body = JacksonUtil.OBJECT_MAPPER.createObjectNode();
body.put("serviceId", serviceIdStr);
body.put("sessionId", sessionIdStr);
body.put("requestId", requestIdStr);
body.put("response", msg.getData());
EdgeEvent edgeEvent = EdgeUtils.constructEdgeEvent(ctx.getTenantId(), edgeId, EdgeEventType.DEVICE,
EdgeEventActionType.RPC_CALL, deviceId, JacksonUtil.OBJECT_MAPPER.valueToTree(body));
ListenableFuture<Void> future = ctx.getEdgeEventService().saveAsync(edgeEvent);
Futures.addCallback(future, new FutureCallback<>() {
@Override
public void onSuccess(Void result) {
ctx.onEdgeEventUpdate(ctx.getTenantId(), edgeId);
ctx.tellSuccess(msg);
}
@Override
public void onFailure(Throwable t) {
ctx.tellFailure(msg, t);
}
}, ctx.getDbCallbackExecutor());
}
} }

119
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/rpc/TbSendRPCReplyNodeTest.java

@ -0,0 +1,119 @@
/**
* Copyright © 2016-2022 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.rule.engine.rpc;
import com.google.common.util.concurrent.SettableFuture;
import org.junit.Before;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.Mock;
import org.mockito.Mockito;
import org.mockito.junit.MockitoJUnitRunner;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ListeningExecutor;
import org.thingsboard.rule.engine.api.RuleEngineRpcService;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.session.SessionMsgType;
import org.thingsboard.server.dao.edge.EdgeEventService;
import java.util.UUID;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
@RunWith(MockitoJUnitRunner.class)
public class TbSendRPCReplyNodeTest {
private static final String DUMMY_SERVICE_ID = "testServiceId";
private static final int DUMMY_REQUEST_ID = 0;
private static final UUID DUMMY_SESSION_ID = UUID.randomUUID();
private static final String DUMMY_DATA = "{\"key\":\"value\"}";
TbSendRPCReplyNode node;
private final TenantId tenantId = TenantId.fromUUID(UUID.randomUUID());
private final DeviceId deviceId = new DeviceId(UUID.randomUUID());
@Mock
private TbContext ctx;
@Mock
private RuleEngineRpcService rpcService;
@Mock
private EdgeEventService edgeEventService;
@Mock
private ListeningExecutor listeningExecutor;
@Before
public void setUp() throws TbNodeException {
node = new TbSendRPCReplyNode();
TbSendRpcReplyNodeConfiguration config = new TbSendRpcReplyNodeConfiguration().defaultConfiguration();
node.init(ctx, new TbNodeConfiguration(JacksonUtil.valueToTree(config)));
}
@Test
public void sendReplyToTransport() {
Mockito.when(ctx.getRpcService()).thenReturn(rpcService);
TbMsg msg = TbMsg.newMsg(SessionMsgType.POST_TELEMETRY_REQUEST.name(), deviceId, getDefaultMetadata(),
TbMsgDataType.JSON, DUMMY_DATA, null, null);
node.onMsg(ctx, msg);
verify(rpcService).sendRpcReplyToDevice(DUMMY_SERVICE_ID, DUMMY_SESSION_ID, DUMMY_REQUEST_ID, DUMMY_DATA);
verify(edgeEventService, never()).saveAsync(any());
}
@Test
public void sendReplyToEdgeQueue() {
Mockito.when(ctx.getTenantId()).thenReturn(tenantId);
Mockito.when(ctx.getEdgeEventService()).thenReturn(edgeEventService);
Mockito.when(edgeEventService.saveAsync(any())).thenReturn(SettableFuture.create());
Mockito.when(ctx.getDbCallbackExecutor()).thenReturn(listeningExecutor);
TbMsgMetaData defaultMetadata = getDefaultMetadata();
defaultMetadata.putValue(DataConstants.EDGE_ID, UUID.randomUUID().toString());
defaultMetadata.putValue(DataConstants.DEVICE_ID, UUID.randomUUID().toString());
TbMsg msg = TbMsg.newMsg(SessionMsgType.POST_TELEMETRY_REQUEST.name(), deviceId, defaultMetadata,
TbMsgDataType.JSON, DUMMY_DATA, null, null);
node.onMsg(ctx, msg);
verify(edgeEventService).saveAsync(any());
verify(rpcService, never()).sendRpcReplyToDevice(DUMMY_SERVICE_ID, DUMMY_SESSION_ID, DUMMY_REQUEST_ID, DUMMY_DATA);
}
private TbMsgMetaData getDefaultMetadata() {
TbSendRpcReplyNodeConfiguration config = new TbSendRpcReplyNodeConfiguration().defaultConfiguration();
TbMsgMetaData metadata = new TbMsgMetaData();
metadata.putValue(config.getServiceIdMetaDataAttribute(), DUMMY_SERVICE_ID);
metadata.putValue(config.getSessionIdMetaDataAttribute(), DUMMY_SESSION_ID.toString());
metadata.putValue(config.getRequestIdMetaDataAttribute(), Integer.toString(DUMMY_REQUEST_ID));
return metadata;
}
}

2
tools/src/main/java/org/thingsboard/client/tools/migrator/README.md

@ -63,7 +63,7 @@ Tool execution time depends on DB size, CPU resources and Disk throughput
* Note that this this part works only for single node Cassandra Cluster. If you have more nodes - it is better to use `sstableloader` tool. * Note that this this part works only for single node Cassandra Cluster. If you have more nodes - it is better to use `sstableloader` tool.
1. [Optional] install Cassandra on the instance 1. [Optional] install Cassandra on the instance
2. [Optional] Using `cqlsh` create `thingsboard` keyspace and requred tables from this files `schema-ts.cql` and `schema-ts-latest.cql` using `source` command 2. [Optional] Using `cqlsh` create `thingsboard` keyspace and requred tables from this files `schema-keyspace.cql`, `schema-ts.cql` and `schema-ts-latest.cql` using `source` command
3. Stop Cassandra 3. Stop Cassandra
4. Look at `/var/lib/cassandra/data/thingsboard` and check for names of data folders 4. Look at `/var/lib/cassandra/data/thingsboard` and check for names of data folders
5. Copy generated SSTable files into cassandra data dir using next command: 5. Copy generated SSTable files into cassandra data dir using next command:

4
ui-ngx/src/app/core/http/attribute.service.ts

@ -43,7 +43,7 @@ export class AttributeService {
public deleteEntityAttributes(entityId: EntityId, attributeScope: AttributeScope, attributes: Array<AttributeData>, public deleteEntityAttributes(entityId: EntityId, attributeScope: AttributeScope, attributes: Array<AttributeData>,
config?: RequestConfig): Observable<any> { config?: RequestConfig): Observable<any> {
const keys = attributes.map(attribute => encodeURI(attribute.key)).join(','); const keys = attributes.map(attribute => encodeURIComponent(attribute.key)).join(',');
return this.http.delete(`/api/plugins/telemetry/${entityId.entityType}/${entityId.id}/${attributeScope}` + return this.http.delete(`/api/plugins/telemetry/${entityId.entityType}/${entityId.id}/${attributeScope}` +
`?keys=${keys}`, `?keys=${keys}`,
defaultHttpOptionsFromConfig(config)); defaultHttpOptionsFromConfig(config));
@ -51,7 +51,7 @@ export class AttributeService {
public deleteEntityTimeseries(entityId: EntityId, timeseries: Array<AttributeData>, deleteAllDataForKeys = false, public deleteEntityTimeseries(entityId: EntityId, timeseries: Array<AttributeData>, deleteAllDataForKeys = false,
startTs?: number, endTs?: number, config?: RequestConfig): Observable<any> { startTs?: number, endTs?: number, config?: RequestConfig): Observable<any> {
const keys = timeseries.map(attribute => encodeURI(attribute.key)).join(','); const keys = timeseries.map(attribute => encodeURIComponent(attribute.key)).join(',');
let url = `/api/plugins/telemetry/${entityId.entityType}/${entityId.id}/timeseries/delete` + let url = `/api/plugins/telemetry/${entityId.entityType}/${entityId.id}/timeseries/delete` +
`?keys=${keys}&deleteAllDataForKeys=${deleteAllDataForKeys}`; `?keys=${keys}&deleteAllDataForKeys=${deleteAllDataForKeys}`;
if (isDefinedAndNotNull(startTs)) { if (isDefinedAndNotNull(startTs)) {

22
ui-ngx/src/app/modules/home/components/dashboard-page/layout/manage-dashboard-layouts-dialog.component.html

@ -65,12 +65,15 @@
mat-raised-button mat-raised-button
color="primary" color="primary"
class="tb-layout-button" class="tb-layout-button"
[matTooltip]="layoutButtonText('main')" (mouseover)="mainLayoutTooltip.show()"
matTooltipPosition="above" (mouseleave)="mainLayoutTooltip.hide()"
matTooltipClass="tb-layout-button-tooltip"
(click)="setFixedLayout('main')" (click)="setFixedLayout('main')"
[ngClass]="layoutButtonClass('main', true)"> [ngClass]="layoutButtonClass('main', true)">
<span>{{ (layoutsFormGroup.value.right ? 'layout.left' : 'layout.main') | translate }}</span> <span [matTooltip]="layoutButtonText('main')"
#mainLayoutTooltip="matTooltip"
matTooltipPosition="above">
{{ (layoutsFormGroup.value.right ? 'layout.left' : 'layout.main') | translate }}
</span>
</button> </button>
<div fxFlex class="tb-layout-preview-element tb-layout-preview-input" *ngIf="showPreviewInputs('main')"> <div fxFlex class="tb-layout-preview-element tb-layout-preview-input" *ngIf="showPreviewInputs('main')">
<input *ngIf="layoutsFormGroup.get('type').value !== layoutWidthType.FIXED" <input *ngIf="layoutsFormGroup.get('type').value !== layoutWidthType.FIXED"
@ -106,12 +109,15 @@
mat-raised-button mat-raised-button
color="primary" color="primary"
class="tb-layout-button tb-layout-button-right" class="tb-layout-button tb-layout-button-right"
[matTooltip]="layoutButtonText('right')" (mouseover)="rightLayoutTooltip.show()"
matTooltipPosition="above" (mouseleave)="rightLayoutTooltip.hide()"
matTooltipClass="tb-layout-button-tooltip"
(click)="setFixedLayout('right')" (click)="setFixedLayout('right')"
[ngClass]="layoutButtonClass('right')"> [ngClass]="layoutButtonClass('right')">
<span>{{ 'layout.right' | translate }}</span> <span [matTooltip]="layoutButtonText('right')"
#rightLayoutTooltip="matTooltip"
matTooltipPosition="above">
{{ 'layout.right' | translate }}
</span>
</button> </button>
<div fxFlex class="tb-layout-preview-element tb-layout-preview-input" *ngIf="showPreviewInputs('right')"> <div fxFlex class="tb-layout-preview-element tb-layout-preview-input" *ngIf="showPreviewInputs('right')">
<input *ngIf="layoutsFormGroup.get('type').value !== layoutWidthType.FIXED" <input *ngIf="layoutsFormGroup.get('type').value !== layoutWidthType.FIXED"

4
ui-ngx/src/app/modules/home/components/dashboard-page/layout/manage-dashboard-layouts-dialog.component.scss

@ -178,8 +178,4 @@ $tb-warn: mat.get-color-from-palette(map-get($tb-theme, warn), text);
width: 160px; width: 160px;
text-align: center; text-align: center;
} }
.tb-layout-button-tooltip {
margin: 30px 40px -35px -50px;
}
} }

31
ui-ngx/src/app/modules/home/components/dashboard-page/layout/manage-dashboard-layouts-dialog.component.ts

@ -14,7 +14,7 @@
/// limitations under the License. /// limitations under the License.
/// ///
import { Component, ElementRef, Inject, SkipSelf, ViewChild } from '@angular/core'; import { Component, Inject, SkipSelf, ViewChild } from '@angular/core';
import { ErrorStateMatcher } from '@angular/material/core'; import { ErrorStateMatcher } from '@angular/material/core';
import { MAT_DIALOG_DATA, MatDialog, MatDialogRef } from '@angular/material/dialog'; import { MAT_DIALOG_DATA, MatDialog, MatDialogRef } from '@angular/material/dialog';
import { Store } from '@ngrx/store'; import { Store } from '@ngrx/store';
@ -85,8 +85,7 @@ export class ManageDashboardLayoutsDialogComponent extends DialogComponent<Manag
private utils: UtilsService, private utils: UtilsService,
private dashboardUtils: DashboardUtilsService, private dashboardUtils: DashboardUtilsService,
private translate: TranslateService, private translate: TranslateService,
private dialog: MatDialog, private dialog: MatDialog) {
private elementRef: ElementRef) {
super(store, router, dialogRef); super(store, router, dialogRef);
this.layouts = this.data.layouts; this.layouts = this.data.layouts;
@ -96,10 +95,13 @@ export class ManageDashboardLayoutsDialogComponent extends DialogComponent<Manag
right: [isDefined(this.layouts.right)], right: [isDefined(this.layouts.right)],
sliderPercentage: [50], sliderPercentage: [50],
sliderFixed: [this.layoutFixedSize.MIN], sliderFixed: [this.layoutFixedSize.MIN],
leftWidthPercentage: [50, [Validators.min(this.layoutPercentageSize.MIN), Validators.max(this.layoutPercentageSize.MAX), Validators.required]], leftWidthPercentage: [50,
rightWidthPercentage: [50, [Validators.min(this.layoutPercentageSize.MIN), Validators.max(this.layoutPercentageSize.MAX), Validators.required]], [Validators.min(this.layoutPercentageSize.MIN), Validators.max(this.layoutPercentageSize.MAX), Validators.required]],
rightWidthPercentage: [50,
[Validators.min(this.layoutPercentageSize.MIN), Validators.max(this.layoutPercentageSize.MAX), Validators.required]],
type: [LayoutWidthType.PERCENTAGE], type: [LayoutWidthType.PERCENTAGE],
fixedWidth: [this.layoutFixedSize.MIN, [Validators.min(this.layoutFixedSize.MIN), Validators.max(this.layoutFixedSize.MAX), Validators.required]], fixedWidth: [this.layoutFixedSize.MIN,
[Validators.min(this.layoutFixedSize.MIN), Validators.max(this.layoutFixedSize.MAX), Validators.required]],
fixedLayout: ['main', []] fixedLayout: ['main', []]
} }
); );
@ -293,22 +295,7 @@ export class ManageDashboardLayoutsDialogComponent extends DialogComponent<Manag
} }
setFixedLayout(layout: string): void { setFixedLayout(layout: string): void {
const layoutButtons = this.elementRef.nativeElement.querySelectorAll('.tb-layout-button'); if (this.layoutsFormGroup.get('type').value === LayoutWidthType.FIXED && this.layoutsFormGroup.get('right').value) {
if (layoutButtons?.length) {
let elementToDisable: HTMLButtonElement;
if (layout === 'right') {
elementToDisable = layoutButtons[0];
} else {
elementToDisable = layoutButtons[1];
}
elementToDisable.disabled = true;
setTimeout(() => {
elementToDisable.disabled = false;
}, 250);
}
if (this.layoutsFormGroup.get('type').value === LayoutWidthType.FIXED) {
this.layoutsFormGroup.get('fixedLayout').setValue(layout); this.layoutsFormGroup.get('fixedLayout').setValue(layout);
} }
} }

Loading…
Cancel
Save