Browse Source

updated BaseAttributeService, InternalTelemetryService, RuleEngineTelemetryService to store old methods for backward compatibility

pull/9850/head
dashevchenko 3 years ago
parent
commit
56f9dce242
  1. 100
      application/src/main/data/upgrade/3.6.3/schema_update.sql
  2. 2
      application/src/main/java/org/thingsboard/server/controller/TelemetryController.java
  3. 5
      application/src/main/java/org/thingsboard/server/service/device/ClaimDevicesServiceImpl.java
  4. 5
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/BaseTelemetryProcessor.java
  6. 5
      application/src/main/java/org/thingsboard/server/service/entitiy/entityview/DefaultTbEntityViewService.java
  7. 60
      application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java
  8. 5
      application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java
  9. 4
      application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
  10. 4
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java
  11. 31
      application/src/main/java/org/thingsboard/server/service/subscription/TbAttributeSubscriptionScope.java
  12. 3
      application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java
  13. 95
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  14. 7
      application/src/main/java/org/thingsboard/server/service/telemetry/InternalTelemetryService.java
  15. 23
      application/src/test/java/org/thingsboard/server/service/state/DefaultDeviceStateServiceTest.java
  16. 22
      common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java
  17. 6
      dao/src/main/java/org/thingsboard/server/dao/attributes/AttributeUtils.java
  18. 57
      dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java
  19. 36
      dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java
  20. 2
      dao/src/main/resources/sql/schema-entities.sql
  21. 2
      dao/src/main/resources/sql/schema-timescale.sql
  22. 2
      dao/src/main/resources/sql/schema-ts-psql.sql
  23. 37
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java
  24. 13
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCopyAttributesToEntityViewNode.java
  25. 5
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/math/TbMathNode.java
  26. 18
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java
  27. 16
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgDeleteAttributesNode.java
  28. 4
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java
  29. 3
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/profile/DeviceStateTest.java
  30. 9
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNodeTest.java
  31. 3
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgDeleteAttributesNodeTest.java

100
application/src/main/data/upgrade/3.6.3/schema_update.sql

@ -21,14 +21,11 @@ $$
BEGIN BEGIN
-- in case of running the upgrade script a second time: -- in case of running the upgrade script a second time:
IF EXISTS(SELECT 1 FROM information_schema.columns WHERE table_name = 'attribute_kv' and column_name='entity_type') THEN IF EXISTS(SELECT 1 FROM information_schema.columns WHERE table_name = 'attribute_kv' and column_name='entity_type') THEN
IF EXISTS(SELECT 1 FROM pg_indexes WHERE indexname = 'idx_attribute_kv_by_key_and_last_update_ts') THEN ALTER INDEX IF EXISTS idx_attribute_kv_by_key_and_last_update_ts RENAME TO idx_attribute_kv_by_key_and_last_update_ts_old;
ALTER INDEX idx_attribute_kv_by_key_and_last_update_ts RENAME TO idx_attribute_kv_by_key_and_last_update_ts_old;
END IF;
IF EXISTS(SELECT 1 FROM pg_constraint WHERE conname = 'attribute_kv_pkey') THEN IF EXISTS(SELECT 1 FROM pg_constraint WHERE conname = 'attribute_kv_pkey') THEN
ALTER TABLE attribute_kv RENAME CONSTRAINT attribute_kv_pkey TO attribute_kv_pkey_old; ALTER TABLE attribute_kv RENAME CONSTRAINT attribute_kv_pkey TO attribute_kv_pkey_old;
END IF; END IF;
ALTER TABLE attribute_kv ALTER TABLE attribute_kv RENAME TO attribute_kv_old;
RENAME TO attribute_kv_old;
CREATE TABLE IF NOT EXISTS attribute_kv CREATE TABLE IF NOT EXISTS attribute_kv
( (
entity_id uuid, entity_id uuid,
@ -43,6 +40,8 @@ $$
CONSTRAINT attribute_kv_pkey PRIMARY KEY (entity_id, attribute_type, attribute_key) CONSTRAINT attribute_kv_pkey PRIMARY KEY (entity_id, attribute_type, attribute_key)
); );
END IF; END IF;
DROP VIEW IF EXISTS device_info_view;
DROP VIEW IF EXISTS device_info_active_attribute_view;
END; END;
$$; $$;
@ -51,14 +50,12 @@ DO
$$ $$
BEGIN BEGIN
IF EXISTS(SELECT 1 FROM information_schema.tables WHERE table_name = 'ts_kv_dictionary') THEN IF EXISTS(SELECT 1 FROM information_schema.tables WHERE table_name = 'ts_kv_dictionary') THEN
ALTER TABLE ts_kv_dictionary ALTER TABLE ts_kv_dictionary RENAME CONSTRAINT ts_key_id_pkey TO key_dictionary_id_pkey;
RENAME CONSTRAINT ts_key_id_pkey TO key_id_pkey; ALTER TABLE ts_kv_dictionary RENAME TO key_dictionary;
ALTER TABLE ts_kv_dictionary
RENAME TO key_dictionary;
ELSE CREATE TABLE IF NOT EXISTS key_dictionary( ELSE CREATE TABLE IF NOT EXISTS key_dictionary(
key varchar(255) NOT NULL, key varchar(255) NOT NULL,
key_id serial UNIQUE, key_id serial UNIQUE,
CONSTRAINT key_id_pkey PRIMARY KEY (key) CONSTRAINT key_dictionary_id_pkey PRIMARY KEY (key)
); );
END IF; END IF;
END; END;
@ -83,83 +80,40 @@ $$ LANGUAGE plpgsql;
-- insert keys into key_dictionary -- insert keys into key_dictionary
DO DO
$$ $$
DECLARE BEGIN
insert_record RECORD; IF EXISTS(SELECT 1 FROM information_schema.tables WHERE table_name = 'attribute_kv_old') THEN
key_cursor refcursor; INSERT INTO key_dictionary(key) SELECT DISTINCT attribute_key FROM attribute_kv_old ON CONFLICT DO NOTHING;
BEGIN END IF;
IF EXISTS(SELECT 1 FROM information_schema.tables WHERE table_name = 'attribute_kv_old') THEN END;
OPEN key_cursor FOR SELECT DISTINCT attribute_key
FROM attribute_kv_old
ORDER BY attribute_key;
LOOP
FETCH key_cursor INTO insert_record;
EXIT WHEN NOT FOUND;
IF NOT EXISTS(SELECT key FROM key_dictionary WHERE key = insert_record.attribute_key) THEN
INSERT INTO key_dictionary(key) VALUES (insert_record.attribute_key);
END IF;
END LOOP;
CLOSE key_cursor;
END IF;
END;
$$; $$;
-- create procedure to migrate all rows from attribute_kv_old to attribute_kv -- migrate attributes from attribute_kv_old to attribute_kv
CREATE OR REPLACE PROCEDURE insert_into_attribute_kv(IN path_to_file varchar) DO
LANGUAGE plpgsql AS
$$ $$
DECLARE DECLARE
row_num_old integer; row_num_old integer;
row_num integer; row_num integer;
attribute_scope_array text[];
BEGIN BEGIN
attribute_scope_array := ARRAY['SERVER_SCOPE', 'CLIENT_SCOPE', 'SHARED_SCOPE'];
IF EXISTS(SELECT 1 FROM information_schema.tables WHERE table_name = 'attribute_kv_old') THEN IF EXISTS(SELECT 1 FROM information_schema.tables WHERE table_name = 'attribute_kv_old') THEN
EXECUTE format('COPY (SELECT records.entity_id AS entity_id, INSERT INTO attribute_kv(entity_id, attribute_type, attribute_key, bool_v, str_v, long_v, dbl_v, json_v, last_update_ts)
to_attribute_type_id(records.attribute_type) AS attribute_type, SELECT a.entity_id, to_attribute_type_id(a.attribute_type), k.key_id, a.bool_v, a.str_v, a.long_v, a.dbl_v, a.json_v, a.last_update_ts
records.attribute_key AS attribute_key, FROM attribute_kv_old a INNER JOIN key_dictionary k ON (a.attribute_key = k.key)
records.bool_v AS bool_v, WHERE a.attribute_type IN ('SERVER_SCOPE', 'CLIENT_SCOPE', 'SHARED_SCOPE');
records.str_v AS str_v,
records.long_v AS long_v,
records.dbl_v AS dbl_v,
records.json_v AS json_v,
records.last_update_ts AS last_update_ts
FROM (SELECT entity_id,
attribute_type,
key_id AS attribute_key,
bool_v,
str_v,
long_v,
dbl_v,
json_v,
last_update_ts
FROM attribute_kv_old INNER JOIN key_dictionary ON (attribute_kv_old.attribute_key = key_dictionary.key)
WHERE attribute_type= ANY(%L)) AS records) TO %L;', attribute_scope_array, path_to_file);
EXECUTE format('COPY attribute_kv FROM %L', path_to_file);
SELECT COUNT(*) INTO row_num_old FROM attribute_kv_old; SELECT COUNT(*) INTO row_num_old FROM attribute_kv_old;
SELECT COUNT(*) INTO row_num FROM attribute_kv; SELECT COUNT(*) INTO row_num FROM attribute_kv;
RAISE NOTICE 'Migrated % of % rows', row_num, row_num_old; RAISE NOTICE 'Migrated % of % rows', row_num, row_num_old;
IF row_num != 0 THEN
DROP TABLE IF EXISTS attribute_kv_old;
ELSE
RAISE EXCEPTION 'Table attribute_kv is empty';
END IF;
CREATE INDEX IF NOT EXISTS idx_attribute_kv_by_key_and_last_update_ts ON attribute_kv(entity_id, attribute_key, last_update_ts desc);
END IF; END IF;
EXCEPTION EXCEPTION
WHEN others THEN WHEN others THEN
ROLLBACK; ROLLBACK;
RAISE EXCEPTION 'Error during COPY: %', SQLERRM; RAISE EXCEPTION 'Error during COPY: %', SQLERRM;
END END
$$; $$;
CREATE OR REPLACE PROCEDURE drop_attribute_kv_old_table()
LANGUAGE plpgsql AS
$$
DECLARE
row_num integer;
BEGIN
SELECT COUNT(*) INTO row_num FROM attribute_kv;
IF row_num != 0 then
DROP TABLE IF EXISTS attribute_kv_old;
DROP PROCEDURE IF EXISTS insert_into_attribute_kv(IN path_to_file varchar);
ELSE
RAISE EXCEPTION 'Table attribute_kv is empty';
END IF;
RETURN;
END;
$$;

