diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml
index f8f71e0813..42191e57a9 100644
--- a/application/src/main/resources/thingsboard.yml
+++ b/application/src/main/resources/thingsboard.yml
@@ -143,6 +143,9 @@ cassandra:
concurrent_limit: "${CASSANDRA_QUERY_CONCURRENT_LIMIT:1000}"
permit_max_wait_time: "${PERMIT_MAX_WAIT_TIME:120000}"
rate_limit_print_interval_ms: "${CASSANDRA_QUERY_RATE_LIMIT_PRINT_MS:10000}"
+ tenant_rate_limits:
+ enabled: "${CASSANDRA_QUERY_TENANT_RATE_LIMITS_ENABLED:false}"
+ configuration: "${CASSANDRA_QUERY_TENANT_RATE_LIMITS_VALUE:1000:1,30000:60}"
# SQL configuration parameters
sql:
diff --git a/common/message/pom.xml b/common/message/pom.xml
index d914d4f21d..228cd1cbbe 100644
--- a/common/message/pom.xml
+++ b/common/message/pom.xml
@@ -60,6 +60,10 @@
ch.qos.logback
logback-classic
+
+ com.github.vladimir-bukhtoyarov
+ bucket4j-core
+
com.google.protobuf
protobuf-java
diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TbTransportRateLimits.java b/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java
similarity index 89%
rename from common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TbTransportRateLimits.java
rename to common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java
index d598734a5a..de6d1b8871 100644
--- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/TbTransportRateLimits.java
+++ b/common/message/src/main/java/org/thingsboard/server/common/msg/tools/TbRateLimits.java
@@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.thingsboard.server.common.transport.service;
+package org.thingsboard.server.common.msg.tools;
import io.github.bucket4j.Bandwidth;
import io.github.bucket4j.Bucket4j;
@@ -25,10 +25,10 @@ import java.time.Duration;
/**
* Created by ashvayka on 22.10.18.
*/
-class TbTransportRateLimits {
+public class TbRateLimits {
private final LocalBucket bucket;
- public TbTransportRateLimits(String limitsConfiguration) {
+ public TbRateLimits(String limitsConfiguration) {
LocalBucketBuilder builder = Bucket4j.builder();
boolean initialized = false;
for (String limitSrc : limitsConfiguration.split(",")) {
@@ -46,7 +46,7 @@ class TbTransportRateLimits {
}
- boolean tryConsume() {
+ public boolean tryConsume() {
return bucket.tryConsume(1);
}
diff --git a/common/transport/transport-api/pom.xml b/common/transport/transport-api/pom.xml
index 4ed6ed784d..3538e468d4 100644
--- a/common/transport/transport-api/pom.xml
+++ b/common/transport/transport-api/pom.xml
@@ -99,10 +99,6 @@
com.google.protobuf
protobuf-java
-
- com.github.vladimir-bukhtoyarov
- bucket4j-core
-
diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/AbstractTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/AbstractTransportService.java
index 4af594c421..c9f681fd3f 100644
--- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/AbstractTransportService.java
+++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/AbstractTransportService.java
@@ -20,6 +20,7 @@ import org.springframework.beans.factory.annotation.Value;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
+import org.thingsboard.server.common.msg.tools.TbRateLimits;
import org.thingsboard.server.common.transport.SessionMsgListener;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
@@ -58,8 +59,8 @@ public abstract class AbstractTransportService implements TransportService {
private ConcurrentMap sessions = new ConcurrentHashMap<>();
//TODO: Implement cleanup of this maps.
- private ConcurrentMap perTenantLimits = new ConcurrentHashMap<>();
- private ConcurrentMap perDeviceLimits = new ConcurrentHashMap<>();
+ private ConcurrentMap perTenantLimits = new ConcurrentHashMap<>();
+ private ConcurrentMap perDeviceLimits = new ConcurrentHashMap<>();
@Override
public void registerAsyncSession(TransportProtos.SessionInfoProto sessionInfo, SessionMsgListener listener) {
@@ -204,7 +205,7 @@ public abstract class AbstractTransportService implements TransportService {
return true;
}
TenantId tenantId = new TenantId(new UUID(sessionInfo.getTenantIdMSB(), sessionInfo.getTenantIdLSB()));
- TbTransportRateLimits rateLimits = perTenantLimits.computeIfAbsent(tenantId, id -> new TbTransportRateLimits(perTenantLimitsConf));
+ TbRateLimits rateLimits = perTenantLimits.computeIfAbsent(tenantId, id -> new TbRateLimits(perTenantLimitsConf));
if (!rateLimits.tryConsume()) {
if (callback != null) {
callback.onError(new TbRateLimitsException(EntityType.TENANT));
@@ -215,7 +216,7 @@ public abstract class AbstractTransportService implements TransportService {
return false;
}
DeviceId deviceId = new DeviceId(new UUID(sessionInfo.getDeviceIdMSB(), sessionInfo.getDeviceIdLSB()));
- rateLimits = perDeviceLimits.computeIfAbsent(deviceId, id -> new TbTransportRateLimits(perDevicesLimitsConf));
+ rateLimits = perDeviceLimits.computeIfAbsent(deviceId, id -> new TbRateLimits(perDevicesLimitsConf));
if (!rateLimits.tryConsume()) {
if (callback != null) {
callback.onError(new TbRateLimitsException(EntityType.DEVICE));
@@ -271,8 +272,8 @@ public abstract class AbstractTransportService implements TransportService {
public void init() {
if (rateLimitEnabled) {
//Just checking the configuration parameters
- new TbTransportRateLimits(perTenantLimitsConf);
- new TbTransportRateLimits(perDevicesLimitsConf);
+ new TbRateLimits(perTenantLimitsConf);
+ new TbRateLimits(perDevicesLimitsConf);
}
this.schedulerExecutor = Executors.newSingleThreadScheduledExecutor();
this.transportCallbackExecutor = new ThreadPoolExecutor(0, 20, 60L, TimeUnit.SECONDS, new SynchronousQueue<>());
diff --git a/dao/src/main/java/org/thingsboard/server/dao/Dao.java b/dao/src/main/java/org/thingsboard/server/dao/Dao.java
index c82379043b..041131dc7f 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/Dao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/Dao.java
@@ -16,20 +16,21 @@
package org.thingsboard.server.dao;
import com.google.common.util.concurrent.ListenableFuture;
+import org.thingsboard.server.common.data.id.TenantId;
import java.util.List;
import java.util.UUID;
public interface Dao {
- List find();
+ List find(TenantId tenantId);
- T findById(UUID id);
+ T findById(TenantId tenantId, UUID id);
- ListenableFuture findByIdAsync(UUID id);
+ ListenableFuture findByIdAsync(TenantId tenantId, UUID id);
- T save(T t);
+ T save(TenantId tenantId, T t);
- boolean removeById(UUID id);
+ boolean removeById(TenantId tenantId, UUID id);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java
index fd625eab3f..5d5456fb44 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDao.java
@@ -33,9 +33,9 @@ public interface AlarmDao extends Dao {
ListenableFuture findLatestByOriginatorAndType(TenantId tenantId, EntityId originator, String type);
- ListenableFuture findAlarmByIdAsync(UUID key);
+ ListenableFuture findAlarmByIdAsync(TenantId tenantId, UUID key);
- Alarm save(Alarm alarm);
+ Alarm save(TenantId tenantId, Alarm alarm);
- ListenableFuture> findAlarms(AlarmQuery query);
+ ListenableFuture> findAlarms(TenantId tenantId, AlarmQuery query);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java
index 6638818207..aace8327e9 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmService.java
@@ -35,17 +35,17 @@ public interface AlarmService {
Alarm createOrUpdateAlarm(Alarm alarm);
- ListenableFuture ackAlarm(AlarmId alarmId, long ackTs);
+ ListenableFuture ackAlarm(TenantId tenantId, AlarmId alarmId, long ackTs);
- ListenableFuture clearAlarm(AlarmId alarmId, JsonNode details, long ackTs);
+ ListenableFuture clearAlarm(TenantId tenantId, AlarmId alarmId, JsonNode details, long ackTs);
- ListenableFuture findAlarmByIdAsync(AlarmId alarmId);
+ ListenableFuture findAlarmByIdAsync(TenantId tenantId, AlarmId alarmId);
- ListenableFuture findAlarmInfoByIdAsync(AlarmId alarmId);
+ ListenableFuture findAlarmInfoByIdAsync(TenantId tenantId, AlarmId alarmId);
- ListenableFuture> findAlarms(AlarmQuery query);
+ ListenableFuture> findAlarms(TenantId tenantId, AlarmQuery query);
- AlarmSeverity findHighestAlarmSeverity(EntityId entityId, AlarmSearchStatus alarmSearchStatus,
+ AlarmSeverity findHighestAlarmSeverity(TenantId tenantId, EntityId entityId, AlarmSearchStatus alarmSearchStatus,
AlarmStatus alarmStatus);
ListenableFuture findLatestByOriginatorAndType(TenantId tenantId, EntityId originator, String type);
diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java
index 23fe85a802..d0698b6986 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/BaseAlarmService.java
@@ -91,7 +91,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
@Override
public Alarm createOrUpdateAlarm(Alarm alarm) {
- alarmDataValidator.validate(alarm);
+ alarmDataValidator.validate(alarm, Alarm::getTenantId);
try {
if (alarm.getStartTs() == 0L) {
alarm.setStartTs(System.currentTimeMillis());
@@ -120,7 +120,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
private Alarm createAlarm(Alarm alarm) throws InterruptedException, ExecutionException {
log.debug("New Alarm : {}", alarm);
- Alarm saved = alarmDao.save(alarm);
+ Alarm saved = alarmDao.save(alarm.getTenantId(), alarm);
createAlarmRelations(saved);
return saved;
}
@@ -129,17 +129,17 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
if (alarm.isPropagate()) {
EntityRelationsQuery query = new EntityRelationsQuery();
query.setParameters(new RelationsSearchParameters(alarm.getOriginator(), EntitySearchDirection.TO, Integer.MAX_VALUE));
- List parentEntities = relationService.findByQuery(query).get().stream().map(r -> r.getFrom()).collect(Collectors.toList());
+ List parentEntities = relationService.findByQuery(alarm.getTenantId(), query).get().stream().map(EntityRelation::getFrom).collect(Collectors.toList());
for (EntityId parentId : parentEntities) {
- createAlarmRelation(parentId, alarm.getId(), alarm.getStatus(), true);
+ createAlarmRelation(alarm.getTenantId(), parentId, alarm.getId(), alarm.getStatus(), true);
}
}
- createAlarmRelation(alarm.getOriginator(), alarm.getId(), alarm.getStatus(), true);
+ createAlarmRelation(alarm.getTenantId(), alarm.getOriginator(), alarm.getId(), alarm.getStatus(), true);
}
private ListenableFuture updateAlarm(Alarm update) {
- alarmDataValidator.validate(update);
- return getAndUpdate(update.getId(), new Function() {
+ alarmDataValidator.validate(update, Alarm::getTenantId);
+ return getAndUpdate(update.getTenantId(), update.getId(), new Function() {
@Nullable
@Override
public Alarm apply(@Nullable Alarm alarm) {
@@ -157,7 +157,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
AlarmStatus newStatus = newAlarm.getStatus();
boolean oldPropagate = oldAlarm.isPropagate();
boolean newPropagate = newAlarm.isPropagate();
- Alarm result = alarmDao.save(merge(oldAlarm, newAlarm));
+ Alarm result = alarmDao.save(newAlarm.getTenantId(), merge(oldAlarm, newAlarm));
if (!oldPropagate && newPropagate) {
try {
createAlarmRelations(result);
@@ -172,8 +172,8 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
}
@Override
- public ListenableFuture ackAlarm(AlarmId alarmId, long ackTime) {
- return getAndUpdate(alarmId, new Function() {
+ public ListenableFuture ackAlarm(TenantId tenantId, AlarmId alarmId, long ackTime) {
+ return getAndUpdate(tenantId, alarmId, new Function() {
@Nullable
@Override
public Boolean apply(@Nullable Alarm alarm) {
@@ -184,7 +184,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
AlarmStatus newStatus = oldStatus.isCleared() ? AlarmStatus.CLEARED_ACK : AlarmStatus.ACTIVE_ACK;
alarm.setStatus(newStatus);
alarm.setAckTs(ackTime);
- alarmDao.save(alarm);
+ alarmDao.save(alarm.getTenantId(), alarm);
updateRelations(alarm, oldStatus, newStatus);
return true;
}
@@ -193,8 +193,8 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
}
@Override
- public ListenableFuture clearAlarm(AlarmId alarmId, JsonNode details, long clearTime) {
- return getAndUpdate(alarmId, new Function() {
+ public ListenableFuture clearAlarm(TenantId tenantId, AlarmId alarmId, JsonNode details, long clearTime) {
+ return getAndUpdate(tenantId, alarmId, new Function() {
@Nullable
@Override
public Boolean apply(@Nullable Alarm alarm) {
@@ -208,7 +208,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
if (details != null) {
alarm.setDetails(details);
}
- alarmDao.save(alarm);
+ alarmDao.save(alarm.getTenantId(), alarm);
updateRelations(alarm, oldStatus, newStatus);
return true;
}
@@ -217,21 +217,21 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
}
@Override
- public ListenableFuture findAlarmByIdAsync(AlarmId alarmId) {
+ public ListenableFuture findAlarmByIdAsync(TenantId tenantId, AlarmId alarmId) {
log.trace("Executing findAlarmById [{}]", alarmId);
validateId(alarmId, "Incorrect alarmId " + alarmId);
- return alarmDao.findAlarmByIdAsync(alarmId.getId());
+ return alarmDao.findAlarmByIdAsync(tenantId, alarmId.getId());
}
@Override
- public ListenableFuture findAlarmInfoByIdAsync(AlarmId alarmId) {
+ public ListenableFuture findAlarmInfoByIdAsync(TenantId tenantId, AlarmId alarmId) {
log.trace("Executing findAlarmInfoByIdAsync [{}]", alarmId);
validateId(alarmId, "Incorrect alarmId " + alarmId);
- return Futures.transformAsync(alarmDao.findAlarmByIdAsync(alarmId.getId()),
+ return Futures.transformAsync(alarmDao.findAlarmByIdAsync(tenantId, alarmId.getId()),
a -> {
AlarmInfo alarmInfo = new AlarmInfo(a);
return Futures.transform(
- entityService.fetchEntityNameAsync(alarmInfo.getOriginator()), originatorName -> {
+ entityService.fetchEntityNameAsync(tenantId, alarmInfo.getOriginator()), originatorName -> {
alarmInfo.setOriginatorName(originatorName);
return alarmInfo;
}
@@ -240,14 +240,14 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
}
@Override
- public ListenableFuture> findAlarms(AlarmQuery query) {
- ListenableFuture> alarms = alarmDao.findAlarms(query);
+ public ListenableFuture> findAlarms(TenantId tenantId, AlarmQuery query) {
+ ListenableFuture> alarms = alarmDao.findAlarms(tenantId, query);
if (query.getFetchOriginator() != null && query.getFetchOriginator().booleanValue()) {
alarms = Futures.transformAsync(alarms, input -> {
List> alarmFutures = new ArrayList<>(input.size());
for (AlarmInfo alarmInfo : input) {
alarmFutures.add(Futures.transform(
- entityService.fetchEntityNameAsync(alarmInfo.getOriginator()), originatorName -> {
+ entityService.fetchEntityNameAsync(tenantId, alarmInfo.getOriginator()), originatorName -> {
if (originatorName == null) {
originatorName = "Deleted";
}
@@ -269,7 +269,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
}
@Override
- public AlarmSeverity findHighestAlarmSeverity(EntityId entityId, AlarmSearchStatus alarmSearchStatus,
+ public AlarmSeverity findHighestAlarmSeverity(TenantId tenantId, EntityId entityId, AlarmSearchStatus alarmSearchStatus,
AlarmStatus alarmStatus) {
TimePageLink nextPageLink = new TimePageLink(100);
boolean hasNext = true;
@@ -279,7 +279,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
query = new AlarmQuery(entityId, nextPageLink, alarmSearchStatus, alarmStatus, false);
List alarms;
try {
- alarms = alarmDao.findAlarms(query).get();
+ alarms = alarmDao.findAlarms(tenantId, query).get();
} catch (ExecutionException | InterruptedException e) {
log.warn("Failed to find highest alarm severity. EntityId: [{}], AlarmSearchStatus: [{}], AlarmStatus: [{}]",
entityId, alarmSearchStatus, alarmStatus);
@@ -312,14 +312,14 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
}
}
- private void deleteRelation(EntityRelation alarmRelation) throws ExecutionException, InterruptedException {
+ private void deleteRelation(TenantId tenantId, EntityRelation alarmRelation) throws ExecutionException, InterruptedException {
log.debug("Deleting Alarm relation: {}", alarmRelation);
- relationService.deleteRelationAsync(alarmRelation).get();
+ relationService.deleteRelationAsync(tenantId, alarmRelation).get();
}
- private void createRelation(EntityRelation alarmRelation) throws ExecutionException, InterruptedException {
+ private void createRelation(TenantId tenantId, EntityRelation alarmRelation) throws ExecutionException, InterruptedException {
log.debug("Creating Alarm relation: {}", alarmRelation);
- relationService.saveRelationAsync(alarmRelation).get();
+ relationService.saveRelationAsync(tenantId, alarmRelation).get();
}
private Alarm merge(Alarm existing, Alarm alarm) {
@@ -344,10 +344,10 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
private void updateRelations(Alarm alarm, AlarmStatus oldStatus, AlarmStatus newStatus) {
try {
- List relations = relationService.findByToAsync(alarm.getId(), RelationTypeGroup.ALARM).get();
+ List relations = relationService.findByToAsync(alarm.getTenantId(), alarm.getId(), RelationTypeGroup.ALARM).get();
Set parents = relations.stream().map(EntityRelation::getFrom).collect(Collectors.toSet());
for (EntityId parentId : parents) {
- updateAlarmRelation(parentId, alarm.getId(), oldStatus, newStatus);
+ updateAlarmRelation(alarm.getTenantId(), parentId, alarm.getId(), oldStatus, newStatus);
}
} catch (ExecutionException | InterruptedException e) {
log.warn("[{}] Failed to update relations. Old status: [{}], New status: [{}]", alarm.getId(), oldStatus, newStatus);
@@ -355,39 +355,39 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
}
}
- private void createAlarmRelation(EntityId entityId, EntityId alarmId, AlarmStatus status, boolean createAnyRelation) {
+ private void createAlarmRelation(TenantId tenantId, EntityId entityId, EntityId alarmId, AlarmStatus status, boolean createAnyRelation) {
try {
if (createAnyRelation) {
- createRelation(new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + AlarmSearchStatus.ANY.name(), RelationTypeGroup.ALARM));
+ createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + AlarmSearchStatus.ANY.name(), RelationTypeGroup.ALARM));
}
- createRelation(new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.name(), RelationTypeGroup.ALARM));
- createRelation(new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getClearSearchStatus().name(), RelationTypeGroup.ALARM));
- createRelation(new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getAckSearchStatus().name(), RelationTypeGroup.ALARM));
+ createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.name(), RelationTypeGroup.ALARM));
+ createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getClearSearchStatus().name(), RelationTypeGroup.ALARM));
+ createRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getAckSearchStatus().name(), RelationTypeGroup.ALARM));
} catch (ExecutionException | InterruptedException e) {
log.warn("[{}] Failed to create relation. Status: [{}]", alarmId, status);
throw new RuntimeException(e);
}
}
- private void deleteAlarmRelation(EntityId entityId, EntityId alarmId, AlarmStatus status) {
+ private void deleteAlarmRelation(TenantId tenantId, EntityId entityId, EntityId alarmId, AlarmStatus status) {
try {
- deleteRelation(new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.name(), RelationTypeGroup.ALARM));
- deleteRelation(new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getClearSearchStatus().name(), RelationTypeGroup.ALARM));
- deleteRelation(new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getAckSearchStatus().name(), RelationTypeGroup.ALARM));
+ deleteRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.name(), RelationTypeGroup.ALARM));
+ deleteRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getClearSearchStatus().name(), RelationTypeGroup.ALARM));
+ deleteRelation(tenantId, new EntityRelation(entityId, alarmId, ALARM_RELATION_PREFIX + status.getAckSearchStatus().name(), RelationTypeGroup.ALARM));
} catch (ExecutionException | InterruptedException e) {
log.warn("[{}] Failed to delete relation. Status: [{}]", alarmId, status);
throw new RuntimeException(e);
}
}
- private void updateAlarmRelation(EntityId entityId, EntityId alarmId, AlarmStatus oldStatus, AlarmStatus newStatus) {
- deleteAlarmRelation(entityId, alarmId, oldStatus);
- createAlarmRelation(entityId, alarmId, newStatus, false);
+ private void updateAlarmRelation(TenantId tenantId, EntityId entityId, EntityId alarmId, AlarmStatus oldStatus, AlarmStatus newStatus) {
+ deleteAlarmRelation(tenantId, entityId, alarmId, oldStatus);
+ createAlarmRelation(tenantId, entityId, alarmId, newStatus, false);
}
- private ListenableFuture getAndUpdate(AlarmId alarmId, Function function) {
+ private ListenableFuture getAndUpdate(TenantId tenantId, AlarmId alarmId, Function function) {
validateId(alarmId, "Alarm id should be specified!");
- ListenableFuture entity = alarmDao.findAlarmByIdAsync(alarmId.getId());
+ ListenableFuture entity = alarmDao.findAlarmByIdAsync(tenantId, alarmId.getId());
return Futures.transform(entity, function, readResultsProcessingExecutor);
}
@@ -395,7 +395,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
new DataValidator() {
@Override
- protected void validateDataImpl(Alarm alarm) {
+ protected void validateDataImpl(TenantId tenantId, Alarm alarm) {
if (StringUtils.isEmpty(alarm.getType())) {
throw new DataValidationException("Alarm type should be specified!");
}
@@ -411,7 +411,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
if (alarm.getTenantId() == null) {
throw new DataValidationException("Alarm should be assigned to tenant!");
} else {
- Tenant tenant = tenantDao.findById(alarm.getTenantId().getId());
+ Tenant tenant = tenantDao.findById(alarm.getTenantId(), alarm.getTenantId().getId());
if (tenant == null) {
throw new DataValidationException("Alarm is referencing to non-existent tenant!");
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/alarm/CassandraAlarmDao.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/CassandraAlarmDao.java
index 646c0556ae..ed8666fa81 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/alarm/CassandraAlarmDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/CassandraAlarmDao.java
@@ -73,9 +73,9 @@ public class CassandraAlarmDao extends CassandraAbstractModelDao> findAlarms(AlarmQuery query) {
+ public ListenableFuture> findAlarms(TenantId tenantId, AlarmQuery query) {
log.trace("Try to find alarms by entity [{}], searchStatus [{}], status [{}] and pageLink [{}]", query.getAffectedEntityId(), query.getSearchStatus(), query.getStatus(), query.getPageLink());
EntityId affectedEntity = query.getAffectedEntityId();
String searchStatusName;
@@ -104,12 +104,12 @@ public class CassandraAlarmDao extends CassandraAbstractModelDao> relations = relationDao.findRelations(affectedEntity, relationType, RelationTypeGroup.ALARM, EntityType.ALARM, query.getPageLink());
+ ListenableFuture> relations = relationDao.findRelations(tenantId, affectedEntity, relationType, RelationTypeGroup.ALARM, EntityType.ALARM, query.getPageLink());
return Futures.transformAsync(relations, input -> {
List> alarmFutures = new ArrayList<>(input.size());
for (EntityRelation relation : input) {
alarmFutures.add(Futures.transform(
- findAlarmByIdAsync(relation.getTo().getId()),
+ findAlarmByIdAsync(tenantId, relation.getTo().getId()),
AlarmInfo::new));
}
return Futures.successfulAsList(alarmFutures);
@@ -117,11 +117,11 @@ public class CassandraAlarmDao extends CassandraAbstractModelDao findAlarmByIdAsync(UUID key) {
+ public ListenableFuture findAlarmByIdAsync(TenantId tenantId, UUID key) {
log.debug("Get alarm by id {}", key);
Select.Where query = select().from(ALARM_BY_ID_VIEW_NAME).where(eq(ModelConstants.ID_PROPERTY, key));
query.limit(1);
log.trace("Execute query {}", query);
- return findOneByStatementAsync(query);
+ return findOneByStatementAsync(tenantId, query);
}
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java
index 3aa23d1504..be989d5c66 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java
@@ -18,6 +18,7 @@ package org.thingsboard.server.dao.asset;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.EntitySubtype;
import org.thingsboard.server.common.data.asset.Asset;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.dao.Dao;
@@ -37,7 +38,7 @@ public interface AssetDao extends Dao {
* @param asset the asset object
* @return saved asset object
*/
- Asset save(Asset asset);
+ Asset save(TenantId tenantId, Asset asset);
/**
* Find assets by tenantId and page link.
diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetService.java b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetService.java
index e373f0e27a..3ecbbd2d8d 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetService.java
@@ -30,19 +30,19 @@ import java.util.Optional;
public interface AssetService {
- Asset findAssetById(AssetId assetId);
+ Asset findAssetById(TenantId tenantId, AssetId assetId);
- ListenableFuture findAssetByIdAsync(AssetId assetId);
+ ListenableFuture findAssetByIdAsync(TenantId tenantId, AssetId assetId);
Asset findAssetByTenantIdAndName(TenantId tenantId, String name);
Asset saveAsset(Asset asset);
- Asset assignAssetToCustomer(AssetId assetId, CustomerId customerId);
+ Asset assignAssetToCustomer(TenantId tenantId, AssetId assetId, CustomerId customerId);
- Asset unassignAssetFromCustomer(AssetId assetId);
+ Asset unassignAssetFromCustomer(TenantId tenantId, AssetId assetId);
- void deleteAsset(AssetId assetId);
+ void deleteAsset(TenantId tenantId, AssetId assetId);
TextPageData findAssetsByTenantId(TenantId tenantId, TextPageLink pageLink);
@@ -60,7 +60,7 @@ public interface AssetService {
void unassignCustomerAssets(TenantId tenantId, CustomerId customerId);
- ListenableFuture> findAssetsByQuery(AssetSearchQuery query);
+ ListenableFuture> findAssetsByQuery(TenantId tenantId, AssetSearchQuery query);
ListenableFuture> findAssetTypesByTenantId(TenantId tenantId);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java
index ebd056048b..0fca8c5de6 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java
@@ -86,17 +86,17 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
private CacheManager cacheManager;
@Override
- public Asset findAssetById(AssetId assetId) {
+ public Asset findAssetById(TenantId tenantId, AssetId assetId) {
log.trace("Executing findAssetById [{}]", assetId);
validateId(assetId, INCORRECT_ASSET_ID + assetId);
- return assetDao.findById(assetId.getId());
+ return assetDao.findById(tenantId, assetId.getId());
}
@Override
- public ListenableFuture findAssetByIdAsync(AssetId assetId) {
+ public ListenableFuture findAssetByIdAsync(TenantId tenantId, AssetId assetId) {
log.trace("Executing findAssetById [{}]", assetId);
validateId(assetId, INCORRECT_ASSET_ID + assetId);
- return assetDao.findByIdAsync(assetId.getId());
+ return assetDao.findByIdAsync(tenantId, assetId.getId());
}
@Cacheable(cacheNames = ASSET_CACHE, key = "{#tenantId, #name}")
@@ -112,31 +112,31 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
@Override
public Asset saveAsset(Asset asset) {
log.trace("Executing saveAsset [{}]", asset);
- assetValidator.validate(asset);
- return assetDao.save(asset);
+ assetValidator.validate(asset, Asset::getTenantId);
+ return assetDao.save(asset.getTenantId(), asset);
}
@Override
- public Asset assignAssetToCustomer(AssetId assetId, CustomerId customerId) {
- Asset asset = findAssetById(assetId);
+ public Asset assignAssetToCustomer(TenantId tenantId, AssetId assetId, CustomerId customerId) {
+ Asset asset = findAssetById(tenantId, assetId);
asset.setCustomerId(customerId);
return saveAsset(asset);
}
@Override
- public Asset unassignAssetFromCustomer(AssetId assetId) {
- Asset asset = findAssetById(assetId);
+ public Asset unassignAssetFromCustomer(TenantId tenantId, AssetId assetId) {
+ Asset asset = findAssetById(tenantId, assetId);
asset.setCustomerId(null);
return saveAsset(asset);
}
@Override
- public void deleteAsset(AssetId assetId) {
+ public void deleteAsset(TenantId tenantId, AssetId assetId) {
log.trace("Executing deleteAsset [{}]", assetId);
validateId(assetId, INCORRECT_ASSET_ID + assetId);
- deleteEntityRelations(assetId);
+ deleteEntityRelations(tenantId, assetId);
- Asset asset = assetDao.findById(assetId.getId());
+ Asset asset = assetDao.findById(tenantId, assetId.getId());
try {
List entityViews = entityViewService.findEntityViewsByTenantIdAndEntityIdAsync(asset.getTenantId(), assetId).get();
if (entityViews != null && !entityViews.isEmpty()) {
@@ -153,7 +153,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
Cache cache = cacheManager.getCache(ASSET_CACHE);
cache.evict(list);
- assetDao.removeById(assetId.getId());
+ assetDao.removeById(tenantId, assetId.getId());
}
@Override
@@ -187,7 +187,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
public void deleteAssetsByTenantId(TenantId tenantId) {
log.trace("Executing deleteAssetsByTenantId, tenantId [{}]", tenantId);
validateId(tenantId, INCORRECT_TENANT_ID + tenantId);
- tenantAssetsRemover.removeEntities(tenantId);
+ tenantAssetsRemover.removeEntities(tenantId, tenantId);
}
@Override
@@ -225,24 +225,24 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
log.trace("Executing unassignCustomerAssets, tenantId [{}], customerId [{}]", tenantId, customerId);
validateId(tenantId, INCORRECT_TENANT_ID + tenantId);
validateId(customerId, INCORRECT_CUSTOMER_ID + customerId);
- new CustomerAssetsUnassigner(tenantId).removeEntities(customerId);
+ customerAssetsUnasigner.removeEntities(tenantId, customerId);
}
@Override
- public ListenableFuture> findAssetsByQuery(AssetSearchQuery query) {
- ListenableFuture> relations = relationService.findByQuery(query.toEntitySearchQuery());
+ public ListenableFuture> findAssetsByQuery(TenantId tenantId, AssetSearchQuery query) {
+ ListenableFuture> relations = relationService.findByQuery(tenantId, query.toEntitySearchQuery());
ListenableFuture> assets = Futures.transformAsync(relations, r -> {
EntitySearchDirection direction = query.toEntitySearchQuery().getParameters().getDirection();
List> futures = new ArrayList<>();
for (EntityRelation relation : r) {
EntityId entityId = direction == EntitySearchDirection.FROM ? relation.getTo() : relation.getFrom();
if (entityId.getEntityType() == EntityType.ASSET) {
- futures.add(findAssetByIdAsync(new AssetId(entityId.getId())));
+ futures.add(findAssetByIdAsync(tenantId, new AssetId(entityId.getId())));
}
}
return Futures.successfulAsList(futures);
});
- assets = Futures.transform(assets, (Function, List>)assetList ->
+ assets = Futures.transform(assets, assetList ->
assetList == null ? Collections.emptyList() : assetList.stream().filter(asset -> query.getAssetTypes().contains(asset.getType())).collect(Collectors.toList())
);
return assets;
@@ -254,7 +254,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
validateId(tenantId, INCORRECT_TENANT_ID + tenantId);
ListenableFuture> tenantAssetTypes = assetDao.findTenantAssetTypesAsync(tenantId.getId());
return Futures.transform(tenantAssetTypes,
- (Function, List>) assetTypes -> {
+ assetTypes -> {
assetTypes.sort(Comparator.comparing(EntitySubtype::getType));
return assetTypes;
});
@@ -264,7 +264,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
new DataValidator() {
@Override
- protected void validateCreate(Asset asset) {
+ protected void validateCreate(TenantId tenantId, Asset asset) {
assetDao.findAssetsByTenantIdAndName(asset.getTenantId().getId(), asset.getName()).ifPresent(
d -> {
throw new DataValidationException("Asset with such name already exists!");
@@ -273,7 +273,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
}
@Override
- protected void validateUpdate(Asset asset) {
+ protected void validateUpdate(TenantId tenantId, Asset asset) {
assetDao.findAssetsByTenantIdAndName(asset.getTenantId().getId(), asset.getName()).ifPresent(
d -> {
if (!d.getId().equals(asset.getId())) {
@@ -284,7 +284,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
}
@Override
- protected void validateDataImpl(Asset asset) {
+ protected void validateDataImpl(TenantId tenantId, Asset asset) {
if (StringUtils.isEmpty(asset.getType())) {
throw new DataValidationException("Asset type should be specified!");
}
@@ -294,7 +294,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
if (asset.getTenantId() == null) {
throw new DataValidationException("Asset should be assigned to tenant!");
} else {
- Tenant tenant = tenantDao.findById(asset.getTenantId().getId());
+ Tenant tenant = tenantDao.findById(tenantId, asset.getTenantId().getId());
if (tenant == null) {
throw new DataValidationException("Asset is referencing to non-existent tenant!");
}
@@ -302,7 +302,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
if (asset.getCustomerId() == null) {
asset.setCustomerId(new CustomerId(NULL_UUID));
} else if (!asset.getCustomerId().getId().equals(NULL_UUID)) {
- Customer customer = customerDao.findById(asset.getCustomerId().getId());
+ Customer customer = customerDao.findById(tenantId, asset.getCustomerId().getId());
if (customer == null) {
throw new DataValidationException("Can't assign asset to non-existent customer!");
}
@@ -314,35 +314,29 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
};
private PaginatedRemover tenantAssetsRemover =
- new PaginatedRemover() {
+ new PaginatedRemover() {
- @Override
- protected List findEntities(TenantId id, TextPageLink pageLink) {
- return assetDao.findAssetsByTenantId(id.getId(), pageLink);
- }
-
- @Override
- protected void removeEntity(Asset entity) {
- deleteAsset(new AssetId(entity.getId().getId()));
- }
- };
+ @Override
+ protected List findEntities(TenantId tenantId, TenantId id, TextPageLink pageLink) {
+ return assetDao.findAssetsByTenantId(id.getId(), pageLink);
+ }
- class CustomerAssetsUnassigner extends PaginatedRemover {
+ @Override
+ protected void removeEntity(TenantId tenantId, Asset entity) {
+ deleteAsset(tenantId, new AssetId(entity.getId().getId()));
+ }
+ };
- private TenantId tenantId;
-
- CustomerAssetsUnassigner(TenantId tenantId) {
- this.tenantId = tenantId;
- }
+ private PaginatedRemover customerAssetsUnasigner = new PaginatedRemover() {
@Override
- protected List findEntities(CustomerId id, TextPageLink pageLink) {
+ protected List findEntities(TenantId tenantId, CustomerId id, TextPageLink pageLink) {
return assetDao.findAssetsByTenantIdAndCustomerId(tenantId.getId(), id.getId(), pageLink);
}
@Override
- protected void removeEntity(Asset entity) {
- unassignAssetFromCustomer(new AssetId(entity.getId().getId()));
+ protected void removeEntity(TenantId tenantId, Asset entity) {
+ unassignAssetFromCustomer(tenantId, new AssetId(entity.getId().getId()));
}
- }
+ };
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/CassandraAssetDao.java b/dao/src/main/java/org/thingsboard/server/dao/asset/CassandraAssetDao.java
index 80a6f43f75..48c8950c94 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/asset/CassandraAssetDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/asset/CassandraAssetDao.java
@@ -28,6 +28,7 @@ import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.EntitySubtype;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.asset.Asset;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.model.EntitySubtypeEntity;
@@ -77,19 +78,19 @@ public class CassandraAssetDao extends CassandraAbstractSearchTextDao findAssetsByTenantId(UUID tenantId, TextPageLink pageLink) {
log.debug("Try to find assets by tenantId [{}] and pageLink [{}]", tenantId, pageLink);
- List assetEntities = findPageWithTextSearch(ASSET_BY_TENANT_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
+ List assetEntities = findPageWithTextSearch(new TenantId(tenantId), ASSET_BY_TENANT_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
Collections.singletonList(eq(ASSET_TENANT_ID_PROPERTY, tenantId)), pageLink);
log.trace("Found assets [{}] by tenantId [{}] and pageLink [{}]", assetEntities, tenantId, pageLink);
@@ -99,7 +100,7 @@ public class CassandraAssetDao extends CassandraAbstractSearchTextDao findAssetsByTenantIdAndType(UUID tenantId, String type, TextPageLink pageLink) {
log.debug("Try to find assets by tenantId [{}], type [{}] and pageLink [{}]", tenantId, type, pageLink);
- List assetEntities = findPageWithTextSearch(ASSET_BY_TENANT_BY_TYPE_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
+ List assetEntities = findPageWithTextSearch(new TenantId(tenantId), ASSET_BY_TENANT_BY_TYPE_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
Arrays.asList(eq(ASSET_TYPE_PROPERTY, type),
eq(ASSET_TENANT_ID_PROPERTY, tenantId)), pageLink);
log.trace("Found assets [{}] by tenantId [{}], type [{}] and pageLink [{}]", assetEntities, tenantId, type, pageLink);
@@ -112,13 +113,13 @@ public class CassandraAssetDao extends CassandraAbstractSearchTextDao findAssetsByTenantIdAndCustomerId(UUID tenantId, UUID customerId, TextPageLink pageLink) {
log.debug("Try to find assets by tenantId [{}], customerId[{}] and pageLink [{}]", tenantId, customerId, pageLink);
- List assetEntities = findPageWithTextSearch(ASSET_BY_CUSTOMER_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
+ List assetEntities = findPageWithTextSearch(new TenantId(tenantId), ASSET_BY_CUSTOMER_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
Arrays.asList(eq(ASSET_CUSTOMER_ID_PROPERTY, customerId),
eq(ASSET_TENANT_ID_PROPERTY, tenantId)),
pageLink);
@@ -130,7 +131,7 @@ public class CassandraAssetDao extends CassandraAbstractSearchTextDao findAssetsByTenantIdAndCustomerIdAndType(UUID tenantId, UUID customerId, String type, TextPageLink pageLink) {
log.debug("Try to find assets by tenantId [{}], customerId [{}], type [{}] and pageLink [{}]", tenantId, customerId, type, pageLink);
- List assetEntities = findPageWithTextSearch(ASSET_BY_CUSTOMER_BY_TYPE_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
+ List assetEntities = findPageWithTextSearch(new TenantId(tenantId), ASSET_BY_CUSTOMER_BY_TYPE_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
Arrays.asList(eq(ASSET_TYPE_PROPERTY, type),
eq(ASSET_CUSTOMER_ID_PROPERTY, customerId),
eq(ASSET_TENANT_ID_PROPERTY, tenantId)),
@@ -148,7 +149,7 @@ public class CassandraAssetDao extends CassandraAbstractSearchTextDao>() {
@Nullable
@Override
diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java
index c9510d276a..a6651648f0 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesDao.java
@@ -17,6 +17,7 @@ package org.thingsboard.server.dao.attributes;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.id.EntityId;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import java.util.Collection;
@@ -28,13 +29,13 @@ import java.util.Optional;
*/
public interface AttributesDao {
- ListenableFuture> find(EntityId entityId, String attributeType, String attributeKey);
+ ListenableFuture> find(TenantId tenantId, EntityId entityId, String attributeType, String attributeKey);
- ListenableFuture> find(EntityId entityId, String attributeType, Collection attributeKey);
+ ListenableFuture> find(TenantId tenantId, EntityId entityId, String attributeType, Collection attributeKey);
- ListenableFuture> findAll(EntityId entityId, String attributeType);
+ ListenableFuture> findAll(TenantId tenantId, EntityId entityId, String attributeType);
- ListenableFuture save(EntityId entityId, String attributeType, AttributeKvEntry attribute);
+ ListenableFuture save(TenantId tenantId, EntityId entityId, String attributeType, AttributeKvEntry attribute);
- ListenableFuture> removeAll(EntityId entityId, String attributeType, List keys);
+ ListenableFuture> removeAll(TenantId tenantId, EntityId entityId, String attributeType, List keys);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java
index fa3a2b16f8..7c6b2e748b 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/attributes/AttributesService.java
+++ b/dao/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 org.thingsboard.server.common.data.id.EntityId;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import java.util.Collection;
@@ -28,13 +29,13 @@ import java.util.Optional;
*/
public interface AttributesService {
- ListenableFuture> find(EntityId entityId, String scope, String attributeKey);
+ ListenableFuture> find(TenantId tenantId, EntityId entityId, String scope, String attributeKey);
- ListenableFuture> find(EntityId entityId, String scope, Collection attributeKeys);
+ ListenableFuture> find(TenantId tenantId, EntityId entityId, String scope, Collection attributeKeys);
- ListenableFuture> findAll(EntityId entityId, String scope);
+ ListenableFuture> findAll(TenantId tenantId, EntityId entityId, String scope);
- ListenableFuture> save(EntityId entityId, String scope, List attributes);
+ ListenableFuture> save(TenantId tenantId, EntityId entityId, String scope, List attributes);
- ListenableFuture> removeAll(EntityId entityId, String scope, List attributeKeys);
+ ListenableFuture> removeAll(TenantId tenantId, EntityId entityId, String scope, List attributeKeys);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java
index 079772bd50..29394bb770 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesService.java
@@ -21,6 +21,7 @@ import com.google.common.util.concurrent.ListenableFuture;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.id.EntityId;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.dao.exception.IncorrectParameterException;
import org.thingsboard.server.dao.service.Validator;
@@ -39,40 +40,40 @@ public class BaseAttributesService implements AttributesService {
private AttributesDao attributesDao;
@Override
- public ListenableFuture> find(EntityId entityId, String scope, String attributeKey) {
+ public ListenableFuture> find(TenantId tenantId, EntityId entityId, String scope, String attributeKey) {
validate(entityId, scope);
Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey);
- return attributesDao.find(entityId, scope, attributeKey);
+ return attributesDao.find(tenantId, entityId, scope, attributeKey);
}
@Override
- public ListenableFuture> find(EntityId entityId, String scope, Collection attributeKeys) {
+ public ListenableFuture> find(TenantId tenantId, EntityId entityId, String scope, Collection attributeKeys) {
validate(entityId, scope);
attributeKeys.forEach(attributeKey -> Validator.validateString(attributeKey, "Incorrect attribute key " + attributeKey));
- return attributesDao.find(entityId, scope, attributeKeys);
+ return attributesDao.find(tenantId, entityId, scope, attributeKeys);
}
@Override
- public ListenableFuture> findAll(EntityId entityId, String scope) {
+ public ListenableFuture> findAll(TenantId tenantId, EntityId entityId, String scope) {
validate(entityId, scope);
- return attributesDao.findAll(entityId, scope);
+ return attributesDao.findAll(tenantId, entityId, scope);
}
@Override
- public ListenableFuture> save(EntityId entityId, String scope, List attributes) {
+ public ListenableFuture> save(TenantId tenantId, EntityId entityId, String scope, List attributes) {
validate(entityId, scope);
attributes.forEach(attribute -> validate(attribute));
List> futures = Lists.newArrayListWithExpectedSize(attributes.size());
for (AttributeKvEntry attribute : attributes) {
- futures.add(attributesDao.save(entityId, scope, attribute));
+ futures.add(attributesDao.save(tenantId, entityId, scope, attribute));
}
return Futures.allAsList(futures);
}
@Override
- public ListenableFuture> removeAll(EntityId entityId, String scope, List keys) {
+ public ListenableFuture> removeAll(TenantId tenantId, EntityId entityId, String scope, List keys) {
validate(entityId, scope);
- return attributesDao.removeAll(entityId, scope, keys);
+ return attributesDao.removeAll(tenantId, entityId, scope, keys);
}
private static void validate(EntityId id, String scope) {
diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java
index a736a69c5e..3674f3d64e 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java
@@ -28,6 +28,7 @@ import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.id.EntityId;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
import org.thingsboard.server.dao.model.ModelConstants;
@@ -73,22 +74,22 @@ public class CassandraBaseAttributesDao extends CassandraAbstractAsyncDao implem
}
@Override
- public ListenableFuture> find(EntityId entityId, String attributeType, String attributeKey) {
+ public ListenableFuture> find(TenantId tenantId, EntityId entityId, String attributeType, String attributeKey) {
Select.Where select = select().from(ATTRIBUTES_KV_CF)
.where(eq(ENTITY_TYPE_COLUMN, entityId.getEntityType()))
.and(eq(ENTITY_ID_COLUMN, entityId.getId()))
.and(eq(ATTRIBUTE_TYPE_COLUMN, attributeType))
.and(eq(ATTRIBUTE_KEY_COLUMN, attributeKey));
log.trace("Generated query [{}] for entityId {} and key {}", select, entityId, attributeKey);
- return Futures.transform(executeAsyncRead(select), (Function super ResultSet, ? extends Optional>) input ->
+ return Futures.transform(executeAsyncRead(tenantId, select), (Function super ResultSet, ? extends Optional>) input ->
Optional.ofNullable(convertResultToAttributesKvEntry(attributeKey, input.one()))
, readResultsProcessingExecutor);
}
@Override
- public ListenableFuture> find(EntityId entityId, String attributeType, Collection attributeKeys) {
+ public ListenableFuture> find(TenantId tenantId, EntityId entityId, String attributeType, Collection attributeKeys) {
List>> entries = new ArrayList<>();
- attributeKeys.forEach(attributeKey -> entries.add(find(entityId, attributeType, attributeKey)));
+ attributeKeys.forEach(attributeKey -> entries.add(find(tenantId, entityId, attributeType, attributeKey)));
return Futures.transform(Futures.allAsList(entries), (Function>, ? extends List>) input -> {
List result = new ArrayList<>();
input.stream().filter(opt -> opt.isPresent()).forEach(opt -> result.add(opt.get()));
@@ -98,19 +99,19 @@ public class CassandraBaseAttributesDao extends CassandraAbstractAsyncDao implem
@Override
- public ListenableFuture> findAll(EntityId entityId, String attributeType) {
+ public ListenableFuture> findAll(TenantId tenantId, EntityId entityId, String attributeType) {
Select.Where select = select().from(ATTRIBUTES_KV_CF)
.where(eq(ENTITY_TYPE_COLUMN, entityId.getEntityType()))
.and(eq(ENTITY_ID_COLUMN, entityId.getId()))
.and(eq(ATTRIBUTE_TYPE_COLUMN, attributeType));
log.trace("Generated query [{}] for entityId {} and attributeType {}", select, entityId, attributeType);
- return Futures.transform(executeAsyncRead(select), (Function super ResultSet, ? extends List>) input ->
+ return Futures.transform(executeAsyncRead(tenantId, select), (Function super ResultSet, ? extends List>) input ->
convertResultToAttributesKvEntryList(input)
, readResultsProcessingExecutor);
}
@Override
- public ListenableFuture save(EntityId entityId, String attributeType, AttributeKvEntry attribute) {
+ public ListenableFuture save(TenantId tenantId, EntityId entityId, String attributeType, AttributeKvEntry attribute) {
BoundStatement stmt = getSaveStmt().bind();
stmt.setString(0, entityId.getEntityType().name());
stmt.setUUID(1, entityId.getId());
@@ -137,26 +138,26 @@ public class CassandraBaseAttributesDao extends CassandraAbstractAsyncDao implem
stmt.setToNull(8);
}
log.trace("Generated save stmt [{}] for entityId {} and attributeType {} and attribute", stmt, entityId, attributeType, attribute);
- return getFuture(executeAsyncWrite(stmt), rs -> null);
+ return getFuture(executeAsyncWrite(tenantId, stmt), rs -> null);
}
@Override
- public ListenableFuture> removeAll(EntityId entityId, String attributeType, List keys) {
+ public ListenableFuture> removeAll(TenantId tenantId, EntityId entityId, String attributeType, List keys) {
List> futures = keys
.stream()
- .map(key -> delete(entityId, attributeType, key))
+ .map(key -> delete(tenantId, entityId, attributeType, key))
.collect(Collectors.toList());
return Futures.allAsList(futures);
}
- private ListenableFuture delete(EntityId entityId, String attributeType, String key) {
+ private ListenableFuture delete(TenantId tenantId, EntityId entityId, String attributeType, String key) {
Statement delete = QueryBuilder.delete().all().from(ModelConstants.ATTRIBUTES_KV_CF)
.where(eq(ENTITY_TYPE_COLUMN, entityId.getEntityType()))
.and(eq(ENTITY_ID_COLUMN, entityId.getId()))
.and(eq(ATTRIBUTE_TYPE_COLUMN, attributeType))
.and(eq(ATTRIBUTE_KEY_COLUMN, key));
log.debug("Remove request: {}", delete.toString());
- return getFuture(executeAsyncWrite(delete), rs -> null);
+ return getFuture(executeAsyncWrite(tenantId, delete), rs -> null);
}
private PreparedStatement getSaveStmt() {
diff --git a/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java
index ecb2bd5a27..24c6a274f1 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java
@@ -128,7 +128,7 @@ public class AuditLogServiceImpl implements AuditLogService {
entityName = entity.getName();
} else {
try {
- entityName = entityService.fetchEntityNameAsync(entityId).get();
+ entityName = entityService.fetchEntityNameAsync(tenantId, entityId).get();
} catch (Exception ex) {}
}
if (e != null) {
@@ -315,7 +315,7 @@ public class AuditLogServiceImpl implements AuditLogService {
AuditLog auditLogEntry = createAuditLogEntry(tenantId, entityId, entityName, customerId, userId, userName,
actionType, actionData, actionStatus, actionFailureDetails);
log.trace("Executing logAction [{}]", auditLogEntry);
- auditLogValidator.validate(auditLogEntry);
+ auditLogValidator.validate(auditLogEntry, AuditLog::getTenantId);
List> futures = Lists.newArrayListWithExpectedSize(INSERTS_PER_ENTRY);
futures.add(auditLogDao.savePartitionsByTenantId(auditLogEntry));
futures.add(auditLogDao.saveByTenantId(auditLogEntry));
@@ -331,7 +331,7 @@ public class AuditLogServiceImpl implements AuditLogService {
private DataValidator auditLogValidator =
new DataValidator() {
@Override
- protected void validateDataImpl(AuditLog auditLog) {
+ protected void validateDataImpl(TenantId tenantId, AuditLog auditLog) {
if (auditLog.getEntityId() == null) {
throw new DataValidationException("Entity Id should be specified!");
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/audit/CassandraAuditLogDao.java b/dao/src/main/java/org/thingsboard/server/dao/audit/CassandraAuditLogDao.java
index 764f4687d3..f2b2973996 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/audit/CassandraAuditLogDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/audit/CassandraAuditLogDao.java
@@ -32,6 +32,7 @@ import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.audit.AuditLog;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.dao.DaoUtil;
@@ -142,7 +143,7 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao null);
+ return getFuture(executeAsyncWrite(auditLog.getTenantId(), stmt), rs -> null);
}
@Override
@@ -151,7 +152,7 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao null);
+ return getFuture(executeAsyncWrite(auditLog.getTenantId(), stmt), rs -> null);
}
@Override
@@ -160,7 +161,7 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao null);
+ return getFuture(executeAsyncWrite(auditLog.getTenantId(), stmt), rs -> null);
}
@Override
@@ -169,11 +170,11 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao null);
+ return getFuture(executeAsyncWrite(auditLog.getTenantId(), stmt), rs -> null);
}
private BoundStatement setSaveStmtVariables(BoundStatement stmt, AuditLog auditLog, long partition) {
- stmt.setUUID(0, auditLog.getId().getId())
+ stmt.setUUID(0, auditLog.getId().getId())
.setUUID(1, auditLog.getTenantId().getId())
.setUUID(2, auditLog.getCustomerId().getId())
.setUUID(3, auditLog.getEntityId().getId())
@@ -200,7 +201,7 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao null);
+ return getFuture(executeAsyncWrite(auditLog.getTenantId(), stmt), rs -> null);
}
private PreparedStatement getSaveByTenantStmt() {
@@ -249,7 +250,7 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao findAuditLogsByTenantIdAndEntityId(UUID tenantId, EntityId entityId, TimePageLink pageLink) {
log.trace("Try to find audit logs by tenant [{}], entity [{}] and pageLink [{}]", tenantId, entityId, pageLink);
- List entities = findPageWithTimeSearch(AUDIT_LOG_BY_ENTITY_ID_CF,
+ List entities = findPageWithTimeSearch(new TenantId(tenantId), AUDIT_LOG_BY_ENTITY_ID_CF,
Arrays.asList(eq(ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY, tenantId),
eq(ModelConstants.AUDIT_LOG_ENTITY_TYPE_PROPERTY, entityId.getEntityType()),
eq(ModelConstants.AUDIT_LOG_ENTITY_ID_PROPERTY, entityId.getId())),
@@ -286,7 +287,7 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao findAuditLogsByTenantIdAndCustomerId(UUID tenantId, CustomerId customerId, TimePageLink pageLink) {
log.trace("Try to find audit logs by tenant [{}], customer [{}] and pageLink [{}]", tenantId, customerId, pageLink);
- List entities = findPageWithTimeSearch(AUDIT_LOG_BY_CUSTOMER_ID_CF,
+ List entities = findPageWithTimeSearch(new TenantId(tenantId), AUDIT_LOG_BY_CUSTOMER_ID_CF,
Arrays.asList(eq(ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY, tenantId),
eq(ModelConstants.AUDIT_LOG_CUSTOMER_ID_PROPERTY, customerId.getId())),
pageLink);
@@ -297,7 +298,7 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao findAuditLogsByTenantIdAndUserId(UUID tenantId, UserId userId, TimePageLink pageLink) {
log.trace("Try to find audit logs by tenant [{}], user [{}] and pageLink [{}]", tenantId, userId, pageLink);
- List entities = findPageWithTimeSearch(AUDIT_LOG_BY_USER_ID_CF,
+ List entities = findPageWithTimeSearch(new TenantId(tenantId), AUDIT_LOG_BY_USER_ID_CF,
Arrays.asList(eq(ModelConstants.AUDIT_LOG_TENANT_ID_PROPERTY, tenantId),
eq(ModelConstants.AUDIT_LOG_USER_ID_PROPERTY, userId.getId())),
pageLink);
@@ -339,7 +340,7 @@ public class CassandraAuditLogDao extends CassandraAbstractSearchTimeDao result = componentDescriptorDao.saveIfNotExist(component);
+ public ComponentDescriptor saveComponent(TenantId tenantId, ComponentDescriptor component) {
+ componentValidator.validate(component, data -> new TenantId(EntityId.NULL_UUID));
+ Optional result = componentDescriptorDao.saveIfNotExist(tenantId, component);
if (result.isPresent()) {
return result.get();
} else {
- return componentDescriptorDao.findByClazz(component.getClazz());
+ return componentDescriptorDao.findByClazz(tenantId, component.getClazz());
}
}
@Override
- public ComponentDescriptor findById(ComponentDescriptorId componentId) {
+ public ComponentDescriptor findById(TenantId tenantId, ComponentDescriptorId componentId) {
Validator.validateId(componentId, "Incorrect component id for search request.");
- return componentDescriptorDao.findById(componentId);
+ return componentDescriptorDao.findById(tenantId, componentId);
}
@Override
- public ComponentDescriptor findByClazz(String clazz) {
+ public ComponentDescriptor findByClazz(TenantId tenantId, String clazz) {
Validator.validateString(clazz, "Incorrect clazz for search request.");
- return componentDescriptorDao.findByClazz(clazz);
+ return componentDescriptorDao.findByClazz(tenantId, clazz);
}
@Override
- public TextPageData findByTypeAndPageLink(ComponentType type, TextPageLink pageLink) {
+ public TextPageData findByTypeAndPageLink(TenantId tenantId, ComponentType type, TextPageLink pageLink) {
Validator.validatePageLink(pageLink, "Incorrect PageLink object for search plugin components request.");
- List components = componentDescriptorDao.findByTypeAndPageLink(type, pageLink);
+ List components = componentDescriptorDao.findByTypeAndPageLink(tenantId, type, pageLink);
return new TextPageData<>(components, pageLink);
}
@Override
- public TextPageData findByScopeAndTypeAndPageLink(ComponentScope scope, ComponentType type, TextPageLink pageLink) {
+ public TextPageData findByScopeAndTypeAndPageLink(TenantId tenantId, ComponentScope scope, ComponentType type, TextPageLink pageLink) {
Validator.validatePageLink(pageLink, "Incorrect PageLink object for search plugin components request.");
- List components = componentDescriptorDao.findByScopeAndTypeAndPageLink(scope, type, pageLink);
+ List components = componentDescriptorDao.findByScopeAndTypeAndPageLink(tenantId, scope, type, pageLink);
return new TextPageData<>(components, pageLink);
}
@Override
- public void deleteByClazz(String clazz) {
+ public void deleteByClazz(TenantId tenantId, String clazz) {
Validator.validateString(clazz, "Incorrect clazz for delete request.");
- componentDescriptorDao.deleteByClazz(clazz);
+ componentDescriptorDao.deleteByClazz(tenantId, clazz);
}
@Override
- public boolean validate(ComponentDescriptor component, JsonNode configuration) {
+ public boolean validate(TenantId tenantId, ComponentDescriptor component, JsonNode configuration) {
JsonValidator validator = JsonSchemaFactory.byDefault().getValidator();
try {
if (!component.getConfigurationDescriptor().has("schema")) {
@@ -109,7 +111,7 @@ public class BaseComponentDescriptorService implements ComponentDescriptorServic
private DataValidator componentValidator =
new DataValidator() {
@Override
- protected void validateDataImpl(ComponentDescriptor plugin) {
+ protected void validateDataImpl(TenantId tenantId, ComponentDescriptor plugin) {
if (plugin.getType() == null) {
throw new DataValidationException("Component type should be specified!.");
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/component/CassandraBaseComponentDescriptorDao.java b/dao/src/main/java/org/thingsboard/server/dao/component/CassandraBaseComponentDescriptorDao.java
index b5b9f15332..d7a0c7f092 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/component/CassandraBaseComponentDescriptorDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/component/CassandraBaseComponentDescriptorDao.java
@@ -23,6 +23,7 @@ import com.datastax.driver.core.utils.UUIDs;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.id.ComponentDescriptorId;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.plugin.ComponentDescriptor;
import org.thingsboard.server.common.data.plugin.ComponentScope;
@@ -62,10 +63,10 @@ public class CassandraBaseComponentDescriptorDao extends CassandraAbstractSearch
}
@Override
- public Optional saveIfNotExist(ComponentDescriptor component) {
+ public Optional saveIfNotExist(TenantId tenantId, ComponentDescriptor component) {
ComponentDescriptorEntity entity = new ComponentDescriptorEntity(component);
log.debug("Save component entity [{}]", entity);
- Optional result = saveIfNotExist(entity);
+ Optional result = saveIfNotExist(tenantId, entity);
if (log.isTraceEnabled()) {
log.trace("Saved result: [{}] for component entity [{}]", result.isPresent(), result.orElse(null));
} else {
@@ -75,9 +76,9 @@ public class CassandraBaseComponentDescriptorDao extends CassandraAbstractSearch
}
@Override
- public ComponentDescriptor findById(ComponentDescriptorId componentId) {
+ public ComponentDescriptor findById(TenantId tenantId, ComponentDescriptorId componentId) {
log.debug("Search component entity by id [{}]", componentId);
- ComponentDescriptor componentDescriptor = super.findById(componentId.getId());
+ ComponentDescriptor componentDescriptor = super.findById(tenantId, componentId.getId());
if (log.isTraceEnabled()) {
log.trace("Search result: [{}] for component entity [{}]", componentDescriptor != null, componentDescriptor);
} else {
@@ -87,11 +88,11 @@ public class CassandraBaseComponentDescriptorDao extends CassandraAbstractSearch
}
@Override
- public ComponentDescriptor findByClazz(String clazz) {
+ public ComponentDescriptor findByClazz(TenantId tenantId, String clazz) {
log.debug("Search component entity by clazz [{}]", clazz);
Select.Where query = select().from(getColumnFamilyName()).where(eq(ModelConstants.COMPONENT_DESCRIPTOR_CLASS_PROPERTY, clazz));
log.trace("Execute query [{}]", query);
- ComponentDescriptorEntity entity = findOneByStatement(query);
+ ComponentDescriptorEntity entity = findOneByStatement(tenantId, query);
if (log.isTraceEnabled()) {
log.trace("Search result: [{}] for component entity [{}]", entity != null, entity);
} else {
@@ -101,9 +102,9 @@ public class CassandraBaseComponentDescriptorDao extends CassandraAbstractSearch
}
@Override
- public List findByTypeAndPageLink(ComponentType type, TextPageLink pageLink) {
+ public List findByTypeAndPageLink(TenantId tenantId, ComponentType type, TextPageLink pageLink) {
log.debug("Try to find component by type [{}] and pageLink [{}]", type, pageLink);
- List entities = findPageWithTextSearch(ModelConstants.COMPONENT_DESCRIPTOR_BY_TYPE_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
+ List entities = findPageWithTextSearch(tenantId, ModelConstants.COMPONENT_DESCRIPTOR_BY_TYPE_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
Arrays.asList(eq(ModelConstants.COMPONENT_DESCRIPTOR_TYPE_PROPERTY, type)), pageLink);
if (log.isTraceEnabled()) {
log.trace(SEARCH_RESULT, Arrays.toString(entities.toArray()));
@@ -114,9 +115,9 @@ public class CassandraBaseComponentDescriptorDao extends CassandraAbstractSearch
}
@Override
- public List findByScopeAndTypeAndPageLink(ComponentScope scope, ComponentType type, TextPageLink pageLink) {
+ public List findByScopeAndTypeAndPageLink(TenantId tenantId, ComponentScope scope, ComponentType type, TextPageLink pageLink) {
log.debug("Try to find component by scope [{}] and type [{}] and pageLink [{}]", scope, type, pageLink);
- List entities = findPageWithTextSearch(ModelConstants.COMPONENT_DESCRIPTOR_BY_SCOPE_TYPE_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
+ List entities = findPageWithTextSearch(tenantId, ModelConstants.COMPONENT_DESCRIPTOR_BY_SCOPE_TYPE_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
Arrays.asList(eq(ModelConstants.COMPONENT_DESCRIPTOR_TYPE_PROPERTY, type),
eq(ModelConstants.COMPONENT_DESCRIPTOR_SCOPE_PROPERTY, scope.name())), pageLink);
if (log.isTraceEnabled()) {
@@ -127,34 +128,34 @@ public class CassandraBaseComponentDescriptorDao extends CassandraAbstractSearch
return DaoUtil.convertDataList(entities);
}
- public boolean removeById(UUID key) {
+ public boolean removeById(TenantId tenantId, UUID key) {
Statement delete = QueryBuilder.delete().all().from(ModelConstants.COMPONENT_DESCRIPTOR_BY_ID).where(eq(ModelConstants.ID_PROPERTY, key));
log.debug("Remove request: {}", delete.toString());
- return executeWrite(delete).wasApplied();
+ return executeWrite(tenantId, delete).wasApplied();
}
@Override
- public void deleteById(ComponentDescriptorId id) {
+ public void deleteById(TenantId tenantId, ComponentDescriptorId id) {
log.debug("Delete plugin meta-data entity by id [{}]", id);
- boolean result = removeById(id.getId());
+ boolean result = removeById(tenantId, id.getId());
log.debug("Delete result: [{}]", result);
}
@Override
- public void deleteByClazz(String clazz) {
+ public void deleteByClazz(TenantId tenantId, String clazz) {
log.debug("Delete plugin meta-data entity by id [{}]", clazz);
Statement delete = QueryBuilder.delete().all().from(getColumnFamilyName()).where(eq(ModelConstants.COMPONENT_DESCRIPTOR_CLASS_PROPERTY, clazz));
log.debug("Remove request: {}", delete.toString());
- ResultSet resultSet = executeWrite(delete);
+ ResultSet resultSet = executeWrite(tenantId, delete);
log.debug("Delete result: [{}]", resultSet.wasApplied());
}
- private Optional saveIfNotExist(ComponentDescriptorEntity entity) {
+ private Optional saveIfNotExist(TenantId tenantId, ComponentDescriptorEntity entity) {
if (entity.getId() == null) {
entity.setId(UUIDs.timeBased());
}
- ResultSet rs = executeRead(QueryBuilder.insertInto(getColumnFamilyName())
+ ResultSet rs = executeRead(tenantId, QueryBuilder.insertInto(getColumnFamilyName())
.value(ModelConstants.ID_PROPERTY, entity.getId())
.value(ModelConstants.COMPONENT_DESCRIPTOR_NAME_PROPERTY, entity.getName())
.value(ModelConstants.COMPONENT_DESCRIPTOR_CLASS_PROPERTY, entity.getClazz())
diff --git a/dao/src/main/java/org/thingsboard/server/dao/component/ComponentDescriptorDao.java b/dao/src/main/java/org/thingsboard/server/dao/component/ComponentDescriptorDao.java
index 7e874826d2..9ba0c96ae4 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/component/ComponentDescriptorDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/component/ComponentDescriptorDao.java
@@ -16,6 +16,7 @@
package org.thingsboard.server.dao.component;
import org.thingsboard.server.common.data.id.ComponentDescriptorId;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.plugin.ComponentDescriptor;
import org.thingsboard.server.common.data.plugin.ComponentScope;
@@ -30,18 +31,18 @@ import java.util.Optional;
*/
public interface ComponentDescriptorDao extends Dao {
- Optional saveIfNotExist(ComponentDescriptor component);
+ Optional saveIfNotExist(TenantId tenantId, ComponentDescriptor component);
- ComponentDescriptor findById(ComponentDescriptorId componentId);
+ ComponentDescriptor findById(TenantId tenantId, ComponentDescriptorId componentId);
- ComponentDescriptor findByClazz(String clazz);
+ ComponentDescriptor findByClazz(TenantId tenantId, String clazz);
- List findByTypeAndPageLink(ComponentType type, TextPageLink pageLink);
+ List findByTypeAndPageLink(TenantId tenantId, ComponentType type, TextPageLink pageLink);
- List findByScopeAndTypeAndPageLink(ComponentScope scope, ComponentType type, TextPageLink pageLink);
+ List findByScopeAndTypeAndPageLink(TenantId tenantId, ComponentScope scope, ComponentType type, TextPageLink pageLink);
- void deleteById(ComponentDescriptorId componentId);
+ void deleteById(TenantId tenantId, ComponentDescriptorId componentId);
- void deleteByClazz(String clazz);
+ void deleteByClazz(TenantId tenantId, String clazz);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/component/ComponentDescriptorService.java b/dao/src/main/java/org/thingsboard/server/dao/component/ComponentDescriptorService.java
index bd101def49..63cc499811 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/component/ComponentDescriptorService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/component/ComponentDescriptorService.java
@@ -17,6 +17,7 @@ package org.thingsboard.server.dao.component;
import com.fasterxml.jackson.databind.JsonNode;
import org.thingsboard.server.common.data.id.ComponentDescriptorId;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TextPageData;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.plugin.ComponentDescriptor;
@@ -28,18 +29,18 @@ import org.thingsboard.server.common.data.plugin.ComponentType;
*/
public interface ComponentDescriptorService {
- ComponentDescriptor saveComponent(ComponentDescriptor component);
+ ComponentDescriptor saveComponent(TenantId tenantId, ComponentDescriptor component);
- ComponentDescriptor findById(ComponentDescriptorId componentId);
+ ComponentDescriptor findById(TenantId tenantId, ComponentDescriptorId componentId);
- ComponentDescriptor findByClazz(String clazz);
+ ComponentDescriptor findByClazz(TenantId tenantId, String clazz);
- TextPageData findByTypeAndPageLink(ComponentType type, TextPageLink pageLink);
+ TextPageData findByTypeAndPageLink(TenantId tenantId, ComponentType type, TextPageLink pageLink);
- TextPageData findByScopeAndTypeAndPageLink(ComponentScope scope, ComponentType type, TextPageLink pageLink);
+ TextPageData findByScopeAndTypeAndPageLink(TenantId tenantId, ComponentScope scope, ComponentType type, TextPageLink pageLink);
- boolean validate(ComponentDescriptor component, JsonNode configuration);
+ boolean validate(TenantId tenantId, ComponentDescriptor component, JsonNode configuration);
- void deleteByClazz(String clazz);
+ void deleteByClazz(TenantId tenantId, String clazz);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/customer/CassandraCustomerDao.java b/dao/src/main/java/org/thingsboard/server/dao/customer/CassandraCustomerDao.java
index 598f98ac47..f1b929b59d 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/customer/CassandraCustomerDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/customer/CassandraCustomerDao.java
@@ -19,6 +19,7 @@ import com.datastax.driver.core.querybuilder.Select;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.Customer;
+import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.model.ModelConstants;
@@ -36,6 +37,7 @@ import static com.datastax.driver.core.querybuilder.QueryBuilder.select;
import static org.thingsboard.server.dao.model.ModelConstants.CUSTOMER_BY_TENANT_AND_TITLE_VIEW_NAME;
import static org.thingsboard.server.dao.model.ModelConstants.CUSTOMER_TENANT_ID_PROPERTY;
import static org.thingsboard.server.dao.model.ModelConstants.CUSTOMER_TITLE_PROPERTY;
+
@Component
@Slf4j
@NoSqlDao
@@ -54,9 +56,9 @@ public class CassandraCustomerDao extends CassandraAbstractSearchTextDao findCustomersByTenantId(UUID tenantId, TextPageLink pageLink) {
log.debug("Try to find customers by tenantId [{}] and pageLink [{}]", tenantId, pageLink);
- List customerEntities = findPageWithTextSearch(ModelConstants.CUSTOMER_BY_TENANT_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
+ List customerEntities = findPageWithTextSearch(new TenantId(tenantId), ModelConstants.CUSTOMER_BY_TENANT_AND_SEARCH_TEXT_COLUMN_FAMILY_NAME,
Arrays.asList(eq(ModelConstants.CUSTOMER_TENANT_ID_PROPERTY, tenantId)),
- pageLink);
+ pageLink);
log.trace("Found customers [{}] by tenantId [{}] and pageLink [{}]", customerEntities, tenantId, pageLink);
return DaoUtil.convertDataList(customerEntities);
}
@@ -67,7 +69,7 @@ public class CassandraCustomerDao extends CassandraAbstractSearchTextDao {
* @param customer the customer object
* @return saved customer object
*/
- Customer save(Customer customer);
+ Customer save(TenantId tenantId, Customer customer);
/**
* Find customers by tenant id and page link.
diff --git a/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerService.java b/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerService.java
index 4b702913b4..7471c545f0 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerService.java
@@ -26,15 +26,15 @@ import java.util.Optional;
public interface CustomerService {
- Customer findCustomerById(CustomerId customerId);
+ Customer findCustomerById(TenantId tenantId, CustomerId customerId);
Optional findCustomerByTenantIdAndTitle(TenantId tenantId, String title);
- ListenableFuture findCustomerByIdAsync(CustomerId customerId);
+ ListenableFuture findCustomerByIdAsync(TenantId tenantId, CustomerId customerId);
Customer saveCustomer(Customer customer);
- void deleteCustomer(CustomerId customerId);
+ void deleteCustomer(TenantId tenantId, CustomerId customerId);
Customer findOrCreatePublicCustomer(TenantId tenantId);
diff --git a/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java
index a9b8bfe181..1bbf273357 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerServiceImpl.java
@@ -77,10 +77,10 @@ public class CustomerServiceImpl extends AbstractEntityService implements Custom
private DashboardService dashboardService;
@Override
- public Customer findCustomerById(CustomerId customerId) {
+ public Customer findCustomerById(TenantId tenantId, CustomerId customerId) {
log.trace("Executing findCustomerById [{}]", customerId);
Validator.validateId(customerId, INCORRECT_CUSTOMER_ID + customerId);
- return customerDao.findById(customerId.getId());
+ return customerDao.findById(tenantId, customerId.getId());
}
@Override
@@ -91,36 +91,36 @@ public class CustomerServiceImpl extends AbstractEntityService implements Custom
}
@Override
- public ListenableFuture findCustomerByIdAsync(CustomerId customerId) {
+ public ListenableFuture findCustomerByIdAsync(TenantId tenantId, CustomerId customerId) {
log.trace("Executing findCustomerByIdAsync [{}]", customerId);
validateId(customerId, INCORRECT_CUSTOMER_ID + customerId);
- return customerDao.findByIdAsync(customerId.getId());
+ return customerDao.findByIdAsync(tenantId, customerId.getId());
}
@Override
public Customer saveCustomer(Customer customer) {
log.trace("Executing saveCustomer [{}]", customer);
- customerValidator.validate(customer);
- Customer savedCustomer = customerDao.save(customer);
- dashboardService.updateCustomerDashboards(savedCustomer.getId());
+ customerValidator.validate(customer, Customer::getTenantId);
+ Customer savedCustomer = customerDao.save(customer.getTenantId(), customer);
+ dashboardService.updateCustomerDashboards(savedCustomer.getTenantId(), savedCustomer.getId());
return savedCustomer;
}
@Override
- public void deleteCustomer(CustomerId customerId) {
+ public void deleteCustomer(TenantId tenantId, CustomerId customerId) {
log.trace("Executing deleteCustomer [{}]", customerId);
Validator.validateId(customerId, INCORRECT_CUSTOMER_ID + customerId);
- Customer customer = findCustomerById(customerId);
+ Customer customer = findCustomerById(tenantId, customerId);
if (customer == null) {
throw new IncorrectParameterException("Unable to delete non-existent customer.");
}
- dashboardService.unassignCustomerDashboards(customerId);
+ dashboardService.unassignCustomerDashboards(tenantId, customerId);
entityViewService.unassignCustomerEntityViews(customer.getTenantId(), customerId);
assetService.unassignCustomerAssets(customer.getTenantId(), customerId);
deviceService.unassignCustomerDevices(customer.getTenantId(), customerId);
userService.deleteCustomerUsers(customer.getTenantId(), customerId);
- deleteEntityRelations(customerId);
- customerDao.removeById(customerId.getId());
+ deleteEntityRelations(tenantId, customerId);
+ customerDao.removeById(tenantId, customerId.getId());
}
@Override
@@ -139,7 +139,7 @@ public class CustomerServiceImpl extends AbstractEntityService implements Custom
} catch (IOException e) {
throw new IncorrectParameterException("Unable to create public customer.", e);
}
- return customerDao.save(publicCustomer);
+ return customerDao.save(tenantId, publicCustomer);
}
}
@@ -156,14 +156,14 @@ public class CustomerServiceImpl extends AbstractEntityService implements Custom
public void deleteCustomersByTenantId(TenantId tenantId) {
log.trace("Executing deleteCustomersByTenantId, tenantId [{}]", tenantId);
Validator.validateId(tenantId, "Incorrect tenantId " + tenantId);
- customersByTenantRemover.removeEntities(tenantId);
+ customersByTenantRemover.removeEntities(tenantId, tenantId);
}
private DataValidator customerValidator =
new DataValidator() {
@Override
- protected void validateCreate(Customer customer) {
+ protected void validateCreate(TenantId tenantId, Customer customer) {
customerDao.findCustomersByTenantIdAndTitle(customer.getTenantId().getId(), customer.getTitle()).ifPresent(
c -> {
throw new DataValidationException("Customer with such title already exists!");
@@ -172,7 +172,7 @@ public class CustomerServiceImpl extends AbstractEntityService implements Custom
}
@Override
- protected void validateUpdate(Customer customer) {
+ protected void validateUpdate(TenantId tenantId, Customer customer) {
customerDao.findCustomersByTenantIdAndTitle(customer.getTenantId().getId(), customer.getTitle()).ifPresent(
c -> {
if (!c.getId().equals(customer.getId())) {
@@ -183,7 +183,7 @@ public class CustomerServiceImpl extends AbstractEntityService implements Custom
}
@Override
- protected void validateDataImpl(Customer customer) {
+ protected void validateDataImpl(TenantId tenantId, Customer customer) {
if (StringUtils.isEmpty(customer.getTitle())) {
throw new DataValidationException("Customer title should be specified!");
}
@@ -196,7 +196,7 @@ public class CustomerServiceImpl extends AbstractEntityService implements Custom
if (customer.getTenantId() == null) {
throw new DataValidationException("Customer should be assigned to tenant!");
} else {
- Tenant tenant = tenantDao.findById(customer.getTenantId().getId());
+ Tenant tenant = tenantDao.findById(tenantId, customer.getTenantId().getId());
if (tenant == null) {
throw new DataValidationException("Customer is referencing to non-existent tenant!");
}
@@ -208,13 +208,13 @@ public class CustomerServiceImpl extends AbstractEntityService implements Custom
new PaginatedRemover() {
@Override
- protected List findEntities(TenantId id, TextPageLink pageLink) {
+ protected List