2
application/src/main/java/org/thingsboard/server/controller/TelemetryController.java

@ -631,7 +631,7 @@ public class TelemetryController extends BaseController {
} }
SecurityUser user = getCurrentUser(); SecurityUser user = getCurrentUser();
return accessValidator.validateEntityAndCallback(getCurrentUser(), Operation.WRITE_ATTRIBUTES, entityIdSrc, (result, tenantId, entityId) -> { return accessValidator.validateEntityAndCallback(getCurrentUser(), Operation.WRITE_ATTRIBUTES, entityIdSrc, (result, tenantId, entityId) -> {
tsSubService.saveAndNotify(tenantId, entityId, scope.name(), attributes, new FutureCallback<Void>() { tsSubService.saveAndNotify(tenantId, entityId, scope, attributes, new FutureCallback<Void>() {
@Override @Override
public void onSuccess(@Nullable Void tmp) { public void onSuccess(@Nullable Void tmp) {
logAttributesUpdated(user, entityId, scope, attributes, null); logAttributesUpdated(user, entityId, scope, attributes, null);

5
application/src/main/java/org/thingsboard/server/service/device/ClaimDevicesServiceImpl.java

@ -31,7 +31,6 @@ import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
@ -187,7 +186,7 @@ public class ClaimDevicesServiceImpl implements ClaimDevicesService {
} }
SettableFuture<ReclaimResult> result = SettableFuture.create(); SettableFuture<ReclaimResult> result = SettableFuture.create();
telemetryService.saveAndNotify( telemetryService.saveAndNotify(
tenantId, savedDevice.getId(), DataConstants.SERVER_SCOPE, Collections.singletonList( tenantId, savedDevice.getId(), AttributeScope.SERVER_SCOPE, Collections.singletonList(
new BaseAttributeKvEntry(new BooleanDataEntry(CLAIM_ATTRIBUTE_NAME, true), System.currentTimeMillis()) new BaseAttributeKvEntry(new BooleanDataEntry(CLAIM_ATTRIBUTE_NAME, true), System.currentTimeMillis())
), ),
new FutureCallback<>() { new FutureCallback<>() {
@ -230,7 +229,7 @@ public class ClaimDevicesServiceImpl implements ClaimDevicesService {
} }
SettableFuture<Void> result = SettableFuture.create(); SettableFuture<Void> result = SettableFuture.create();
telemetryService.deleteAndNotify(device.getTenantId(), telemetryService.deleteAndNotify(device.getTenantId(),
device.getId(), DataConstants.SERVER_SCOPE, Arrays.asList(CLAIM_ATTRIBUTE_NAME, CLAIM_DATA_ATTRIBUTE_NAME), new FutureCallback<>() { device.getId(), AttributeScope.SERVER_SCOPE, Arrays.asList(CLAIM_ATTRIBUTE_NAME, CLAIM_DATA_ATTRIBUTE_NAME), new FutureCallback<>() {
@Override @Override
public void onSuccess(@Nullable Void tmp) { public void onSuccess(@Nullable Void tmp) {
result.set(tmp); result.set(tmp);

5
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java

@ -29,6 +29,7 @@ import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.ResourceUtils; import org.thingsboard.server.common.data.ResourceUtils;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
@ -412,7 +413,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
Collections.singletonList(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry(key, value))), Collections.singletonList(new BasicTsKvEntry(System.currentTimeMillis(), new LongDataEntry(key, value))),
new AttributeSaveCallback(tenantId, edgeId, key, value)); new AttributeSaveCallback(tenantId, edgeId, key, value));
} else { } else {
tsSubService.saveAttrAndNotify(tenantId, edgeId, DataConstants.SERVER_SCOPE, key, value, new AttributeSaveCallback(tenantId, edgeId, key, value)); tsSubService.saveAttrAndNotify(tenantId, edgeId, AttributeScope.SERVER_SCOPE, key, value, new AttributeSaveCallback(tenantId, edgeId, key, value));
} }
} }
@ -424,7 +425,7 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
Collections.singletonList(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry(key, value))), Collections.singletonList(new BasicTsKvEntry(System.currentTimeMillis(), new BooleanDataEntry(key, value))),
new AttributeSaveCallback(tenantId, edgeId, key, value)); new AttributeSaveCallback(tenantId, edgeId, key, value));
} else { } else {
tsSubService.saveAttrAndNotify(tenantId, edgeId, DataConstants.SERVER_SCOPE, key, value, new AttributeSaveCallback(tenantId, edgeId, key, value)); tsSubService.saveAttrAndNotify(tenantId, edgeId, AttributeScope.SERVER_SCOPE, key, value, new AttributeSaveCallback(tenantId, edgeId, key, value));
} }
} }

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/telemetry/BaseTelemetryProcessor.java

@ -252,7 +252,7 @@ public abstract class BaseTelemetryProcessor extends BaseEdgeProcessor {
JsonObject json = JsonUtils.getJsonObject(msg.getKvList()); JsonObject json = JsonUtils.getJsonObject(msg.getKvList());
List<AttributeKvEntry> attributes = new ArrayList<>(JsonConverter.convertToAttributes(json)); List<AttributeKvEntry> attributes = new ArrayList<>(JsonConverter.convertToAttributes(json));
String scope = metaData.getValue("scope"); String scope = metaData.getValue("scope");
tsSubService.saveAndNotify(tenantId, entityId, scope, attributes, new FutureCallback<Void>() { tsSubService.saveAndNotify(tenantId, entityId, AttributeScope.valueOf(scope), attributes, new FutureCallback<Void>() {
@Override @Override
public void onSuccess(@Nullable Void tmp) { public void onSuccess(@Nullable Void tmp) {
var defaultQueueAndRuleChain = getDefaultQueueNameAndRuleChainId(tenantId, entityId); var defaultQueueAndRuleChain = getDefaultQueueNameAndRuleChainId(tenantId, entityId);

5
application/src/main/java/org/thingsboard/server/service/entitiy/entityview/DefaultTbEntityViewService.java

@ -26,7 +26,6 @@ import org.springframework.stereotype.Service;
import org.springframework.util.ConcurrentReferenceHashMap; import org.springframework.util.ConcurrentReferenceHashMap;
import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.User;
@ -288,7 +287,7 @@ public class DefaultTbEntityViewService extends AbstractTbEntityService implemen
(startTime == 0 && endTime > lastUpdateTs) || (startTime == 0 && endTime > lastUpdateTs) ||
(startTime < lastUpdateTs && endTime > lastUpdateTs); (startTime < lastUpdateTs && endTime > lastUpdateTs);
}).collect(Collectors.toList()); }).collect(Collectors.toList());
tsSubService.saveAndNotify(entityView.getTenantId(), entityId, scope.name(), attributes, new FutureCallback<Void>() { tsSubService.saveAndNotify(entityView.getTenantId(), entityId, scope, attributes, new FutureCallback<Void>() {
@Override @Override
public void onSuccess(@Nullable Void tmp) { public void onSuccess(@Nullable Void tmp) {
try { try {
@ -356,7 +355,7 @@ public class DefaultTbEntityViewService extends AbstractTbEntityService implemen
EntityViewId entityId = entityView.getId(); EntityViewId entityId = entityView.getId();
SettableFuture<Void> resultFuture = SettableFuture.create(); SettableFuture<Void> resultFuture = SettableFuture.create();
if (keys != null && !keys.isEmpty()) { if (keys != null && !keys.isEmpty()) {
tsSubService.deleteAndNotify(entityView.getTenantId(), entityId, scope.name(), keys, new FutureCallback<Void>() { tsSubService.deleteAndNotify(entityView.getTenantId(), entityId, scope, keys, new FutureCallback<Void>() {
@Override @Override
public void onSuccess(@Nullable Void tmp) { public void onSuccess(@Nullable Void tmp) {
try { try {

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

@ -779,31 +779,7 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
}); });
break; break;
case "3.6.3": case "3.6.3":
updateSchema("3.6.3", 3006003, "3.7.0", 3007000, connection -> { updateSchema("3.6.3", 3006003, "3.7.0", 3007000, null);
try {
Path pathToTempAttributeKvFile;
if (SystemUtils.IS_OS_WINDOWS) {
pathToTempAttributeKvFile = createTempFileWindows("attribute_kv_temp",".sql");
} else {
pathToTempAttributeKvFile = createTempFile("attribute_kv", "attribute_kv_temp.sql");
}
executeQuery(connection, "call insert_into_attribute_kv('" + pathToTempAttributeKvFile + "')");
// remove attribute_kv_old
executeQuery(connection, "call drop_attribute_kv_old_table()");
//create index for new table attribute_kv
executeQuery(connection, "CREATE INDEX IF NOT EXISTS idx_attribute_kv_by_key_and_last_update_ts ON attribute_kv(entity_id, attribute_key, last_update_ts desc);");
// remove temp files
boolean deleteTsKvFile = Files.deleteIfExists(pathToTempAttributeKvFile);
if (deleteTsKvFile) {
log.info("Successfully deleted the temp file for attribute_kv table upgrade!");
}
} catch (Exception e) {
log.error("Failed updating schema!!!", e);
}
});
break; break;
default: default:
throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion); throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion);
@ -829,40 +805,6 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService
} }
} }
private static Path createTempFile(String tempDirectoryName, String tempFileName) throws IOException {
Path pathToTempAttributeKvFile;
Path tempDirPath = Files.createTempDirectory(tempDirectoryName);
File tempDirAsFile = tempDirPath.toFile();
boolean writable = tempDirAsFile.setWritable(true, false);
boolean readable = tempDirAsFile.setReadable(true, false);
boolean executable = tempDirAsFile.setExecutable(true, false);
pathToTempAttributeKvFile = tempDirPath.resolve(tempFileName).toAbsolutePath();
if (!(writable && readable && executable)) {
throw new RuntimeException("Failed to grant write permissions for the: " + tempDirPath + "folder!");
}
return pathToTempAttributeKvFile;
}
private static Path createTempFileWindows(String prefix, String suffix) throws IOException {
Path pathToTempAttributeKvFile;
log.info("Lookup for environment variable: {} ...", THINGSBOARD_WINDOWS_UPGRADE_DIR);
Path pathToDir;
String thingsboardWindowsUpgradeDir = System.getenv("THINGSBOARD_WINDOWS_UPGRADE_DIR");
if (StringUtils.isNotEmpty(thingsboardWindowsUpgradeDir)) {
log.info("Environment variable: {} was found!", THINGSBOARD_WINDOWS_UPGRADE_DIR);
pathToDir = Paths.get(thingsboardWindowsUpgradeDir);
} else {
log.info("Failed to lookup environment variable: {}", THINGSBOARD_WINDOWS_UPGRADE_DIR);
pathToDir = Paths.get(PATH_TO_USERS_PUBLIC_FOLDER);
}
log.info("Directory: {} will be used for creation temporary upgrade files!", pathToDir);
Path attributeKvFile = Files.createTempFile(pathToDir, prefix, suffix);
pathToTempAttributeKvFile = attributeKvFile.toAbsolutePath();
return pathToTempAttributeKvFile;
}
private void runSchemaUpdateScript(Connection connection, String version) throws Exception { private void runSchemaUpdateScript(Connection connection, String version) throws Exception {
Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", version, SCHEMA_UPDATE_SQL); Path schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", version, SCHEMA_UPDATE_SQL);
loadSql(schemaUpdateFile, connection); loadSql(schemaUpdateFile, connection);

5
application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java

@ -20,6 +20,7 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.context.annotation.Lazy; import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.rule.engine.api.RuleEngineTelemetryService; import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg;
import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
@ -335,7 +336,7 @@ public class DefaultOtaPackageStateService implements OtaPackageStateService {
remove(device, otaPackageType, attrToRemove); remove(device, otaPackageType, attrToRemove);
telemetryService.saveAndNotify(tenantId, deviceId, DataConstants.SHARED_SCOPE, attributes, new FutureCallback<>() { telemetryService.saveAndNotify(tenantId, deviceId, AttributeScope.SHARED_SCOPE, attributes, new FutureCallback<>() {
@Override @Override
public void onSuccess(@Nullable Void tmp) { public void onSuccess(@Nullable Void tmp) {
log.trace("[{}] Success save attributes with target firmware!", deviceId); log.trace("[{}] Success save attributes with target firmware!", deviceId);
@ -353,7 +354,7 @@ public class DefaultOtaPackageStateService implements OtaPackageStateService {
} }
private void remove(Device device, OtaPackageType otaPackageType, List<String> attributesKeys) { private void remove(Device device, OtaPackageType otaPackageType, List<String> attributesKeys) {
telemetryService.deleteAndNotify(device.getTenantId(), device.getId(), DataConstants.SHARED_SCOPE, attributesKeys, telemetryService.deleteAndNotify(device.getTenantId(), device.getId(), AttributeScope.SHARED_SCOPE, attributesKeys,
new FutureCallback<>() { new FutureCallback<>() {
@Override @Override
public void onSuccess(@Nullable Void tmp) { public void onSuccess(@Nullable Void tmp) {

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

@ -806,7 +806,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
Collections.singletonList(new BasicTsKvEntry(getCurrentTimeMillis(), new LongDataEntry(key, value))), Collections.singletonList(new BasicTsKvEntry(getCurrentTimeMillis(), new LongDataEntry(key, value))),
new TelemetrySaveCallback<>(deviceId, key, value)); new TelemetrySaveCallback<>(deviceId, key, value));
} else { } else {
tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value)); tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, AttributeScope.SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value));
} }
} }
@ -817,7 +817,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService<Dev
Collections.singletonList(new BasicTsKvEntry(getCurrentTimeMillis(), new BooleanDataEntry(key, value))), Collections.singletonList(new BasicTsKvEntry(getCurrentTimeMillis(), new BooleanDataEntry(key, value))),
new TelemetrySaveCallback<>(deviceId, key, value)); new TelemetrySaveCallback<>(deviceId, key, value));
} else { } else {
tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value)); tsSubService.saveAttrAndNotify(TenantId.SYS_TENANT_ID, deviceId, AttributeScope.SERVER_SCOPE, key, value, new TelemetrySaveCallback<>(deviceId, key, value));
} }
} }

4
application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java

@ -449,8 +449,8 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer
} }
final Map<String, Long> keyStates = subscription.getKeyStates(); final Map<String, Long> keyStates = subscription.getKeyStates();
AttributeScope scope; AttributeScope scope;
if (subscription.getScope() != null && !TbAttributeSubscriptionScope.ANY_SCOPE.equals(subscription.getScope())) { if (subscription.getScope() != null && subscription.getScope().getAttributeScope() != null) {
scope = AttributeScope.valueOf(subscription.getScope().name()); scope = subscription.getScope().getAttributeScope();
} else { } else {
scope = AttributeScope.CLIENT_SCOPE; scope = AttributeScope.CLIENT_SCOPE;
} }

31
application/src/main/java/org/thingsboard/server/service/subscription/TbAttributeSubscriptionScope.java

@ -15,8 +15,37 @@
*/ */
package org.thingsboard.server.service.subscription; package org.thingsboard.server.service.subscription;
import org.thingsboard.server.common.data.AttributeScope;
public enum TbAttributeSubscriptionScope { public enum TbAttributeSubscriptionScope {
ANY_SCOPE, CLIENT_SCOPE, SHARED_SCOPE, SERVER_SCOPE ANY_SCOPE(),
CLIENT_SCOPE(AttributeScope.CLIENT_SCOPE),
SHARED_SCOPE(AttributeScope.SHARED_SCOPE),
SERVER_SCOPE(AttributeScope.SERVER_SCOPE);
private final AttributeScope attributeScope;
TbAttributeSubscriptionScope() {
this.attributeScope = null;
}
TbAttributeSubscriptionScope(AttributeScope attributeScope) {
this.attributeScope = attributeScope;
}
public AttributeScope getAttributeScope() {
return attributeScope;
}
public static TbAttributeSubscriptionScope of(AttributeScope attributeScope) {
for (TbAttributeSubscriptionScope scope : TbAttributeSubscriptionScope.values()) {
if (attributeScope == scope.getAttributeScope()) {
return scope;
}
}
throw new IllegalArgumentException("Unknown AttributeScope: " + attributeScope.name());
}
} }

3
application/src/main/java/org/thingsboard/server/service/sync/ie/importing/csv/AbstractBulkImportService.java

@ -29,6 +29,7 @@ import org.springframework.security.core.context.SecurityContextHolder;
import org.thingsboard.common.util.DonAsynchron; import org.thingsboard.common.util.DonAsynchron;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.HasAdditionalInfo; import org.thingsboard.server.common.data.HasAdditionalInfo;
import org.thingsboard.server.common.data.HasTenantId; import org.thingsboard.server.common.data.HasTenantId;
@ -230,7 +231,7 @@ public abstract class AbstractBulkImportService<E extends HasId<? extends Entity
List<AttributeKvEntry> attributes = new ArrayList<>(JsonConverter.convertToAttributes(kvsEntry.getValue())); List<AttributeKvEntry> attributes = new ArrayList<>(JsonConverter.convertToAttributes(kvsEntry.getValue()));
accessValidator.validateEntityAndCallback(user, Operation.WRITE_ATTRIBUTES, entity.getId(), (result, tenantId, entityId) -> { accessValidator.validateEntityAndCallback(user, Operation.WRITE_ATTRIBUTES, entity.getId(), (result, tenantId, entityId) -> {
tsSubscriptionService.saveAndNotify(tenantId, entityId, scope, attributes, new FutureCallback<>() { tsSubscriptionService.saveAndNotify(tenantId, entityId, AttributeScope.valueOf(scope), attributes, new FutureCallback<>() {
@Override @Override
public void onSuccess(Void unused) { public void onSuccess(Void unused) {

95
application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java

@ -242,19 +242,38 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
saveAndNotify(tenantId, entityId, scope, attributes, true, callback); saveAndNotify(tenantId, entityId, scope, attributes, true, callback);
} }
@Override
public void saveAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes, FutureCallback<Void> callback) {
saveAndNotify(tenantId, entityId, scope, attributes, true, callback);
}
@Override @Override
public void saveAndNotify(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback) { public void saveAndNotify(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback) {
checkInternalEntity(entityId); checkInternalEntity(entityId);
saveAndNotifyInternal(tenantId, entityId, scope, attributes, notifyDevice, callback); saveAndNotifyInternal(tenantId, entityId, scope, attributes, notifyDevice, callback);
} }
@Override
public void saveAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback) {
checkInternalEntity(entityId);
saveAndNotifyInternal(tenantId, entityId, scope, attributes, notifyDevice, callback);
}
@Override @Override
public void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback) { public void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback) {
ListenableFuture<List<String>> saveFuture = attrService.save(tenantId, entityId, AttributeScope.valueOf(scope), attributes); ListenableFuture<List<String>> saveFuture = attrService.save(tenantId, entityId, scope, attributes);
addVoidCallback(saveFuture, callback); addVoidCallback(saveFuture, callback);
addWsCallback(saveFuture, success -> onAttributesUpdate(tenantId, entityId, scope, attributes, notifyDevice)); addWsCallback(saveFuture, success -> onAttributesUpdate(tenantId, entityId, scope, attributes, notifyDevice));
} }
@Override
public void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback) {
ListenableFuture<List<String>> saveFuture = attrService.save(tenantId, entityId, scope, attributes);
addVoidCallback(saveFuture, callback);
addWsCallback(saveFuture, success -> onAttributesUpdate(tenantId, entityId, scope.name(), attributes, notifyDevice));
}
@Override @Override
public void saveLatestAndNotify(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, FutureCallback<Void> callback) { public void saveLatestAndNotify(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, FutureCallback<Void> callback) {
checkInternalEntity(entityId); checkInternalEntity(entityId);
@ -274,19 +293,38 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
deleteAndNotifyInternal(tenantId, entityId, scope, keys, false, callback); deleteAndNotifyInternal(tenantId, entityId, scope, keys, false, callback);
} }
@Override
public void deleteAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> keys, FutureCallback<Void> callback) {
checkInternalEntity(entityId);
deleteAndNotifyInternal(tenantId, entityId, scope, keys, false, callback);
}
@Override @Override
public void deleteAndNotify(TenantId tenantId, EntityId entityId, String scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback) { public void deleteAndNotify(TenantId tenantId, EntityId entityId, String scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback) {
checkInternalEntity(entityId); checkInternalEntity(entityId);
deleteAndNotifyInternal(tenantId, entityId, scope, keys, notifyDevice, callback); deleteAndNotifyInternal(tenantId, entityId, scope, keys, notifyDevice, callback);
} }
@Override
public void deleteAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback) {
checkInternalEntity(entityId);
deleteAndNotifyInternal(tenantId, entityId, scope, keys, notifyDevice, callback);
}
@Override @Override
public void deleteAndNotifyInternal(TenantId tenantId, EntityId entityId, String scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback) { public void deleteAndNotifyInternal(TenantId tenantId, EntityId entityId, String scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback) {
ListenableFuture<List<String>> deleteFuture = attrService.removeAll(tenantId, entityId, AttributeScope.valueOf(scope), keys); ListenableFuture<List<String>> deleteFuture = attrService.removeAll(tenantId, entityId, scope, keys);
addVoidCallback(deleteFuture, callback); addVoidCallback(deleteFuture, callback);
addWsCallback(deleteFuture, success -> onAttributesDelete(tenantId, entityId, scope, keys, notifyDevice)); addWsCallback(deleteFuture, success -> onAttributesDelete(tenantId, entityId, scope, keys, notifyDevice));
} }
@Override
public void deleteAndNotifyInternal(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback) {
ListenableFuture<List<String>> deleteFuture = attrService.removeAll(tenantId, entityId, scope, keys);
addVoidCallback(deleteFuture, callback);
addWsCallback(deleteFuture, success -> onAttributesDelete(tenantId, entityId, scope.name(), keys, notifyDevice));
}
@Override @Override
public void deleteLatest(TenantId tenantId, EntityId entityId, List<String> keys, FutureCallback<Void> callback) { public void deleteLatest(TenantId tenantId, EntityId entityId, List<String> keys, FutureCallback<Void> callback) {
checkInternalEntity(entityId); checkInternalEntity(entityId);
@ -328,24 +366,49 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
, System.currentTimeMillis())), callback); , System.currentTimeMillis())), callback);
} }
@Override
public void saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, long value, FutureCallback<Void> callback) {
saveAndNotify(tenantId, entityId, scope, Collections.singletonList(new BaseAttributeKvEntry(new LongDataEntry(key, value)
, System.currentTimeMillis())), callback);
}
@Override @Override
public void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, String value, FutureCallback<Void> callback) { public void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, String value, FutureCallback<Void> callback) {
saveAndNotify(tenantId, entityId, scope, Collections.singletonList(new BaseAttributeKvEntry(new StringDataEntry(key, value) saveAndNotify(tenantId, entityId, scope, Collections.singletonList(new BaseAttributeKvEntry(new StringDataEntry(key, value)
, System.currentTimeMillis())), callback); , System.currentTimeMillis())), callback);
} }
@Override
public void saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, String value, FutureCallback<Void> callback) {
saveAndNotify(tenantId, entityId, scope, Collections.singletonList(new BaseAttributeKvEntry(new StringDataEntry(key, value)
, System.currentTimeMillis())), callback);
}
@Override @Override
public void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, double value, FutureCallback<Void> callback) { public void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, double value, FutureCallback<Void> callback) {
saveAndNotify(tenantId, entityId, scope, Collections.singletonList(new BaseAttributeKvEntry(new DoubleDataEntry(key, value) saveAndNotify(tenantId, entityId, scope, Collections.singletonList(new BaseAttributeKvEntry(new DoubleDataEntry(key, value)
, System.currentTimeMillis())), callback); , System.currentTimeMillis())), callback);
} }
@Override
public void saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, double value, FutureCallback<Void> callback) {
saveAndNotify(tenantId, entityId, scope, Collections.singletonList(new BaseAttributeKvEntry(new DoubleDataEntry(key, value)
, System.currentTimeMillis())), callback);
}
@Override @Override
public void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, boolean value, FutureCallback<Void> callback) { public void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, boolean value, FutureCallback<Void> callback) {
saveAndNotify(tenantId, entityId, scope, Collections.singletonList(new BaseAttributeKvEntry(new BooleanDataEntry(key, value) saveAndNotify(tenantId, entityId, scope, Collections.singletonList(new BaseAttributeKvEntry(new BooleanDataEntry(key, value)
, System.currentTimeMillis())), callback); , System.currentTimeMillis())), callback);
} }
@Override
public void saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, boolean value, FutureCallback<Void> callback) {
saveAndNotify(tenantId, entityId, scope, Collections.singletonList(new BaseAttributeKvEntry(new BooleanDataEntry(key, value)
, System.currentTimeMillis())), callback);
}
@Override @Override
public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, long value) { public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, long value) {
SettableFuture<Void> future = SettableFuture.create(); SettableFuture<Void> future = SettableFuture.create();
@ -353,6 +416,13 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
return future; return future;
} }
@Override
public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, long value) {
SettableFuture<Void> future = SettableFuture.create();
saveAttrAndNotify(tenantId, entityId, scope, key, value, new VoidFutureCallback(future));
return future;
}
@Override @Override
public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, String value) { public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, String value) {
SettableFuture<Void> future = SettableFuture.create(); SettableFuture<Void> future = SettableFuture.create();
@ -360,6 +430,13 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
return future; return future;
} }
@Override
public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, String value) {
SettableFuture<Void> future = SettableFuture.create();
saveAttrAndNotify(tenantId, entityId, scope, key, value, new VoidFutureCallback(future));
return future;
}
@Override @Override
public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, double value) { public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, double value) {
SettableFuture<Void> future = SettableFuture.create(); SettableFuture<Void> future = SettableFuture.create();
@ -367,6 +444,13 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
return future; return future;
} }
@Override
public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, double value) {
SettableFuture<Void> future = SettableFuture.create();
saveAttrAndNotify(tenantId, entityId, scope, key, value, new VoidFutureCallback(future));
return future;
}
@Override @Override
public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, boolean value) { public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, boolean value) {
SettableFuture<Void> future = SettableFuture.create(); SettableFuture<Void> future = SettableFuture.create();
@ -374,6 +458,13 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
return future; return future;
} }
@Override
public ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, boolean value) {
SettableFuture<Void> future = SettableFuture.create();
saveAttrAndNotify(tenantId, entityId, scope, key, value, new VoidFutureCallback(future));
return future;
}
private void onAttributesUpdate(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice) { private void onAttributesUpdate(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice) {
forwardToSubscriptionManagerService(tenantId, entityId, subscriptionManagerService -> { forwardToSubscriptionManagerService(tenantId, entityId, subscriptionManagerService -> {
subscriptionManagerService.onAttributesUpdate(tenantId, entityId, scope, attributes, notifyDevice, TbCallback.EMPTY); subscriptionManagerService.onAttributesUpdate(tenantId, entityId, scope, attributes, notifyDevice, TbCallback.EMPTY);

7
application/src/main/java/org/thingsboard/server/service/telemetry/InternalTelemetryService.java

@ -17,6 +17,7 @@ package org.thingsboard.server.service.telemetry;
import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.FutureCallback;
import org.thingsboard.rule.engine.api.RuleEngineTelemetryService; import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
import org.thingsboard.server.common.data.AttributeScope;
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.AttributeKvEntry; import org.thingsboard.server.common.data.kv.AttributeKvEntry;
@ -33,12 +34,18 @@ public interface InternalTelemetryService extends RuleEngineTelemetryService {
void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, long ttl, FutureCallback<Integer> callback); void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, long ttl, FutureCallback<Integer> callback);
@Deprecated(since = "3.7.0")
void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback); void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback);
void saveAndNotifyInternal(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback);
void saveLatestAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, FutureCallback<Void> callback); void saveLatestAndNotifyInternal(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, FutureCallback<Void> callback);
@Deprecated(since = "3.7.0")
void deleteAndNotifyInternal(TenantId tenantId, EntityId entityId, String scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback); void deleteAndNotifyInternal(TenantId tenantId, EntityId entityId, String scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback);
void deleteAndNotifyInternal(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback);
void deleteLatestInternal(TenantId tenantId, EntityId entityId, List<String> keys, FutureCallback<Void> callback); void deleteLatestInternal(TenantId tenantId, EntityId entityId, List<String> keys, FutureCallback<Void> callback);
} }

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

@ -26,6 +26,7 @@ import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension; import org.mockito.junit.jupiter.MockitoExtension;
import org.springframework.test.util.ReflectionTestUtils; import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.DeviceIdInfo; 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.id.TenantId;
@ -251,7 +252,7 @@ public class DefaultDeviceStateServiceTest {
long newTimeout = System.currentTimeMillis() - deviceState.getLastActivityTime() + increase; long newTimeout = System.currentTimeMillis() - deviceState.getLastActivityTime() + increase;
service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout); service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout);
verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), any(), any()); verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(AttributeScope.class), eq(ACTIVITY_STATE), any(), any());
Thread.sleep(defaultTimeout + increase); Thread.sleep(defaultTimeout + increase);
service.checkStates(); service.checkStates();
activityVerify(false); activityVerify(false);
@ -290,7 +291,7 @@ public class DefaultDeviceStateServiceTest {
long newTimeout = 1; long newTimeout = 1;
Thread.sleep(newTimeout); Thread.sleep(newTimeout);
verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), any(), any()); verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(AttributeScope.class), eq(ACTIVITY_STATE), any(), any());
} }
@Test @Test
@ -311,7 +312,7 @@ public class DefaultDeviceStateServiceTest {
service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis()); service.onDeviceActivity(tenantId, deviceId, System.currentTimeMillis());
activityVerify(true); activityVerify(true);
verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), any(), any()); verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(AttributeScope.class), eq(ACTIVITY_STATE), any(), any());
long newTimeout = 1; long newTimeout = 1;
Thread.sleep(newTimeout); Thread.sleep(newTimeout);
@ -352,11 +353,11 @@ public class DefaultDeviceStateServiceTest {
long newTimeout = 1; long newTimeout = 1;
service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout); service.onDeviceInactivityTimeoutUpdate(tenantId, deviceId, newTimeout);
verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), any(), any()); verify(telemetrySubscriptionService, never()).saveAttrAndNotify(any(), eq(deviceId), any(AttributeScope.class), eq(ACTIVITY_STATE), any(), any());
} }
private void activityVerify(boolean isActive) { private void activityVerify(boolean isActive) {
verify(telemetrySubscriptionService, times(1)).saveAttrAndNotify(any(), eq(deviceId), any(), eq(ACTIVITY_STATE), eq(isActive), any()); verify(telemetrySubscriptionService, times(1)).saveAttrAndNotify(any(), eq(deviceId), any(AttributeScope.class), eq(ACTIVITY_STATE), eq(isActive), any());
} }
@Test @Test
@ -403,19 +404,19 @@ public class DefaultDeviceStateServiceTest {
assertThat(deviceState.isActive()).isEqualTo(true); assertThat(deviceState.isActive()).isEqualTo(true);
assertThat(deviceState.getLastActivityTime()).isEqualTo(lastReportedActivity); assertThat(deviceState.getLastActivityTime()).isEqualTo(lastReportedActivity);
then(telemetrySubscriptionService).should().saveAttrAndNotify( then(telemetrySubscriptionService).should().saveAttrAndNotify(
any(), eq(deviceId), any(), eq(LAST_ACTIVITY_TIME), eq(lastReportedActivity), any() any(), eq(deviceId), any(AttributeScope.class), eq(LAST_ACTIVITY_TIME), eq(lastReportedActivity), any()
); );
assertThat(deviceState.getLastInactivityAlarmTime()).isEqualTo(expectedInactivityAlarmTime); assertThat(deviceState.getLastInactivityAlarmTime()).isEqualTo(expectedInactivityAlarmTime);
if (shouldSetInactivityAlarmTimeToZero) { if (shouldSetInactivityAlarmTimeToZero) {
then(telemetrySubscriptionService).should().saveAttrAndNotify( then(telemetrySubscriptionService).should().saveAttrAndNotify(
any(), eq(deviceId), any(), eq(INACTIVITY_ALARM_TIME), eq(0L), any() any(), eq(deviceId), any(AttributeScope.class), eq(INACTIVITY_ALARM_TIME), eq(0L), any()
); );
} }
if (shouldUpdateActivityStateToActive) { if (shouldUpdateActivityStateToActive) {
then(telemetrySubscriptionService).should().saveAttrAndNotify( then(telemetrySubscriptionService).should().saveAttrAndNotify(
eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(SERVER_SCOPE), eq(ACTIVITY_STATE), eq(true), any() eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(AttributeScope.SERVER_SCOPE), eq(ACTIVITY_STATE), eq(true), any()
); );
var msgCaptor = ArgumentCaptor.forClass(TbMsg.class); var msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
@ -497,7 +498,7 @@ public class DefaultDeviceStateServiceTest {
assertThat(deviceState.isActive()).isEqualTo(expectedActivityState); assertThat(deviceState.isActive()).isEqualTo(expectedActivityState);
if (activityState && !expectedActivityState) { if (activityState && !expectedActivityState) {
then(telemetrySubscriptionService).should().saveAttrAndNotify( then(telemetrySubscriptionService).should().saveAttrAndNotify(
any(), eq(deviceId), any(), eq(ACTIVITY_STATE), eq(false), any() any(), eq(deviceId), any(AttributeScope.class), eq(ACTIVITY_STATE), eq(false), any()
); );
} }
} }
@ -594,7 +595,7 @@ public class DefaultDeviceStateServiceTest {
if (shouldUpdateActivityStateToInactive) { if (shouldUpdateActivityStateToInactive) {
then(telemetrySubscriptionService).should().saveAttrAndNotify( then(telemetrySubscriptionService).should().saveAttrAndNotify(
eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(SERVER_SCOPE), eq(ACTIVITY_STATE), eq(false), any() eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(AttributeScope.SERVER_SCOPE), eq(ACTIVITY_STATE), eq(false), any()
); );
var msgCaptor = ArgumentCaptor.forClass(TbMsg.class); var msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
@ -611,7 +612,7 @@ public class DefaultDeviceStateServiceTest {
assertThat(actualNotification.isActive()).isFalse(); assertThat(actualNotification.isActive()).isFalse();
then(telemetrySubscriptionService).should().saveAttrAndNotify( then(telemetrySubscriptionService).should().saveAttrAndNotify(
eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(SERVER_SCOPE), eq(TenantId.SYS_TENANT_ID), eq(deviceId), eq(AttributeScope.SERVER_SCOPE),
eq(INACTIVITY_ALARM_TIME), eq(expectedLastInactivityAlarmTime), any() eq(INACTIVITY_ALARM_TIME), eq(expectedLastInactivityAlarmTime), any()
); );
} }

22
common/dao-api/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java

@ -17,6 +17,7 @@ package org.thingsboard.server.dao.attributes;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
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;
@ -31,20 +32,41 @@ import java.util.Optional;
*/ */
public interface AttributesService { public interface AttributesService {
@Deprecated(since = "3.7.0")
ListenableFuture<Optional<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, String scope, String attributeKey);
ListenableFuture<Optional<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, String attributeKey); ListenableFuture<Optional<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, String attributeKey);
@Deprecated(since = "3.7.0")
ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, String scope, Collection<String> attributeKeys);
ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, Collection<String> attributeKeys); ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, Collection<String> attributeKeys);
@Deprecated(since = "3.7.0")
ListenableFuture<List<AttributeKvEntry>> findAll(TenantId tenantId, EntityId entityId, String scope);
ListenableFuture<List<AttributeKvEntry>> findAll(TenantId tenantId, EntityId entityId, AttributeScope scope); ListenableFuture<List<AttributeKvEntry>> findAll(TenantId tenantId, EntityId entityId, AttributeScope scope);
@Deprecated(since = "3.7.0")
ListenableFuture<List<String>> save(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes);
ListenableFuture<List<String>> save(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes); ListenableFuture<List<String>> save(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes);
@Deprecated(since = "3.7.0")
ListenableFuture<String> save(TenantId tenantId, EntityId entityId, String scope, AttributeKvEntry attribute);
ListenableFuture<String> save(TenantId tenantId, EntityId entityId, AttributeScope scope, AttributeKvEntry attribute); ListenableFuture<String> save(TenantId tenantId, EntityId entityId, AttributeScope scope, AttributeKvEntry attribute);
@Deprecated(since = "3.7.0")
ListenableFuture<List<String>> removeAll(TenantId tenantId, EntityId entityId, String scope, List<String> attributeKeys);
ListenableFuture<List<String>> removeAll(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> attributeKeys); ListenableFuture<List<String>> removeAll(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> attributeKeys);
List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); List<String> findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId);
@Deprecated(since = "3.7.0")
List<String> findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List<EntityId> entityIds);
List<String> findAllKeysByEntityIds(TenantId tenantId, List<EntityId> entityIds); List<String> findAllKeysByEntityIds(TenantId tenantId, List<EntityId> entityIds);
} }

6
dao/src/main/java/org/thingsboard/server/dao/attributes/AttributeUtils.java

@ -26,6 +26,12 @@ import java.util.List;
public class AttributeUtils { public class AttributeUtils {
@Deprecated(since = "3.7.0")
public static void validate(EntityId id, String scope) {
Validator.validateId(id.getId(), "Incorrect id " + id);
Validator.validateString(scope, "Incorrect scope " + scope);
}
public static void validate(EntityId id, AttributeScope scope) { public static void validate(EntityId id, AttributeScope scope) {
Validator.validateId(id.getId(), "Incorrect id " + id); Validator.validateId(id.getId(), "Incorrect id " + id);
Validator.checkNotNull(scope, "Incorrect scope " + scope); Validator.checkNotNull(scope, "Incorrect scope " + scope);

57
dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java

@ -23,6 +23,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.context.annotation.Primary; import org.springframework.context.annotation.Primary;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
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;
@ -53,6 +54,13 @@ public class BaseAttributesService implements AttributesService {
this.attributesDao = attributesDao; this.attributesDao = attributesDao;
} }
@Override
public ListenableFuture<Optional<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, String scope, String attributeKey) {
validate(entityId, scope);
Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey);
return Futures.immediateFuture(attributesDao.find(tenantId, entityId, AttributeScope.valueOf(scope), attributeKey));
}
@Override @Override
public ListenableFuture<Optional<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, String attributeKey) { public ListenableFuture<Optional<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, String attributeKey) {
validate(entityId, scope); validate(entityId, scope);
@ -60,17 +68,30 @@ public class BaseAttributesService implements AttributesService {
return Futures.immediateFuture(attributesDao.find(tenantId, entityId, scope, attributeKey)); return Futures.immediateFuture(attributesDao.find(tenantId, entityId, scope, attributeKey));
} }
@Override
public ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, String scope, Collection<String> attributeKeys) {
validate(entityId, scope);
attributeKeys.forEach(attributeKey -> Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey));
return Futures.immediateFuture(attributesDao.find(tenantId, entityId, AttributeScope.valueOf(scope), attributeKeys));
}
@Override @Override
public ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, Collection<String> attributeKeys) { public ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, Collection<String> attributeKeys) {
validate(entityId, scope); validate(entityId, scope);
attributeKeys.forEach(attributeKey -> Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey)); attributeKeys.forEach(attributeKey -> Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey));
return Futures.immediateFuture(attributesDao.find(tenantId, entityId, scope, attributeKeys)); return Futures.immediateFuture(attributesDao.find(tenantId, entityId, scope, attributeKeys));
}
@Override
public ListenableFuture<List<AttributeKvEntry>> findAll(TenantId tenantId, EntityId entityId, String scope) {
validate(entityId, scope);
return Futures.immediateFuture(attributesDao.findAll(tenantId, entityId, AttributeScope.valueOf(scope)));
} }
@Override @Override
public ListenableFuture<List<AttributeKvEntry>> findAll(TenantId tenantId, EntityId entityId, AttributeScope scope) { public ListenableFuture<List<AttributeKvEntry>> findAll(TenantId tenantId, EntityId entityId, AttributeScope scope) {
validate(entityId, scope); validate(entityId, scope);
return Futures.immediateFuture(attributesDao.findAll(tenantId, entityId, scope)); return Futures.immediateFuture(attributesDao.findAll(tenantId, entityId, scope));
} }
@Override @Override
@ -78,29 +99,55 @@ public class BaseAttributesService implements AttributesService {
return attributesDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId); return attributesDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId);
} }
@Override
public List<String> findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List<EntityId> entityIds) {
return attributesDao.findAllKeysByEntityIds(tenantId, entityIds);
}
@Override @Override
public List<String> findAllKeysByEntityIds(TenantId tenantId, List<EntityId> entityIds) { public List<String> findAllKeysByEntityIds(TenantId tenantId, List<EntityId> entityIds) {
return attributesDao.findAllKeysByEntityIds(tenantId, entityIds); return attributesDao.findAllKeysByEntityIds(tenantId, entityIds);
} }
@Override
public ListenableFuture<String> save(TenantId tenantId, EntityId entityId, String scope, AttributeKvEntry attribute) {
validate(entityId, scope);
AttributeUtils.validate(attribute, valueNoXssValidation);
return attributesDao.save(tenantId, entityId, AttributeScope.valueOf(scope), attribute);
}
@Override @Override
public ListenableFuture<String> save(TenantId tenantId, EntityId entityId, AttributeScope scope, AttributeKvEntry attribute) { public ListenableFuture<String> save(TenantId tenantId, EntityId entityId, AttributeScope scope, AttributeKvEntry attribute) {
validate(entityId, scope); validate(entityId, scope);
AttributeUtils.validate(attribute, valueNoXssValidation); AttributeUtils.validate(attribute, valueNoXssValidation);
return attributesDao.save(tenantId, entityId, scope, attribute); return attributesDao.save(tenantId, entityId, scope, attribute);
}
@Override
public ListenableFuture<List<String>> save(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes) {
validate(entityId, scope);
AttributeUtils.validate(attributes, valueNoXssValidation);
List<ListenableFuture<String>> saveFutures = attributes.stream().map(attribute -> attributesDao.save(tenantId, entityId, AttributeScope.valueOf(scope), attribute)).collect(Collectors.toList());
return Futures.allAsList(saveFutures);
} }
@Override @Override
public ListenableFuture<List<String>> save(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes) { public ListenableFuture<List<String>> save(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes) {
validate(entityId, scope); validate(entityId, scope);
AttributeUtils.validate(attributes, valueNoXssValidation); AttributeUtils.validate(attributes, valueNoXssValidation);
List<ListenableFuture<String>> saveFutures = attributes.stream().map(attribute -> attributesDao.save(tenantId, entityId, scope, attribute)).collect(Collectors.toList()); List<ListenableFuture<String>> saveFutures = attributes.stream().map(attribute -> attributesDao.save(tenantId, entityId, scope, attribute)).collect(Collectors.toList());
return Futures.allAsList(saveFutures); return Futures.allAsList(saveFutures);
} }
@Override
public ListenableFuture<List<String>> removeAll(TenantId tenantId, EntityId entityId, String scope, List<String> attributeKeys) {
validate(entityId, scope);
return Futures.allAsList(attributesDao.removeAll(tenantId, entityId, AttributeScope.valueOf(scope), attributeKeys));
}
@Override @Override
public ListenableFuture<List<String>> removeAll(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> attributeKeys) { public ListenableFuture<List<String>> removeAll(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> attributeKeys) {
validate(entityId, scope); validate(entityId, scope);
return Futures.allAsList(attributesDao.removeAll(tenantId, entityId, scope, attributeKeys)); return Futures.allAsList(attributesDao.removeAll(tenantId, entityId, scope, attributeKeys));
} }
} }

36
dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java

@ -27,6 +27,7 @@ import org.springframework.stereotype.Service;
import org.thingsboard.server.cache.TbCacheValueWrapper; import org.thingsboard.server.cache.TbCacheValueWrapper;
import org.thingsboard.server.cache.TbTransactionalCache; import org.thingsboard.server.cache.TbTransactionalCache;
import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
@ -108,6 +109,11 @@ public class CachedAttributesService implements AttributesService {
} }
@Override
public ListenableFuture<Optional<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, String scope, String attributeKey) {
return find(tenantId, entityId, AttributeScope.valueOf(scope), attributeKey);
}
@Override @Override
public ListenableFuture<Optional<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, String attributeKey) { public ListenableFuture<Optional<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, String attributeKey) {
validate(entityId, scope); validate(entityId, scope);
@ -137,6 +143,11 @@ public class CachedAttributesService implements AttributesService {
} }
} }
@Override
public ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, String scope, Collection<String> attributeKeys) {
return find(tenantId, entityId, AttributeScope.valueOf(scope), attributeKeys);
}
@Override @Override
public ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, final Collection<String> attributeKeysNonUnique) { public ListenableFuture<List<AttributeKvEntry>> find(TenantId tenantId, EntityId entityId, AttributeScope scope, final Collection<String> attributeKeysNonUnique) {
validate(entityId, scope); validate(entityId, scope);
@ -204,6 +215,11 @@ public class CachedAttributesService implements AttributesService {
return cachedAttributes; return cachedAttributes;
} }
@Override
public ListenableFuture<List<AttributeKvEntry>> findAll(TenantId tenantId, EntityId entityId, String scope) {
return findAll(tenantId, entityId, AttributeScope.valueOf(scope));
}
@Override @Override
public ListenableFuture<List<AttributeKvEntry>> findAll(TenantId tenantId, EntityId entityId, AttributeScope scope) { public ListenableFuture<List<AttributeKvEntry>> findAll(TenantId tenantId, EntityId entityId, AttributeScope scope) {
validate(entityId, scope); validate(entityId, scope);
@ -215,11 +231,21 @@ public class CachedAttributesService implements AttributesService {
return attributesDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId); return attributesDao.findAllKeysByDeviceProfileId(tenantId, deviceProfileId);
} }
@Override
public List<String> findAllKeysByEntityIds(TenantId tenantId, EntityType entityType, List<EntityId> entityIds) {
return findAllKeysByEntityIds(tenantId, entityIds);
}
@Override @Override
public List<String> findAllKeysByEntityIds(TenantId tenantId, List<EntityId> entityIds) { public List<String> findAllKeysByEntityIds(TenantId tenantId, List<EntityId> entityIds) {
return attributesDao.findAllKeysByEntityIds(tenantId, entityIds); return attributesDao.findAllKeysByEntityIds(tenantId, entityIds);
} }
@Override
public ListenableFuture<String> save(TenantId tenantId, EntityId entityId, String scope, AttributeKvEntry attribute) {
return save(tenantId, entityId, AttributeScope.valueOf(scope), attribute);
}
@Override @Override
public ListenableFuture<String> save(TenantId tenantId, EntityId entityId, AttributeScope scope, AttributeKvEntry attribute) { public ListenableFuture<String> save(TenantId tenantId, EntityId entityId, AttributeScope scope, AttributeKvEntry attribute) {
validate(entityId, scope); validate(entityId, scope);
@ -228,6 +254,11 @@ public class CachedAttributesService implements AttributesService {
return Futures.transform(future, key -> evict(entityId, scope, attribute, key), cacheExecutor); return Futures.transform(future, key -> evict(entityId, scope, attribute, key), cacheExecutor);
} }
@Override
public ListenableFuture<List<String>> save(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes) {
return save(tenantId, entityId, scope, attributes);
}
@Override @Override
public ListenableFuture<List<String>> save(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes) { public ListenableFuture<List<String>> save(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes) {
validate(entityId, scope); validate(entityId, scope);
@ -249,6 +280,11 @@ public class CachedAttributesService implements AttributesService {
return key; return key;
} }
@Override
public ListenableFuture<List<String>> removeAll(TenantId tenantId, EntityId entityId, String scope, List<String> attributeKeys) {
return removeAll(tenantId, entityId, AttributeScope.valueOf(scope), attributeKeys);
}
@Override @Override
public ListenableFuture<List<String>> removeAll(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> attributeKeys) { public ListenableFuture<List<String>> removeAll(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> attributeKeys) {
validate(entityId, scope); validate(entityId, scope);

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

@ -551,7 +551,7 @@ CREATE TABLE IF NOT EXISTS key_dictionary
( (
key varchar(255) NOT NULL, key varchar(255) NOT NULL,
key_id serial UNIQUE, key_id serial UNIQUE,
CONSTRAINT key_id_pkey PRIMARY KEY (key) CONSTRAINT key_dictionary_id_pkey PRIMARY KEY (key)
); );
CREATE TABLE IF NOT EXISTS oauth2_params ( CREATE TABLE IF NOT EXISTS oauth2_params (

2
dao/src/main/resources/sql/schema-timescale.sql

@ -31,7 +31,7 @@ CREATE TABLE IF NOT EXISTS ts_kv (
CREATE TABLE IF NOT EXISTS key_dictionary ( CREATE TABLE IF NOT EXISTS key_dictionary (
key varchar(255) NOT NULL, key varchar(255) NOT NULL,
key_id serial UNIQUE, key_id serial UNIQUE,
CONSTRAINT key_id_pkey PRIMARY KEY (key) CONSTRAINT key_dictionary_id_pkey PRIMARY KEY (key)
); );
CREATE TABLE IF NOT EXISTS ts_kv_latest ( CREATE TABLE IF NOT EXISTS ts_kv_latest (

2
dao/src/main/resources/sql/schema-ts-psql.sql

@ -31,7 +31,7 @@ CREATE TABLE IF NOT EXISTS key_dictionary
( (
key varchar(255) NOT NULL, key varchar(255) NOT NULL,
key_id serial UNIQUE, key_id serial UNIQUE,
CONSTRAINT key_id_pkey PRIMARY KEY (key) CONSTRAINT key_dictionary_id_pkey PRIMARY KEY (key)
); );
CREATE OR REPLACE PROCEDURE drop_partitions_by_system_ttl(IN partition_type varchar, IN system_ttl bigint, INOUT deleted bigint) CREATE OR REPLACE PROCEDURE drop_partitions_by_system_ttl(IN partition_type varchar, IN system_ttl bigint, INOUT deleted bigint)

37
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineTelemetryService.java

@ -17,6 +17,7 @@ package org.thingsboard.rule.engine.api;
import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
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;
@ -40,32 +41,68 @@ public interface RuleEngineTelemetryService {
void saveWithoutLatestAndNotify(TenantId tenantId, CustomerId id, EntityId entityId, List<TsKvEntry> ts, long ttl, FutureCallback<Void> callback); void saveWithoutLatestAndNotify(TenantId tenantId, CustomerId id, EntityId entityId, List<TsKvEntry> ts, long ttl, FutureCallback<Void> callback);
@Deprecated(since = "3.7.0")
void saveAndNotify(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, FutureCallback<Void> callback); void saveAndNotify(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, FutureCallback<Void> callback);
void saveAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes, FutureCallback<Void> callback);
@Deprecated(since = "3.7.0")
void saveAndNotify(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback); void saveAndNotify(TenantId tenantId, EntityId entityId, String scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback);
void saveAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, List<AttributeKvEntry> attributes, boolean notifyDevice, FutureCallback<Void> callback);
void saveLatestAndNotify(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, FutureCallback<Void> callback); void saveLatestAndNotify(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts, FutureCallback<Void> callback);
@Deprecated(since = "3.7.0")
ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, long value); ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, long value);
ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, long value);
@Deprecated(since = "3.7.0")
ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, String value); ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, String value);
ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, String value);
@Deprecated(since = "3.7.0")
ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, double value); ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, double value);
ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, double value);
@Deprecated(since = "3.7.0")
ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, boolean value); ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, boolean value);
ListenableFuture<Void> saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, boolean value);
@Deprecated(since = "3.7.0")
void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, long value, FutureCallback<Void> callback); void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, long value, FutureCallback<Void> callback);
void saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, long value, FutureCallback<Void> callback);
@Deprecated(since = "3.7.0")
void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, String value, FutureCallback<Void> callback); void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, String value, FutureCallback<Void> callback);
void saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, String value, FutureCallback<Void> callback);
@Deprecated(since = "3.7.0")
void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, double value, FutureCallback<Void> callback); void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, double value, FutureCallback<Void> callback);
void saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, double value, FutureCallback<Void> callback);
@Deprecated(since = "3.7.0")
void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, boolean value, FutureCallback<Void> callback); void saveAttrAndNotify(TenantId tenantId, EntityId entityId, String scope, String key, boolean value, FutureCallback<Void> callback);
void saveAttrAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, String key, boolean value, FutureCallback<Void> callback);
@Deprecated(since = "3.7.0")
void deleteAndNotify(TenantId tenantId, EntityId entityId, String scope, List<String> keys, FutureCallback<Void> callback); void deleteAndNotify(TenantId tenantId, EntityId entityId, String scope, List<String> keys, FutureCallback<Void> callback);
void deleteAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> keys, FutureCallback<Void> callback);
@Deprecated(since = "3.7.0")
void deleteAndNotify(TenantId tenantId, EntityId entityId, String scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback); void deleteAndNotify(TenantId tenantId, EntityId entityId, String scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback);
void deleteAndNotify(TenantId tenantId, EntityId entityId, AttributeScope scope, List<String> keys, boolean notifyDevice, FutureCallback<Void> callback);
void deleteLatest(TenantId tenantId, EntityId entityId, List<String> keys, FutureCallback<Void> callback); void deleteLatest(TenantId tenantId, EntityId entityId, List<String> keys, FutureCallback<Void> callback);
void deleteAllLatest(TenantId tenantId, EntityId entityId, FutureCallback<Collection<String>> callback); void deleteAllLatest(TenantId tenantId, EntityId entityId, FutureCallback<Collection<String>> callback);

13
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbCopyAttributesToEntityViewNode.java

@ -29,6 +29,7 @@ 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.AttributeScope;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.AttributeKvEntry;
@ -79,8 +80,8 @@ public class TbCopyAttributesToEntityViewNode implements TbNode {
ACTIVITY_EVENT, INACTIVITY_EVENT, POST_ATTRIBUTES_REQUEST)) { ACTIVITY_EVENT, INACTIVITY_EVENT, POST_ATTRIBUTES_REQUEST)) {
if (!msg.getMetaData().getData().isEmpty()) { if (!msg.getMetaData().getData().isEmpty()) {
long now = System.currentTimeMillis(); long now = System.currentTimeMillis();
String scope = msg.isTypeOf(POST_ATTRIBUTES_REQUEST) ? AttributeScope scope = msg.isTypeOf(POST_ATTRIBUTES_REQUEST) ?
DataConstants.CLIENT_SCOPE : msg.getMetaData().getValue(DataConstants.SCOPE); AttributeScope.CLIENT_SCOPE : AttributeScope.valueOf(msg.getMetaData().getValue(DataConstants.SCOPE));
ListenableFuture<List<EntityView>> entityViewsFuture = ListenableFuture<List<EntityView>> entityViewsFuture =
ctx.getEntityViewService().findEntityViewsByTenantIdAndEntityIdAsync(ctx.getTenantId(), msg.getOriginator()); ctx.getEntityViewService().findEntityViewsByTenantIdAndEntityIdAsync(ctx.getTenantId(), msg.getOriginator());
@ -145,17 +146,17 @@ public class TbCopyAttributesToEntityViewNode implements TbNode {
ctx.enqueueForTellNext(ctx.newMsg(msg.getQueueName(), msg.getType(), entityView.getId(), msg.getCustomerId(), msg.getMetaData(), msg.getData()), SUCCESS); ctx.enqueueForTellNext(ctx.newMsg(msg.getQueueName(), msg.getType(), entityView.getId(), msg.getCustomerId(), msg.getMetaData(), msg.getData()), SUCCESS);
} }
private boolean attributeContainsInEntityView(String scope, String attrKey, EntityView entityView) { private boolean attributeContainsInEntityView(AttributeScope scope, String attrKey, EntityView entityView) {
AttributesEntityView attributesEntityView = entityView.getKeys().getAttributes(); AttributesEntityView attributesEntityView = entityView.getKeys().getAttributes();
List<String> keys = null; List<String> keys = null;
switch (scope) { switch (scope) {
case DataConstants.CLIENT_SCOPE: case CLIENT_SCOPE:
keys = attributesEntityView.getCs(); keys = attributesEntityView.getCs();
break; break;
case DataConstants.SERVER_SCOPE: case SERVER_SCOPE:
keys = attributesEntityView.getSs(); keys = attributesEntityView.getSs();
break; break;
case DataConstants.SHARED_SCOPE: case SHARED_SCOPE:
keys = attributesEntityView.getSh(); keys = attributesEntityView.getSh();
break; break;
} }

5
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/math/TbMathNode.java

@ -34,7 +34,6 @@ 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.AttributeScope; import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
@ -213,11 +212,11 @@ public class TbMathNode implements TbNode {
if (isIntegerResult(mathResultDef, config.getOperation())) { if (isIntegerResult(mathResultDef, config.getOperation())) {
var value = toIntValue(result); var value = toIntValue(result);
return ctx.getTelemetryService().saveAttrAndNotify( return ctx.getTelemetryService().saveAttrAndNotify(
ctx.getTenantId(), msg.getOriginator(), attributeScope.name(), mathResultDef.getKey(), value); ctx.getTenantId(), msg.getOriginator(), attributeScope, mathResultDef.getKey(), value);
} else { } else {
var value = toDoubleValue(mathResultDef, result); var value = toDoubleValue(mathResultDef, result);
return ctx.getTelemetryService().saveAttrAndNotify( return ctx.getTelemetryService().saveAttrAndNotify(
ctx.getTenantId(), msg.getOriginator(), attributeScope.name(), mathResultDef.getKey(), value); ctx.getTenantId(), msg.getOriginator(), attributeScope, mathResultDef.getKey(), value);
} }
} }

18
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java

@ -90,7 +90,7 @@ public class TbMsgAttributesNode implements TbNode {
ctx.tellSuccess(msg); ctx.tellSuccess(msg);
return; return;
} }
String scope = getScope(msg.getMetaData().getValue(SCOPE)); AttributeScope scope = getScope(msg.getMetaData().getValue(SCOPE));
boolean sendAttributesUpdateNotification = checkSendNotification(scope); boolean sendAttributesUpdateNotification = checkSendNotification(scope);
if (!config.isUpdateAttributesOnlyOnValueChange()) { if (!config.isUpdateAttributesOnlyOnValueChange()) {
@ -99,7 +99,7 @@ public class TbMsgAttributesNode implements TbNode {
} }
List<String> keys = newAttributes.stream().map(KvEntry::getKey).collect(Collectors.toList()); List<String> keys = newAttributes.stream().map(KvEntry::getKey).collect(Collectors.toList());
ListenableFuture<List<AttributeKvEntry>> findFuture = ctx.getAttributesService().find(ctx.getTenantId(), msg.getOriginator(), AttributeScope.valueOf(scope), keys); ListenableFuture<List<AttributeKvEntry>> findFuture = ctx.getAttributesService().find(ctx.getTenantId(), msg.getOriginator(), scope, keys);
DonAsynchron.withCallback(findFuture, DonAsynchron.withCallback(findFuture,
currentAttributes -> { currentAttributes -> {
@ -110,7 +110,7 @@ public class TbMsgAttributesNode implements TbNode {
MoreExecutors.directExecutor()); MoreExecutors.directExecutor());
} }
void saveAttr(List<AttributeKvEntry> attributes, TbContext ctx, TbMsg msg, String scope, boolean sendAttributesUpdateNotification) { void saveAttr(List<AttributeKvEntry> attributes, TbContext ctx, TbMsg msg, AttributeScope scope, boolean sendAttributesUpdateNotification) {
if (attributes.isEmpty()) { if (attributes.isEmpty()) {
ctx.tellSuccess(msg); ctx.tellSuccess(msg);
return; return;
@ -122,7 +122,7 @@ public class TbMsgAttributesNode implements TbNode {
attributes, attributes,
config.isNotifyDevice() || checkNotifyDeviceMdValue(msg.getMetaData().getValue(NOTIFY_DEVICE_METADATA_KEY)), config.isNotifyDevice() || checkNotifyDeviceMdValue(msg.getMetaData().getValue(NOTIFY_DEVICE_METADATA_KEY)),
sendAttributesUpdateNotification ? sendAttributesUpdateNotification ?
new AttributesUpdateNodeCallback(ctx, msg, scope, attributes) : new AttributesUpdateNodeCallback(ctx, msg, scope.name(), attributes) :
new TelemetryNodeCallback(ctx, msg) new TelemetryNodeCallback(ctx, msg)
); );
} }
@ -145,8 +145,8 @@ public class TbMsgAttributesNode implements TbNode {
.collect(Collectors.toList()); .collect(Collectors.toList());
} }
private boolean checkSendNotification(String scope) { private boolean checkSendNotification(AttributeScope scope) {
return config.isSendAttributesUpdatedNotification() && !CLIENT_SCOPE.equals(scope); return config.isSendAttributesUpdatedNotification() && AttributeScope.CLIENT_SCOPE != scope;
} }
private boolean checkNotifyDeviceMdValue(String notifyDeviceMdValue) { private boolean checkNotifyDeviceMdValue(String notifyDeviceMdValue) {
@ -154,11 +154,11 @@ public class TbMsgAttributesNode implements TbNode {
return StringUtils.isEmpty(notifyDeviceMdValue) || Boolean.parseBoolean(notifyDeviceMdValue); return StringUtils.isEmpty(notifyDeviceMdValue) || Boolean.parseBoolean(notifyDeviceMdValue);
} }
private String getScope(String mdScopeValue) { private AttributeScope getScope(String mdScopeValue) {
if (StringUtils.isNotEmpty(mdScopeValue)) { if (StringUtils.isNotEmpty(mdScopeValue)) {
return mdScopeValue; return AttributeScope.valueOf(mdScopeValue);
} }
return config.getScope(); return AttributeScope.valueOf(config.getScope());
} }
@Override @Override

16
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgDeleteAttributesNode.java

@ -22,6 +22,7 @@ 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.AttributeScope;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
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;
@ -32,7 +33,6 @@ import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.DataConstants.NOTIFY_DEVICE_METADATA_KEY; import static org.thingsboard.server.common.data.DataConstants.NOTIFY_DEVICE_METADATA_KEY;
import static org.thingsboard.server.common.data.DataConstants.SCOPE; import static org.thingsboard.server.common.data.DataConstants.SCOPE;
import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE;
@Slf4j @Slf4j
@RuleNode( @RuleNode(
@ -69,7 +69,7 @@ public class TbMsgDeleteAttributesNode implements TbNode {
if (keysToDelete.isEmpty()) { if (keysToDelete.isEmpty()) {
ctx.tellSuccess(msg); ctx.tellSuccess(msg);
} else { } else {
String scope = getScope(msg.getMetaData().getValue(SCOPE)); AttributeScope scope = getScope(msg.getMetaData().getValue(SCOPE));
ctx.getTelemetryService().deleteAndNotify( ctx.getTelemetryService().deleteAndNotify(
ctx.getTenantId(), ctx.getTenantId(),
msg.getOriginator(), msg.getOriginator(),
@ -77,21 +77,21 @@ public class TbMsgDeleteAttributesNode implements TbNode {
keysToDelete, keysToDelete,
checkNotifyDevice(msg.getMetaData().getValue(NOTIFY_DEVICE_METADATA_KEY), scope), checkNotifyDevice(msg.getMetaData().getValue(NOTIFY_DEVICE_METADATA_KEY), scope),
config.isSendAttributesDeletedNotification() ? config.isSendAttributesDeletedNotification() ?
new AttributesDeleteNodeCallback(ctx, msg, scope, keysToDelete) : new AttributesDeleteNodeCallback(ctx, msg, scope.name(), keysToDelete) :
new TelemetryNodeCallback(ctx, msg) new TelemetryNodeCallback(ctx, msg)
); );
} }
} }
private String getScope(String mdScopeValue) { private AttributeScope getScope(String mdScopeValue) {
if (StringUtils.isNotEmpty(mdScopeValue)) { if (StringUtils.isNotEmpty(mdScopeValue)) {
return mdScopeValue; return AttributeScope.valueOf(mdScopeValue);
} }
return config.getScope(); return AttributeScope.valueOf(config.getScope());
} }
private boolean checkNotifyDevice(String notifyDeviceMdValue, String scope) { private boolean checkNotifyDevice(String notifyDeviceMdValue, AttributeScope scope) {
return SHARED_SCOPE.equals(scope) && (config.isNotifyDevice() || Boolean.parseBoolean(notifyDeviceMdValue)); return (AttributeScope.SHARED_SCOPE == scope) && (config.isNotifyDevice() || Boolean.parseBoolean(notifyDeviceMdValue));
} }
} }

4
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/math/TbMathNodeTest.java

@ -438,14 +438,14 @@ public class TbMathNodeTest {
TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originator, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().put("a", 5).toString()); TbMsg msg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, originator, TbMsgMetaData.EMPTY, JacksonUtil.newObjectNode().put("a", 5).toString());
when(telemetryService.saveAttrAndNotify(any(), any(), anyString(), anyString(), anyDouble())) when(telemetryService.saveAttrAndNotify(any(), any(), any(AttributeScope.class), anyString(), anyDouble()))
.thenReturn(Futures.immediateFuture(null)); .thenReturn(Futures.immediateFuture(null));
node.onMsg(ctx, msg); node.onMsg(ctx, msg);
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class); ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture()); verify(ctx, timeout(TIMEOUT)).tellSuccess(msgCaptor.capture());
verify(telemetryService, times(1)).saveAttrAndNotify(any(), any(), anyString(), anyString(), anyDouble()); verify(telemetryService, times(1)).saveAttrAndNotify(any(), any(), any(AttributeScope.class), anyString(), anyDouble());
TbMsg resultMsg = msgCaptor.getValue(); TbMsg resultMsg = msgCaptor.getValue();
assertNotNull(resultMsg); assertNotNull(resultMsg);

3
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/profile/DeviceStateTest.java

@ -22,6 +22,7 @@ import org.mockito.ArgumentCaptor;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.rule.engine.api.RuleEngineAlarmService; import org.thingsboard.rule.engine.api.RuleEngineAlarmService;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.alarm.AlarmApiCallResult; import org.thingsboard.server.common.data.alarm.AlarmApiCallResult;
@ -75,7 +76,7 @@ public class DeviceStateTest {
when(ctx.getDeviceService()).thenReturn(mock(DeviceService.class)); when(ctx.getDeviceService()).thenReturn(mock(DeviceService.class));
AttributesService attributesService = mock(AttributesService.class); AttributesService attributesService = mock(AttributesService.class);
when(attributesService.find(any(), any(), any(), anyCollection())).thenReturn(Futures.immediateFuture(Collections.emptyList())); when(attributesService.find(any(), any(), any(AttributeScope.class), anyCollection())).thenReturn(Futures.immediateFuture(Collections.emptyList()));
when(ctx.getAttributesService()).thenReturn(attributesService); when(ctx.getAttributesService()).thenReturn(attributesService);
RuleEngineAlarmService alarmService = mock(RuleEngineAlarmService.class); RuleEngineAlarmService alarmService = mock(RuleEngineAlarmService.class);

9
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNodeTest.java

@ -29,7 +29,7 @@ import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
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.server.common.data.DataConstants; import org.thingsboard.server.common.data.AttributeScope;
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.id.TenantId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.AttributeKvEntry;
@ -54,7 +54,6 @@ import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.BDDMockito.willCallRealMethod; import static org.mockito.BDDMockito.willCallRealMethod;
import static org.mockito.Mockito.mock; import static org.mockito.Mockito.mock;
@ -157,7 +156,7 @@ class TbMsgAttributesNodeTest {
when(ctxMock.getTenantId()).thenReturn(tenantId); when(ctxMock.getTenantId()).thenReturn(tenantId);
when(ctxMock.getTelemetryService()).thenReturn(telemetryServiceMock); when(ctxMock.getTelemetryService()).thenReturn(telemetryServiceMock);
willCallRealMethod().given(node).init(any(TbContext.class), any(TbNodeConfiguration.class)); willCallRealMethod().given(node).init(any(TbContext.class), any(TbNodeConfiguration.class));
willCallRealMethod().given(node).saveAttr(any(), eq(ctxMock), any(TbMsg.class), anyString(), anyBoolean()); willCallRealMethod().given(node).saveAttr(any(), eq(ctxMock), any(TbMsg.class), any(AttributeScope.class), anyBoolean());
node.init(ctxMock, tbNodeConfiguration); node.init(ctxMock, tbNodeConfiguration);
@ -169,12 +168,12 @@ class TbMsgAttributesNodeTest {
var testTbMsg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, deviceId, md, TbMsg.EMPTY_STRING); var testTbMsg = TbMsg.newMsg(TbMsgType.POST_TELEMETRY_REQUEST, deviceId, md, TbMsg.EMPTY_STRING);
List<AttributeKvEntry> testAttrList = List.of(new BaseAttributeKvEntry(0L, new StringDataEntry("testKey", "testValue"))); List<AttributeKvEntry> testAttrList = List.of(new BaseAttributeKvEntry(0L, new StringDataEntry("testKey", "testValue")));
node.saveAttr(testAttrList, ctxMock, testTbMsg, DataConstants.SHARED_SCOPE, false); node.saveAttr(testAttrList, ctxMock, testTbMsg, AttributeScope.SHARED_SCOPE, false);
ArgumentCaptor<Boolean> notifyDeviceCaptor = ArgumentCaptor.forClass(Boolean.class); ArgumentCaptor<Boolean> notifyDeviceCaptor = ArgumentCaptor.forClass(Boolean.class);
verify(telemetryServiceMock, times(1)).saveAndNotify( verify(telemetryServiceMock, times(1)).saveAndNotify(
eq(tenantId), eq(deviceId), eq(DataConstants.SHARED_SCOPE), eq(tenantId), eq(deviceId), eq(AttributeScope.SHARED_SCOPE),
eq(testAttrList), notifyDeviceCaptor.capture(), any() eq(testAttrList), notifyDeviceCaptor.capture(), any()
); );
boolean notifyDevice = notifyDeviceCaptor.getValue(); boolean notifyDevice = notifyDeviceCaptor.getValue();

3
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/telemetry/TbMsgDeleteAttributesNodeTest.java

@ -25,6 +25,7 @@ import org.thingsboard.rule.engine.api.RuleEngineTelemetryService;
import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbContext;
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.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.msg.TbMsgType; import org.thingsboard.server.common.data.msg.TbMsgType;
@ -81,7 +82,7 @@ public class TbMsgDeleteAttributesNodeTest {
callBack.onSuccess(null); callBack.onSuccess(null);
return null; return null;
}).given(telemetryService).deleteAndNotify( }).given(telemetryService).deleteAndNotify(
any(), any(), anyString(), anyList(), anyBoolean(), any()); any(), any(), any(AttributeScope.class), anyList(), anyBoolean(), any());
} }
@AfterEach @AfterEach

Loading…
Cancel
Save