diff --git a/application/pom.xml b/application/pom.xml
index 9364cfc16a..a4602d4b5d 100644
--- a/application/pom.xml
+++ b/application/pom.xml
@@ -241,6 +241,21 @@
mockito-all
test
+
+ org.dbunit
+ dbunit
+ test
+
+
+ com.github.springtestdbunit
+ spring-test-dbunit
+ test
+
+
+ ru.yandex.qatools.embed
+ postgresql-embedded
+ test
+
diff --git a/application/src/main/java/org/thingsboard/server/ThingsboardServerApplication.java b/application/src/main/java/org/thingsboard/server/ThingsboardServerApplication.java
index 4d3c8e9161..2994810d7c 100644
--- a/application/src/main/java/org/thingsboard/server/ThingsboardServerApplication.java
+++ b/application/src/main/java/org/thingsboard/server/ThingsboardServerApplication.java
@@ -16,15 +16,13 @@
package org.thingsboard.server;
import org.springframework.boot.SpringApplication;
-import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
-import org.springframework.boot.autoconfigure.SpringBootApplication;
+import org.springframework.boot.SpringBootConfiguration;
import org.springframework.context.annotation.ComponentScan;
import springfox.documentation.swagger2.annotations.EnableSwagger2;
import java.util.Arrays;
-@EnableAutoConfiguration
-@SpringBootApplication
+@SpringBootConfiguration
@EnableSwagger2
@ComponentScan({"org.thingsboard.server"})
public class ThingsboardServerApplication {
diff --git a/application/src/main/java/org/thingsboard/server/actors/plugin/PluginProcessingContext.java b/application/src/main/java/org/thingsboard/server/actors/plugin/PluginProcessingContext.java
index fcde477821..0e2ffcb792 100644
--- a/application/src/main/java/org/thingsboard/server/actors/plugin/PluginProcessingContext.java
+++ b/application/src/main/java/org/thingsboard/server/actors/plugin/PluginProcessingContext.java
@@ -15,15 +15,7 @@
*/
package org.thingsboard.server.actors.plugin;
-import java.io.IOException;
-import java.util.*;
-import java.util.concurrent.Executor;
-import java.util.concurrent.Executors;
-import java.util.stream.Collectors;
-
-import com.datastax.driver.core.ResultSet;
-import com.datastax.driver.core.ResultSetFuture;
-import com.datastax.driver.core.Row;
+import akka.actor.ActorRef;
import com.google.common.base.Function;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
@@ -45,8 +37,8 @@ import org.thingsboard.server.common.data.rule.RuleMetaData;
import org.thingsboard.server.common.msg.cluster.ServerAddress;
import org.thingsboard.server.extensions.api.device.DeviceAttributesEventNotificationMsg;
import org.thingsboard.server.extensions.api.plugins.PluginApiCallSecurityContext;
-import org.thingsboard.server.extensions.api.plugins.PluginContext;
import org.thingsboard.server.extensions.api.plugins.PluginCallback;
+import org.thingsboard.server.extensions.api.plugins.PluginContext;
import org.thingsboard.server.extensions.api.plugins.msg.PluginToRuleMsg;
import org.thingsboard.server.extensions.api.plugins.msg.TimeoutMsg;
import org.thingsboard.server.extensions.api.plugins.msg.ToDeviceRpcRequest;
@@ -55,9 +47,12 @@ import org.thingsboard.server.extensions.api.plugins.rpc.RpcMsg;
import org.thingsboard.server.extensions.api.plugins.ws.PluginWebsocketSessionRef;
import org.thingsboard.server.extensions.api.plugins.ws.msg.PluginWebsocketMsg;
-import akka.actor.ActorRef;
-
import javax.annotation.Nullable;
+import java.io.IOException;
+import java.util.*;
+import java.util.concurrent.Executor;
+import java.util.concurrent.Executors;
+import java.util.stream.Collectors;
@Slf4j
public final class PluginProcessingContext implements PluginContext {
@@ -95,8 +90,8 @@ public final class PluginProcessingContext implements PluginContext {
@Override
public void saveAttributes(final TenantId tenantId, final EntityId entityId, final String scope, final List attributes, final PluginCallback callback) {
validate(entityId, new ValidationCallback(callback, ctx -> {
- ListenableFuture> rsListFuture = pluginCtx.attributesService.save(entityId, scope, attributes);
- Futures.addCallback(rsListFuture, getListCallback(callback, v -> {
+ ListenableFuture> futures = pluginCtx.attributesService.save(entityId, scope, attributes);
+ Futures.addCallback(futures, getListCallback(callback, v -> {
if (entityId.getEntityType() == EntityType.DEVICE) {
onDeviceAttributesChanged(tenantId, new DeviceId(entityId.getId()), scope, attributes);
}
@@ -108,8 +103,8 @@ public final class PluginProcessingContext implements PluginContext {
@Override
public void removeAttributes(final TenantId tenantId, final EntityId entityId, final String scope, final List keys, final PluginCallback callback) {
validate(entityId, new ValidationCallback(callback, ctx -> {
- ListenableFuture> future = pluginCtx.attributesService.removeAll(entityId, scope, keys);
- Futures.addCallback(future, getCallback(callback, v -> null), executor);
+ ListenableFuture> futures = pluginCtx.attributesService.removeAll(entityId, scope, keys);
+ Futures.addCallback(futures, getCallback(callback, v -> null), executor);
if (entityId.getEntityType() == EntityType.DEVICE) {
onDeviceAttributesDeleted(tenantId, new DeviceId(entityId.getId()), keys.stream().map(key -> new AttributeKey(scope, key)).collect(Collectors.toSet()));
}
@@ -161,7 +156,7 @@ public final class PluginProcessingContext implements PluginContext {
@Override
public void saveTsData(final EntityId entityId, final TsKvEntry entry, final PluginCallback callback) {
validate(entityId, new ValidationCallback(callback, ctx -> {
- ListenableFuture> rsListFuture = pluginCtx.tsService.save(entityId, entry);
+ ListenableFuture> rsListFuture = pluginCtx.tsService.save(entityId, entry);
Futures.addCallback(rsListFuture, getListCallback(callback, v -> null), executor);
}));
}
@@ -174,7 +169,7 @@ public final class PluginProcessingContext implements PluginContext {
@Override
public void saveTsData(final EntityId entityId, final List entries, long ttl, final PluginCallback callback) {
validate(entityId, new ValidationCallback(callback, ctx -> {
- ListenableFuture> rsListFuture = pluginCtx.tsService.save(entityId, entries, ttl);
+ ListenableFuture> rsListFuture = pluginCtx.tsService.save(entityId, entries, ttl);
Futures.addCallback(rsListFuture, getListCallback(callback, v -> null), executor);
}));
}
@@ -191,26 +186,16 @@ public final class PluginProcessingContext implements PluginContext {
@Override
public void loadLatestTimeseries(final EntityId entityId, final PluginCallback> callback) {
validate(entityId, new ValidationCallback(callback, ctx -> {
- ResultSetFuture future = pluginCtx.tsService.findAllLatest(entityId);
- Futures.addCallback(future, getCallback(callback, pluginCtx.tsService::convertResultSetToTsKvEntryList), executor);
+ ListenableFuture> future = pluginCtx.tsService.findAllLatest(entityId);
+ Futures.addCallback(future, getCallback(callback, v -> v), executor);
}));
}
@Override
public void loadLatestTimeseries(final EntityId entityId, final Collection keys, final PluginCallback> callback) {
validate(entityId, new ValidationCallback(callback, ctx -> {
- ListenableFuture> rsListFuture = pluginCtx.tsService.findLatest(entityId, keys);
- Futures.addCallback(rsListFuture, getListCallback(callback, rsList ->
- {
- List result = new ArrayList<>();
- for (ResultSet rs : rsList) {
- Row row = rs.one();
- if (row != null) {
- result.add(pluginCtx.tsService.convertResultToTsKvEntry(row));
- }
- }
- return result;
- }), executor);
+ ListenableFuture> rsListFuture = pluginCtx.tsService.findLatest(entityId, keys);
+ Futures.addCallback(rsListFuture, getCallback(callback, v -> v), executor);
}));
}
@@ -237,10 +222,10 @@ public final class PluginProcessingContext implements PluginContext {
pluginCtx.toDeviceActor(DeviceAttributesEventNotificationMsg.onUpdate(tenantId, deviceId, scope, values));
}
- private FutureCallback> getListCallback(final PluginCallback callback, Function, T> transformer) {
- return new FutureCallback>() {
+ private FutureCallback> getListCallback(final PluginCallback callback, Function, R> transformer) {
+ return new FutureCallback>() {
@Override
- public void onSuccess(@Nullable List result) {
+ public void onSuccess(@Nullable List result) {
pluginCtx.self().tell(PluginCallbackMessage.onSuccess(callback, transformer.apply(result)), ActorRef.noSender());
}
diff --git a/application/src/main/java/org/thingsboard/server/controller/BaseController.java b/application/src/main/java/org/thingsboard/server/controller/BaseController.java
index c63dda7b3a..1040f3a892 100644
--- a/application/src/main/java/org/thingsboard/server/controller/BaseController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/BaseController.java
@@ -15,8 +15,6 @@
*/
package org.thingsboard.server.controller;
-import com.fasterxml.jackson.databind.JsonNode;
-import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired;
diff --git a/application/src/main/java/org/thingsboard/server/controller/CustomerController.java b/application/src/main/java/org/thingsboard/server/controller/CustomerController.java
index 51bc4527bd..091eb0edce 100644
--- a/application/src/main/java/org/thingsboard/server/controller/CustomerController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/CustomerController.java
@@ -67,6 +67,20 @@ public class CustomerController extends BaseController {
}
}
+ @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')")
+ @RequestMapping(value = "/customer/{customerId}/title", method = RequestMethod.GET, produces = "application/text")
+ @ResponseBody
+ public String getCustomerTitleById(@PathVariable("customerId") String strCustomerId) throws ThingsboardException {
+ checkParameter("customerId", strCustomerId);
+ try {
+ CustomerId customerId = new CustomerId(toUUID(strCustomerId));
+ Customer customer = checkCustomerId(customerId);
+ return customer.getTitle();
+ } catch (Exception e) {
+ throw handleException(e);
+ }
+ }
+
@PreAuthorize("hasAuthority('TENANT_ADMIN')")
@RequestMapping(value = "/customer", method = RequestMethod.POST)
@ResponseBody
diff --git a/application/src/main/java/org/thingsboard/server/controller/PluginController.java b/application/src/main/java/org/thingsboard/server/controller/PluginController.java
index 513266400e..5dd14ccf15 100644
--- a/application/src/main/java/org/thingsboard/server/controller/PluginController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/PluginController.java
@@ -25,7 +25,6 @@ import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.data.plugin.PluginMetaData;
import org.thingsboard.server.common.data.security.Authority;
-import org.thingsboard.server.common.data.widget.WidgetsBundle;
import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.exception.ThingsboardException;
diff --git a/application/src/main/java/org/thingsboard/server/service/mail/DefaultMailService.java b/application/src/main/java/org/thingsboard/server/service/mail/DefaultMailService.java
index d80b3b719f..f4dfb096e3 100644
--- a/application/src/main/java/org/thingsboard/server/service/mail/DefaultMailService.java
+++ b/application/src/main/java/org/thingsboard/server/service/mail/DefaultMailService.java
@@ -15,34 +15,31 @@
*/
package org.thingsboard.server.service.mail;
-import java.util.HashMap;
-import java.util.Locale;
-import java.util.Map;
-import java.util.Properties;
-
-import javax.annotation.PostConstruct;
-import javax.mail.internet.MimeMessage;
-
+import com.fasterxml.jackson.databind.JsonNode;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.apache.velocity.app.VelocityEngine;
-import org.springframework.beans.factory.annotation.Qualifier;
-import org.thingsboard.server.exception.ThingsboardErrorCode;
-import org.thingsboard.server.exception.ThingsboardException;
-import org.thingsboard.server.common.data.AdminSettings;
-import org.thingsboard.server.dao.settings.AdminSettingsService;
-import org.thingsboard.server.dao.exception.IncorrectParameterException;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.MessageSource;
import org.springframework.core.NestedRuntimeException;
import org.springframework.mail.javamail.JavaMailSenderImpl;
import org.springframework.mail.javamail.MimeMessageHelper;
import org.springframework.stereotype.Service;
import org.springframework.ui.velocity.VelocityEngineUtils;
+import org.thingsboard.server.common.data.AdminSettings;
+import org.thingsboard.server.dao.exception.IncorrectParameterException;
+import org.thingsboard.server.dao.settings.AdminSettingsService;
+import org.thingsboard.server.exception.ThingsboardErrorCode;
+import org.thingsboard.server.exception.ThingsboardException;
+
+import javax.annotation.PostConstruct;
+import javax.mail.internet.MimeMessage;
+import java.util.HashMap;
+import java.util.Locale;
+import java.util.Map;
+import java.util.Properties;
-import com.fasterxml.jackson.databind.JsonNode;
@Service
@Slf4j
public class DefaultMailService implements MailService {
@@ -69,9 +66,13 @@ public class DefaultMailService implements MailService {
@Override
public void updateMailConfiguration() {
AdminSettings settings = adminSettingsService.findAdminSettingsByKey("mail");
- JsonNode jsonConfig = settings.getJsonValue();
- mailSender = createMailSender(jsonConfig);
- mailFrom = jsonConfig.get("mailFrom").asText();
+ if (settings != null) {
+ JsonNode jsonConfig = settings.getJsonValue();
+ mailSender = createMailSender(jsonConfig);
+ mailFrom = jsonConfig.get("mailFrom").asText();
+ } else {
+ throw new IncorrectParameterException("Failed to date mail configuration. Settings not found!");
+ }
}
private JavaMailSenderImpl createMailSender(JsonNode jsonConfig) {
diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml
index 4b488af7bb..b443dcdb5c 100644
--- a/application/src/main/resources/thingsboard.yml
+++ b/application/src/main/resources/thingsboard.yml
@@ -107,6 +107,7 @@ coap:
# Cassandra driver configuration parameters
cassandra:
+ enabled: "${CASSANDRA_ENABLED:false}"
# Thingsboard cluster name
cluster_name: "${CASSANDRA_CLUSTER_NAME:Thingsboard Cluster}"
# Thingsboard keyspace name
@@ -152,7 +153,7 @@ cassandra:
# Specify partitioning size for timestamp key-value storage. Example MINUTES, HOURS, DAYS, MONTHS
ts_key_value_partitioning: "${TS_KV_PARTITIONING:MONTHS}"
# Specify max data points per request
- min_aggregation_step_ms: "${TS_KV_MIN_AGGREGATION_STEP_MS:100}"
+ min_aggregation_step_ms: "${TS_KV_MIN_AGGREGATION_STEP_MS:1000}"
# Actor system parameters
actors:
@@ -222,3 +223,24 @@ spring.mvc.cors:
allowed-headers: "*"
max-age: "1800"
allow-credentials: "true"
+
+# SQL DAO Configuration
+sql:
+ enabled: "${SQL_ENABLED:true}"
+
+spring:
+ data:
+ jpa:
+ repositories:
+ enabled: "true"
+ jpa:
+ show-sql: "false"
+ generate-ddl: "true"
+ database-platform: "org.hibernate.dialect.PostgreSQLDialect"
+ hibernate:
+ ddl-auto: "validate"
+ datasource:
+ driverClassName: "${SPRING_DRIVER_CLASS_NAME:org.postgresql.Driver}"
+ url: "${SPRING_DATASOURCE_URL:jdbc:postgresql://localhost:5432/thingsboard}"
+ username: "${SPRING_DATASOURCE_USERNAME:postgres}"
+ password: "${SPRING_DATASOURCE_PASSWORD:postgres}"
\ No newline at end of file
diff --git a/application/src/test/java/org/thingsboard/server/ThingsboardApplicationTests.java b/application/src/test/java/org/thingsboard/server/ThingsboardApplicationTests.java
deleted file mode 100644
index 80d0b62205..0000000000
--- a/application/src/test/java/org/thingsboard/server/ThingsboardApplicationTests.java
+++ /dev/null
@@ -1,49 +0,0 @@
-/**
- * Copyright © 2016-2017 The Thingsboard Authors
- *
- * Licensed under the Apache License, Version 2.0 (the "License");
- * you may not use this file except in compliance with the License.
- * You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package org.thingsboard.server;
-
-import org.junit.Test;
-import org.junit.runner.RunWith;
-import org.springframework.test.context.web.WebAppConfiguration;
-import org.springframework.boot.test.IntegrationTest;
-import org.springframework.boot.test.SpringApplicationConfiguration;
-import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
-
-@RunWith(SpringJUnit4ClassRunner.class)
-@SpringApplicationConfiguration(classes = ThingsboardServerApplication.class)
-@WebAppConfiguration
-@IntegrationTest("server.port:0")
-public class ThingsboardApplicationTests {
-
- @Test
- public void contextLoads() {
- String test = "[ \n" +
- " {\n" +
- " \"key\": \"name\",\n" +
- "\t\"type\": \"text\" \n" +
- " },\n" +
- " {\n" +
- "\t\"key\": \"name2\",\n" +
- "\t\"type\": \"color\"\n" +
- " },\n" +
- " {\n" +
- "\t\"key\": \"name3\",\n" +
- "\t\"type\": \"javascript\"\n" +
- " } \n" +
- "]";
- }
-
-}
diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java
index b85f159b21..052b61ac6b 100644
--- a/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java
+++ b/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java
@@ -28,14 +28,9 @@ import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.runner.RunWith;
-import org.mockito.Mockito;
-import org.mockito.invocation.InvocationOnMock;
-import org.mockito.stubbing.Answer;
import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
-import org.springframework.boot.test.IntegrationTest;
-import org.springframework.boot.test.SpringApplicationContextLoader;
-import org.springframework.context.annotation.Bean;
+import org.springframework.boot.test.context.SpringBootContextLoader;
+import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.context.annotation.Configuration;
import org.springframework.http.HttpHeaders;
@@ -48,7 +43,7 @@ import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.ActiveProfiles;
import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.TestPropertySource;
-import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
+import org.springframework.test.context.junit4.SpringRunner;
import org.springframework.test.context.web.WebAppConfiguration;
import org.springframework.test.web.servlet.MockMvc;
import org.springframework.test.web.servlet.ResultActions;
@@ -66,11 +61,9 @@ import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.config.ThingsboardSecurityConfiguration;
-import org.thingsboard.server.exception.ThingsboardException;
-import org.thingsboard.server.service.mail.MailService;
import org.thingsboard.server.service.mail.TestMailService;
-import org.thingsboard.server.service.security.auth.rest.LoginRequest;
import org.thingsboard.server.service.security.auth.jwt.RefreshTokenRequest;
+import org.thingsboard.server.service.security.auth.rest.LoginRequest;
import java.io.IOException;
import java.nio.charset.Charset;
@@ -81,21 +74,18 @@ import java.util.List;
import static org.springframework.security.test.web.servlet.setup.SecurityMockMvcConfigurers.springSecurity;
import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.*;
-import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.header;
-import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath;
-import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
+import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.*;
import static org.springframework.test.web.servlet.setup.MockMvcBuilders.webAppContextSetup;
@ActiveProfiles("test")
-@RunWith(SpringJUnit4ClassRunner.class)
-@ContextConfiguration(classes=AbstractControllerTest.class, loader=SpringApplicationContextLoader.class)
-@TestPropertySource(locations = {"classpath:cassandra-test.properties", "classpath:thingsboard-test.properties"})
+@RunWith(SpringRunner.class)
+@ContextConfiguration(classes = AbstractControllerTest.class, loader = SpringBootContextLoader.class)
+@TestPropertySource(locations = {"classpath:cassandra-test.properties", "classpath:application-test.properties", "classpath:nosql-test.properties"})
@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS)
@Configuration
-@EnableAutoConfiguration
@ComponentScan({"org.thingsboard.server"})
@WebAppConfiguration
-@IntegrationTest("server.port:0")
+@SpringBootTest
public abstract class AbstractControllerTest {
protected static final String TEST_TENANT_NAME = "TEST TENANT";
@@ -113,7 +103,6 @@ public abstract class AbstractControllerTest {
MediaType.APPLICATION_JSON.getSubtype(),
Charset.forName("utf8"));
-
protected MockMvc mockMvc;
protected String token;
@@ -305,7 +294,7 @@ public abstract class AbstractControllerTest {
protected T doPost(String urlTemplate, T content, Class responseClass, String... params) throws Exception {
return readResponse(doPost(urlTemplate, content, params).andExpect(status().isOk()), responseClass);
}
-
+
protected T doDelete(String urlTemplate, Class responseClass, String... params) throws Exception {
return readResponse(doDelete(urlTemplate, params).andExpect(status().isOk()), responseClass);
}
@@ -364,8 +353,8 @@ public abstract class AbstractControllerTest {
ObjectMapper mapper = new ObjectMapper();
return mapper.readerFor(type).readValue(content);
}
-
- class IdComparator> implements Comparator {
+
+ public class IdComparator> implements Comparator {
@Override
public int compare(D o1, D o2) {
return o1.getId().getId().compareTo(o2.getId().getId());
diff --git a/application/src/test/java/org/thingsboard/server/controller/AssetControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/AssetControllerTest.java
index 10a2ff7c8f..d15023b91f 100644
--- a/application/src/test/java/org/thingsboard/server/controller/AssetControllerTest.java
+++ b/application/src/test/java/org/thingsboard/server/controller/AssetControllerTest.java
@@ -28,7 +28,6 @@ import org.thingsboard.server.common.data.*;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.asset.TenantAssetType;
import org.thingsboard.server.common.data.id.CustomerId;
-import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.page.TextPageData;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.security.Authority;
diff --git a/application/src/test/java/org/thingsboard/server/controller/ControllerTestSuite.java b/application/src/test/java/org/thingsboard/server/controller/ControllerTestSuite.java
index 21093b3a0b..98e558fc0a 100644
--- a/application/src/test/java/org/thingsboard/server/controller/ControllerTestSuite.java
+++ b/application/src/test/java/org/thingsboard/server/controller/ControllerTestSuite.java
@@ -18,22 +18,22 @@ package org.thingsboard.server.controller;
import org.cassandraunit.dataset.cql.ClassPathCQLDataSet;
import org.junit.ClassRule;
import org.junit.extensions.cpsuite.ClasspathSuite;
-import org.junit.extensions.cpsuite.ClasspathSuite.ClassnameFilters;
import org.junit.runner.RunWith;
import org.thingsboard.server.dao.CustomCassandraCQLUnit;
import java.util.Arrays;
@RunWith(ClasspathSuite.class)
-@ClassnameFilters({"org.thingsboard.server.controller.*Test"})
+@ClasspathSuite.ClassnameFilters({
+ "org.thingsboard.server.controller.*Test"})
public class ControllerTestSuite {
@ClassRule
public static CustomCassandraCQLUnit cassandraUnit =
- new CustomCassandraCQLUnit(Arrays.asList(
- new ClassPathCQLDataSet("schema.cql", false, false),
- new ClassPathCQLDataSet("system-data.cql", false, false),
- new ClassPathCQLDataSet("system-test.cql", false, false)),
+ new CustomCassandraCQLUnit(
+ Arrays.asList(
+ new ClassPathCQLDataSet("cassandra/schema.cql", false, false),
+ new ClassPathCQLDataSet("cassandra/system-data.cql", false, false),
+ new ClassPathCQLDataSet("system-test.cql", false, false)),
"cassandra-test.yaml", 30000l);
-
}
diff --git a/application/src/test/java/org/thingsboard/server/controller/DashboardControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/DashboardControllerTest.java
index e6eb243db2..dda68a64ef 100644
--- a/application/src/test/java/org/thingsboard/server/controller/DashboardControllerTest.java
+++ b/application/src/test/java/org/thingsboard/server/controller/DashboardControllerTest.java
@@ -23,6 +23,7 @@ import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
+import com.datastax.driver.core.utils.UUIDs;
import org.apache.commons.lang3.RandomStringUtils;
import org.thingsboard.server.common.data.*;
import org.thingsboard.server.common.data.id.CustomerId;
@@ -35,7 +36,6 @@ import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
-import com.datastax.driver.core.utils.UUIDs;
import com.fasterxml.jackson.core.type.TypeReference;
public class DashboardControllerTest extends AbstractControllerTest {
@@ -155,7 +155,7 @@ public class DashboardControllerTest extends AbstractControllerTest {
dashboard.setTitle("My dashboard");
Dashboard savedDashboard = doPost("/api/dashboard", dashboard, Dashboard.class);
- doPost("/api/customer/" + UUIDs.timeBased().toString()
+ doPost("/api/customer/" + UUIDs.timeBased().toString()
+ "/dashboard/" + savedDashboard.getId().getId().toString())
.andExpect(status().isNotFound());
}
diff --git a/application/src/test/java/org/thingsboard/server/controller/DeviceControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/DeviceControllerTest.java
index 5d40a79680..26eaa5dbb7 100644
--- a/application/src/test/java/org/thingsboard/server/controller/DeviceControllerTest.java
+++ b/application/src/test/java/org/thingsboard/server/controller/DeviceControllerTest.java
@@ -23,6 +23,7 @@ import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
+import com.datastax.driver.core.utils.UUIDs;
import org.apache.commons.lang3.RandomStringUtils;
import org.thingsboard.server.common.data.*;
import org.thingsboard.server.common.data.id.CustomerId;
@@ -39,7 +40,6 @@ import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;
-import com.datastax.driver.core.utils.UUIDs;
import com.fasterxml.jackson.core.type.TypeReference;
public class DeviceControllerTest extends AbstractControllerTest {
@@ -215,7 +215,7 @@ public class DeviceControllerTest extends AbstractControllerTest {
device.setType("default");
Device savedDevice = doPost("/api/device", device, Device.class);
- doPost("/api/customer/" + UUIDs.timeBased().toString()
+ doPost("/api/customer/" + UUIDs.timeBased().toString()
+ "/device/" + savedDevice.getId().getId().toString())
.andExpect(status().isNotFound());
}
diff --git a/application/src/test/java/org/thingsboard/server/mqtt/AbstractFeatureIntegrationTest.java b/application/src/test/java/org/thingsboard/server/mqtt/AbstractFeatureIntegrationTest.java
deleted file mode 100644
index 0f449d3310..0000000000
--- a/application/src/test/java/org/thingsboard/server/mqtt/AbstractFeatureIntegrationTest.java
+++ /dev/null
@@ -1,73 +0,0 @@
-/**
- * Copyright © 2016-2017 The Thingsboard Authors
- *
- * Licensed under the Apache License, Version 2.0 (the "License");
- * you may not use this file except in compliance with the License.
- * You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package org.thingsboard.server.mqtt;
-
-import org.junit.runner.RunWith;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
-import org.springframework.boot.test.IntegrationTest;
-import org.springframework.boot.test.SpringApplicationContextLoader;
-import org.springframework.context.annotation.ComponentScan;
-import org.springframework.context.annotation.Configuration;
-import org.springframework.http.converter.HttpMessageConverter;
-import org.springframework.http.converter.json.MappingJackson2HttpMessageConverter;
-import org.springframework.test.annotation.DirtiesContext;
-import org.springframework.test.context.ActiveProfiles;
-import org.springframework.test.context.ContextConfiguration;
-import org.springframework.test.context.TestPropertySource;
-import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
-import org.springframework.test.context.web.WebAppConfiguration;
-import org.springframework.web.context.WebApplicationContext;
-import org.thingsboard.server.mqtt.telemetry.MqttTelemetryIntergrationTest;
-
-import java.util.Arrays;
-
-import static org.junit.Assert.assertNotNull;
-
-/**
- * @author Valerii Sosliuk
- */
-@ActiveProfiles("default")
-@RunWith(SpringJUnit4ClassRunner.class)
-@ContextConfiguration(classes= MqttTelemetryIntergrationTest.class, loader=SpringApplicationContextLoader.class)
-@TestPropertySource(locations = {"classpath:cassandra-test.properties", "classpath:thingsboard-test.properties"})
-@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS)
-@Configuration
-@EnableAutoConfiguration
-@ComponentScan({"org.thingsboard.server"})
-@WebAppConfiguration
-@IntegrationTest("server.port:8080")
-public class AbstractFeatureIntegrationTest {
-
- @SuppressWarnings("rawtypes")
- private HttpMessageConverter mappingJackson2HttpMessageConverter;
-
- @Autowired
- private WebApplicationContext webApplicationContext;
-
- @Autowired
- void setConverters(HttpMessageConverter>[] converters) {
-
- this.mappingJackson2HttpMessageConverter = Arrays.stream(converters)
- .filter(hmc -> hmc instanceof MappingJackson2HttpMessageConverter)
- .findAny()
- .get();
-
- assertNotNull("the JSON message converter must not be null",
- this.mappingJackson2HttpMessageConverter);
- }
-
-}
diff --git a/application/src/test/java/org/thingsboard/server/mqtt/MqttSuite.java b/application/src/test/java/org/thingsboard/server/mqtt/MqttTestSuite.java
similarity index 73%
rename from application/src/test/java/org/thingsboard/server/mqtt/MqttSuite.java
rename to application/src/test/java/org/thingsboard/server/mqtt/MqttTestSuite.java
index a1e415dfa9..cf3bc71532 100644
--- a/application/src/test/java/org/thingsboard/server/mqtt/MqttSuite.java
+++ b/application/src/test/java/org/thingsboard/server/mqtt/MqttTestSuite.java
@@ -23,19 +23,16 @@ import org.thingsboard.server.dao.CustomCassandraCQLUnit;
import java.util.Arrays;
-/**
- * @author Valerii Sosliuk
- */
@RunWith(ClasspathSuite.class)
-@ClasspathSuite.ClassnameFilters({"org.thingsboard.server.mqtt.*.*Test"})
-public class MqttSuite {
+@ClasspathSuite.ClassnameFilters({
+ "org.thingsboard.server.mqtt.*.*Test"})
+public class MqttTestSuite {
@ClassRule
public static CustomCassandraCQLUnit cassandraUnit =
new CustomCassandraCQLUnit(
- Arrays.asList(new ClassPathCQLDataSet("schema.cql", false, false),
- new ClassPathCQLDataSet("system-data.cql", false, false),
- new ClassPathCQLDataSet("demo-data.cql", false, false)),
+ Arrays.asList(
+ new ClassPathCQLDataSet("cassandra/schema.cql", false, false),
+ new ClassPathCQLDataSet("cassandra/system-data.cql", false, false)),
"cassandra-test.yaml", 30000l);
-
}
diff --git a/application/src/test/java/org/thingsboard/server/mqtt/rpc/MqttServerSideRpcIntegrationTest.java b/application/src/test/java/org/thingsboard/server/mqtt/rpc/MqttServerSideRpcIntegrationTest.java
index ac47d7ed64..a3a2355387 100644
--- a/application/src/test/java/org/thingsboard/server/mqtt/rpc/MqttServerSideRpcIntegrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/mqtt/rpc/MqttServerSideRpcIntegrationTest.java
@@ -17,47 +17,67 @@ package org.thingsboard.server.mqtt.rpc;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.paho.client.mqttv3.*;
-import org.junit.Assert;
-import org.junit.Before;
-import org.junit.Ignore;
-import org.junit.Test;
+import org.junit.*;
import org.springframework.http.HttpStatus;
-import org.springframework.http.ResponseEntity;
import org.springframework.web.client.HttpClientErrorException;
-import org.thingsboard.client.tools.RestClient;
import org.thingsboard.server.common.data.Device;
+import org.thingsboard.server.common.data.Tenant;
+import org.thingsboard.server.common.data.User;
+import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.common.data.security.DeviceCredentials;
-import org.thingsboard.server.mqtt.AbstractFeatureIntegrationTest;
+import org.thingsboard.server.controller.AbstractControllerTest;
import java.util.UUID;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
+import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
/**
* @author Valerii Sosliuk
*/
@Slf4j
-public class MqttServerSideRpcIntegrationTest extends AbstractFeatureIntegrationTest {
+public class MqttServerSideRpcIntegrationTest extends AbstractControllerTest {
private static final String MQTT_URL = "tcp://localhost:1883";
- private static final String BASE_URL = "http://localhost:8080";
+ private static final String FAIL_MSG_IF_HTTP_CLIENT_ERROR_NOT_ENCOUNTERED = "HttpClientErrorException expected, but not encountered";
- private static final String USERNAME = "tenant@thingsboard.org";
- private static final String PASSWORD = "tenant";
-
- private RestClient restClient;
+ private Tenant savedTenant;
+ private User tenantAdmin;
@Before
public void beforeTest() throws Exception {
- restClient = new RestClient(BASE_URL);
- restClient.login(USERNAME, PASSWORD);
+ loginSysAdmin();
+
+ Tenant tenant = new Tenant();
+ tenant.setTitle("My tenant");
+ savedTenant = doPost("/api/tenant", tenant, Tenant.class);
+ Assert.assertNotNull(savedTenant);
+
+ tenantAdmin = new User();
+ tenantAdmin.setAuthority(Authority.TENANT_ADMIN);
+ tenantAdmin.setTenantId(savedTenant.getId());
+ tenantAdmin.setEmail("tenant2@thingsboard.org");
+ tenantAdmin.setFirstName("Joe");
+ tenantAdmin.setLastName("Downs");
+
+ createUserAndLogin(tenantAdmin, "testPassword1");
+ }
+
+ @After
+ public void afterTest() throws Exception {
+ loginSysAdmin();
+
+ doDelete("/api/tenant/" + savedTenant.getId().getId().toString())
+ .andExpect(status().isOk());
}
@Test
+ @Ignore
public void testServerMqttOneWayRpc() throws Exception {
Device device = new Device();
device.setName("Test One-Way Server-Side RPC");
+ device.setType("default");
Device savedDevice = getSavedDevice(device);
DeviceCredentials deviceCredentials = getDeviceCredentials(savedDevice);
assertEquals(savedDevice.getId(), deviceCredentials.getDeviceId());
@@ -69,22 +89,22 @@ public class MqttServerSideRpcIntegrationTest extends AbstractFeatureIntegration
MqttConnectOptions options = new MqttConnectOptions();
options.setUserName(accessToken);
- client.connect(options);
- Thread.sleep(3000);
+ client.connect(options).waitForCompletion();
client.subscribe("v1/devices/me/rpc/request/+", 1);
client.setCallback(new TestMqttCallback(client));
String setGpioRequest = "{\"method\":\"setGpio\",\"params\":{\"pin\": \"23\",\"value\": 1}}";
String deviceId = savedDevice.getId().getId().toString();
- ResponseEntity result = restClient.getRestTemplate().postForEntity(BASE_URL + "api/plugins/rpc/oneway/" + deviceId, setGpioRequest, String.class);
- Assert.assertEquals(HttpStatus.OK, result.getStatusCode());
- Assert.assertNull(result.getBody());
+ String result = doPost("api/plugins/rpc/oneway/" + deviceId, setGpioRequest, String.class);
+ Assert.assertNull(result);
}
@Test
+ @Ignore
public void testServerMqttOneWayRpcDeviceOffline() throws Exception {
Device device = new Device();
device.setName("Test One-Way Server-Side RPC Device Offline");
+ device.setType("default");
Device savedDevice = getSavedDevice(device);
DeviceCredentials deviceCredentials = getDeviceCredentials(savedDevice);
assertEquals(savedDevice.getId(), deviceCredentials.getDeviceId());
@@ -94,8 +114,8 @@ public class MqttServerSideRpcIntegrationTest extends AbstractFeatureIntegration
String setGpioRequest = "{\"method\":\"setGpio\",\"params\":{\"pin\": \"23\",\"value\": 1}}";
String deviceId = savedDevice.getId().getId().toString();
try {
- restClient.getRestTemplate().postForEntity(BASE_URL + "api/plugins/rpc/oneway/" + deviceId, setGpioRequest, String.class);
- Assert.fail("HttpClientErrorException expected, but not encountered");
+ doPost("api/plugins/rpc/oneway/" + deviceId, setGpioRequest, String.class);
+ Assert.fail(FAIL_MSG_IF_HTTP_CLIENT_ERROR_NOT_ENCOUNTERED);
} catch (HttpClientErrorException e) {
log.error(e.getMessage(), e);
Assert.assertEquals(HttpStatus.REQUEST_TIMEOUT, e.getStatusCode());
@@ -104,12 +124,13 @@ public class MqttServerSideRpcIntegrationTest extends AbstractFeatureIntegration
}
@Test
+ @Ignore
public void testServerMqttOneWayRpcDeviceDoesNotExist() throws Exception {
String setGpioRequest = "{\"method\":\"setGpio\",\"params\":{\"pin\": \"23\",\"value\": 1}}";
String nonExistentDeviceId = UUID.randomUUID().toString();
try {
- restClient.getRestTemplate().postForEntity(BASE_URL + "api/plugins/rpc/oneway/" + nonExistentDeviceId, setGpioRequest, String.class);
- Assert.fail("HttpClientErrorException expected, but not encountered");
+ doPost("api/plugins/rpc/oneway/" + nonExistentDeviceId, setGpioRequest, String.class);
+ Assert.fail(FAIL_MSG_IF_HTTP_CLIENT_ERROR_NOT_ENCOUNTERED);
} catch (HttpClientErrorException e) {
log.error(e.getMessage(), e);
Assert.assertEquals(HttpStatus.BAD_REQUEST, e.getStatusCode());
@@ -118,10 +139,11 @@ public class MqttServerSideRpcIntegrationTest extends AbstractFeatureIntegration
}
@Test
+ @Ignore
public void testServerMqttTwoWayRpc() throws Exception {
-
Device device = new Device();
device.setName("Test Two-Way Server-Side RPC");
+ device.setType("default");
Device savedDevice = getSavedDevice(device);
DeviceCredentials deviceCredentials = getDeviceCredentials(savedDevice);
assertEquals(savedDevice.getId(), deviceCredentials.getDeviceId());
@@ -133,8 +155,7 @@ public class MqttServerSideRpcIntegrationTest extends AbstractFeatureIntegration
MqttConnectOptions options = new MqttConnectOptions();
options.setUserName(accessToken);
- client.connect(options);
- Thread.sleep(3000);
+ client.connect(options).waitForCompletion();
client.subscribe("v1/devices/me/rpc/request/+", 1);
client.setCallback(new TestMqttCallback(client));
@@ -145,9 +166,11 @@ public class MqttServerSideRpcIntegrationTest extends AbstractFeatureIntegration
}
@Test
+ @Ignore
public void testServerMqttTwoWayRpcDeviceOffline() throws Exception {
Device device = new Device();
device.setName("Test Two-Way Server-Side RPC Device Offline");
+ device.setType("default");
Device savedDevice = getSavedDevice(device);
DeviceCredentials deviceCredentials = getDeviceCredentials(savedDevice);
assertEquals(savedDevice.getId(), deviceCredentials.getDeviceId());
@@ -157,8 +180,8 @@ public class MqttServerSideRpcIntegrationTest extends AbstractFeatureIntegration
String setGpioRequest = "{\"method\":\"setGpio\",\"params\":{\"pin\": \"23\",\"value\": 1}}";
String deviceId = savedDevice.getId().getId().toString();
try {
- restClient.getRestTemplate().postForEntity(BASE_URL + "api/plugins/rpc/twoway/" + deviceId, setGpioRequest, String.class);
- Assert.fail("HttpClientErrorException expected, but not encountered");
+ doPost("api/plugins/rpc/twoway/" + deviceId, setGpioRequest, String.class);
+ Assert.fail(FAIL_MSG_IF_HTTP_CLIENT_ERROR_NOT_ENCOUNTERED);
} catch (HttpClientErrorException e) {
log.error(e.getMessage(), e);
Assert.assertEquals(HttpStatus.REQUEST_TIMEOUT, e.getStatusCode());
@@ -167,12 +190,13 @@ public class MqttServerSideRpcIntegrationTest extends AbstractFeatureIntegration
}
@Test
+ @Ignore
public void testServerMqttTwoWayRpcDeviceDoesNotExist() throws Exception {
String setGpioRequest = "{\"method\":\"setGpio\",\"params\":{\"pin\": \"23\",\"value\": 1}}";
String nonExistentDeviceId = UUID.randomUUID().toString();
try {
- restClient.getRestTemplate().postForEntity(BASE_URL + "api/plugins/rpc/oneway/" + nonExistentDeviceId, setGpioRequest, String.class);
- Assert.fail("HttpClientErrorException expected, but not encountered");
+ doPost("api/plugins/rpc/oneway/" + nonExistentDeviceId, setGpioRequest, String.class);
+ Assert.fail(FAIL_MSG_IF_HTTP_CLIENT_ERROR_NOT_ENCOUNTERED);
} catch (HttpClientErrorException e) {
log.error(e.getMessage(), e);
Assert.assertEquals(HttpStatus.BAD_REQUEST, e.getStatusCode());
@@ -180,16 +204,16 @@ public class MqttServerSideRpcIntegrationTest extends AbstractFeatureIntegration
}
}
- private Device getSavedDevice(Device device) {
- return restClient.getRestTemplate().postForEntity(BASE_URL + "/api/device", device, Device.class).getBody();
+ private Device getSavedDevice(Device device) throws Exception {
+ return doPost("/api/device", device, Device.class);
}
- private DeviceCredentials getDeviceCredentials(Device savedDevice) {
- return restClient.getRestTemplate().getForEntity(BASE_URL + "/api/device/" + savedDevice.getId().getId().toString() + "/credentials", DeviceCredentials.class).getBody();
+ private DeviceCredentials getDeviceCredentials(Device savedDevice) throws Exception {
+ return doGet("/api/device/" + savedDevice.getId().getId().toString() + "/credentials", DeviceCredentials.class);
}
- private String getStringResult(String requestData, String callType, String deviceId) {
- return restClient.getRestTemplate().postForEntity(BASE_URL + "api/plugins/rpc/" + callType + "/" + deviceId, requestData, String.class).getBody();
+ private String getStringResult(String requestData, String callType, String deviceId) throws Exception {
+ return doPost("api/plugins/rpc/" + callType + "/" + deviceId, requestData, String.class);
}
private static class TestMqttCallback implements MqttCallback {
diff --git a/application/src/test/java/org/thingsboard/server/mqtt/telemetry/MqttTelemetryIntergrationTest.java b/application/src/test/java/org/thingsboard/server/mqtt/telemetry/MqttTelemetryIntegrationTest.java
similarity index 70%
rename from application/src/test/java/org/thingsboard/server/mqtt/telemetry/MqttTelemetryIntergrationTest.java
rename to application/src/test/java/org/thingsboard/server/mqtt/telemetry/MqttTelemetryIntegrationTest.java
index d4f139c893..b0b628a8a1 100644
--- a/application/src/test/java/org/thingsboard/server/mqtt/telemetry/MqttTelemetryIntergrationTest.java
+++ b/application/src/test/java/org/thingsboard/server/mqtt/telemetry/MqttTelemetryIntegrationTest.java
@@ -20,12 +20,12 @@ import org.eclipse.paho.client.mqttv3.MqttAsyncClient;
import org.eclipse.paho.client.mqttv3.MqttConnectOptions;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.junit.Before;
+import org.junit.Ignore;
import org.junit.Test;
import org.springframework.web.util.UriComponentsBuilder;
-import org.thingsboard.client.tools.RestClient;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.security.DeviceCredentials;
-import org.thingsboard.server.mqtt.AbstractFeatureIntegrationTest;
+import org.thingsboard.server.controller.AbstractControllerTest;
import java.net.URI;
import java.util.Arrays;
@@ -39,35 +39,32 @@ import static org.junit.Assert.assertNotNull;
* @author Valerii Sosliuk
*/
@Slf4j
-public class MqttTelemetryIntergrationTest extends AbstractFeatureIntegrationTest {
+public class MqttTelemetryIntegrationTest extends AbstractControllerTest {
private static final String MQTT_URL = "tcp://localhost:1883";
- private static final String BASE_URL = "http://localhost:8080";
-
- private static final String USERNAME = "tenant@thingsboard.org";
- private static final String PASSWORD = "tenant";
private Device savedDevice;
-
private String accessToken;
- private RestClient restClient;
@Before
public void beforeTest() throws Exception {
- restClient = new RestClient(BASE_URL);
- restClient.login(USERNAME, PASSWORD);
+ loginTenantAdmin();
Device device = new Device();
device.setName("Test device");
- savedDevice = restClient.getRestTemplate().postForEntity(BASE_URL + "/api/device", device, Device.class).getBody();
+ device.setType("default");
+ savedDevice = doPost("/api/device", device, Device.class);
+
DeviceCredentials deviceCredentials =
- restClient.getRestTemplate().getForEntity(BASE_URL + "/api/device/" + savedDevice.getId().getId().toString() + "/credentials", DeviceCredentials.class).getBody();
+ doGet("/api/device/" + savedDevice.getId().getId().toString() + "/credentials", DeviceCredentials.class);
+
assertEquals(savedDevice.getId(), deviceCredentials.getDeviceId());
accessToken = deviceCredentials.getCredentialsId();
assertNotNull(accessToken);
}
@Test
+ @Ignore
public void testPushMqttRpcData() throws Exception {
String clientId = MqttAsyncClient.generateClientId();
MqttAsyncClient client = new MqttAsyncClient(MQTT_URL, clientId);
@@ -83,13 +80,13 @@ public class MqttTelemetryIntergrationTest extends AbstractFeatureIntegrationTes
String deviceId = savedDevice.getId().getId().toString();
Thread.sleep(1000);
- List keys = restClient.getRestTemplate().getForEntity(BASE_URL + "/api/plugins/telemetry/" + deviceId + "/keys/timeseries", List.class).getBody();
+ Object keys = doGet("/api/plugins/telemetry/" + deviceId + "/keys/timeseries", Object.class);
assertEquals(Arrays.asList("key1", "key2", "key3", "key4"), keys);
- UriComponentsBuilder builder = UriComponentsBuilder.fromHttpUrl(BASE_URL + "/api/plugins/telemetry/" + deviceId + "/values/timeseries")
- .queryParam("keys", String.join(",", keys));
+ UriComponentsBuilder builder = UriComponentsBuilder.fromHttpUrl("/api/plugins/telemetry/" + deviceId + "/values/timeseries")
+ .queryParam("keys", String.join(",", (CharSequence[]) keys));
URI uri = builder.build().encode().toUri();
- Map>> values = restClient.getRestTemplate().getForEntity(uri, Map.class).getBody();
+ Map>> values = doGet(uri.getPath(), Map.class);
assertEquals("value1", values.get("key1").get(0).get("value"));
assertEquals("true", values.get("key2").get(0).get("value"));
diff --git a/application/src/test/java/org/thingsboard/server/system/HttpDeviceApiTest.java b/application/src/test/java/org/thingsboard/server/system/BaseHttpDeviceApiTest.java
similarity index 97%
rename from application/src/test/java/org/thingsboard/server/system/HttpDeviceApiTest.java
rename to application/src/test/java/org/thingsboard/server/system/BaseHttpDeviceApiTest.java
index af8b50fc48..680fa9492d 100644
--- a/application/src/test/java/org/thingsboard/server/system/HttpDeviceApiTest.java
+++ b/application/src/test/java/org/thingsboard/server/system/BaseHttpDeviceApiTest.java
@@ -35,7 +35,7 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.
/**
* @author Andrew Shvayka
*/
-public class HttpDeviceApiTest extends AbstractControllerTest {
+public abstract class BaseHttpDeviceApiTest extends AbstractControllerTest {
private static final AtomicInteger idSeq = new AtomicInteger(new Random(System.currentTimeMillis()).nextInt());
diff --git a/application/src/test/java/org/thingsboard/server/system/SystemTestSuite.java b/application/src/test/java/org/thingsboard/server/system/SystemNoSqlTestSuite.java
similarity index 78%
rename from application/src/test/java/org/thingsboard/server/system/SystemTestSuite.java
rename to application/src/test/java/org/thingsboard/server/system/SystemNoSqlTestSuite.java
index f3297d6f52..54c3bc70fd 100644
--- a/application/src/test/java/org/thingsboard/server/system/SystemTestSuite.java
+++ b/application/src/test/java/org/thingsboard/server/system/SystemNoSqlTestSuite.java
@@ -27,13 +27,14 @@ import java.util.Arrays;
* @author Andrew Shvayka
*/
@RunWith(ClasspathSuite.class)
-@ClasspathSuite.ClassnameFilters({"org.thingsboard.server.system.*Test"})
-public class SystemTestSuite {
+@ClasspathSuite.ClassnameFilters({"org.thingsboard.server.system.*NoSqlTest"})
+public class SystemNoSqlTestSuite {
@ClassRule
public static CustomCassandraCQLUnit cassandraUnit =
- new CustomCassandraCQLUnit(Arrays.asList(
- new ClassPathCQLDataSet("schema.cql", false, false),
- new ClassPathCQLDataSet("system-data.cql", false, false)),
+ new CustomCassandraCQLUnit(
+ Arrays.asList(
+ new ClassPathCQLDataSet("cassandra/schema.cql", false, false),
+ new ClassPathCQLDataSet("cassandra/system-data.cql", false, false)),
"cassandra-test.yaml", 30000l);
}
diff --git a/application/src/test/java/org/thingsboard/server/system/SystemSqlTestSuite.java b/application/src/test/java/org/thingsboard/server/system/SystemSqlTestSuite.java
new file mode 100644
index 0000000000..726b2052aa
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/system/SystemSqlTestSuite.java
@@ -0,0 +1,38 @@
+/**
+ * Copyright © 2016-2017 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.server.system;
+
+import org.junit.ClassRule;
+import org.junit.extensions.cpsuite.ClasspathSuite;
+import org.junit.runner.RunWith;
+import org.thingsboard.server.dao.CustomPostgresUnit;
+
+import java.util.Arrays;
+
+/**
+ * Created by Valerii Sosliuk on 6/27/2017.
+ */
+@RunWith(ClasspathSuite.class)
+@ClasspathSuite.ClassnameFilters({"org.thingsboard.server.system.sql.*SqlTest"})
+public class SystemSqlTestSuite {
+
+ @ClassRule
+ public static CustomPostgresUnit postgresUnit = new CustomPostgresUnit(
+ Arrays.asList("postgres/schema.sql", "postgres/system-data.sql"),
+ "postgres-embedded-test.properties");
+
+
+}
diff --git a/application/src/test/java/org/thingsboard/server/system/nosql/DeviceApiNoSqlTest.java b/application/src/test/java/org/thingsboard/server/system/nosql/DeviceApiNoSqlTest.java
new file mode 100644
index 0000000000..7ee185cd0f
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/system/nosql/DeviceApiNoSqlTest.java
@@ -0,0 +1,27 @@
+/**
+ * Copyright © 2016-2017 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.server.system.nosql;
+
+import org.thingsboard.server.dao.service.DaoNoSqlTest;
+import org.thingsboard.server.dao.util.NoSqlDao;
+import org.thingsboard.server.system.BaseHttpDeviceApiTest;
+
+/**
+ * Created by Valerii Sosliuk on 6/27/2017.
+ */
+@DaoNoSqlTest
+public class DeviceApiNoSqlTest extends BaseHttpDeviceApiTest {
+}
diff --git a/application/src/test/java/org/thingsboard/server/system/sql/DeviceApiSqlTest.java b/application/src/test/java/org/thingsboard/server/system/sql/DeviceApiSqlTest.java
new file mode 100644
index 0000000000..bf533132c7
--- /dev/null
+++ b/application/src/test/java/org/thingsboard/server/system/sql/DeviceApiSqlTest.java
@@ -0,0 +1,26 @@
+/**
+ * Copyright © 2016-2017 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.server.system.sql;
+
+import org.thingsboard.server.dao.service.DaoSqlTest;
+import org.thingsboard.server.system.BaseHttpDeviceApiTest;
+
+/**
+ * Created by Valerii Sosliuk on 6/27/2017.
+ */
+@DaoSqlTest
+public class DeviceApiSqlTest extends BaseHttpDeviceApiTest{
+}
diff --git a/application/src/test/resources/thingsboard-test.properties b/application/src/test/resources/thingsboard-test.properties
deleted file mode 100644
index 203f23972a..0000000000
--- a/application/src/test/resources/thingsboard-test.properties
+++ /dev/null
@@ -1 +0,0 @@
-updates.enabled=false
\ No newline at end of file
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/asset/TenantAssetType.java b/common/data/src/main/java/org/thingsboard/server/common/data/asset/TenantAssetType.java
index bd23f2aa5a..8e0eb0e69e 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/asset/TenantAssetType.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/asset/TenantAssetType.java
@@ -17,6 +17,8 @@ package org.thingsboard.server.common.data.asset;
import org.thingsboard.server.common.data.id.TenantId;
+import java.util.UUID;
+
public class TenantAssetType {
private static final long serialVersionUID = 8057290243855622101L;
@@ -33,6 +35,11 @@ public class TenantAssetType {
this.tenantId = tenantId;
}
+ public TenantAssetType(String type, UUID tenantId) {
+ this.type = type;
+ this.tenantId = new TenantId(tenantId);
+ }
+
public String getType() {
return type;
}
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java
index 8d60f525f4..8ede698241 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java
@@ -15,8 +15,6 @@
*/
package org.thingsboard.server.common.data.kv;
-import java.util.Optional;
-
public interface TsKvQuery {
String getKey();
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelation.java b/common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelation.java
index 8dab589596..478d5c803d 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelation.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelation.java
@@ -18,8 +18,6 @@ package org.thingsboard.server.common.data.relation;
import com.fasterxml.jackson.databind.JsonNode;
import org.thingsboard.server.common.data.id.EntityId;
-import java.util.Objects;
-
public class EntityRelation {
private static final long serialVersionUID = 2807343040519543363L;
diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleMetaData.java b/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleMetaData.java
index 8a0f847617..1451d2c1a6 100644
--- a/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleMetaData.java
+++ b/common/data/src/main/java/org/thingsboard/server/common/data/rule/RuleMetaData.java
@@ -15,15 +15,12 @@
*/
package org.thingsboard.server.common.data.rule;
+import com.fasterxml.jackson.databind.JsonNode;
import lombok.Data;
-import lombok.ToString;
import org.thingsboard.server.common.data.HasName;
import org.thingsboard.server.common.data.SearchTextBased;
-import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.RuleId;
import org.thingsboard.server.common.data.id.TenantId;
-
-import com.fasterxml.jackson.databind.JsonNode;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleState;
@Data
diff --git a/dao/pom.xml b/dao/pom.xml
index cc3804feda..2bb77511bb 100644
--- a/dao/pom.xml
+++ b/dao/pom.xml
@@ -57,14 +57,24 @@
logback-classic
- org.springframework
- spring-test
+ org.postgresql
+ postgresql
junit
junit
test
+
+ org.dbunit
+ dbunit
+ test
+
+
+ com.github.springtestdbunit
+ spring-test-dbunit
+ test
+
org.mockito
mockito-all
@@ -150,6 +160,20 @@
org.bouncycastle
bcprov-jdk15on
+
+ org.springframework.boot
+ spring-boot-starter-data-jpa
+
+
+ org.springframework
+ spring-test
+ test
+
+
+ ru.yandex.qatools.embed
+ postgresql-embedded
+ test
+
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 2703cdc23a..f0580eba3e 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/Dao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/Dao.java
@@ -15,7 +15,6 @@
*/
package org.thingsboard.server.dao;
-import com.datastax.driver.core.ResultSet;
import com.google.common.util.concurrent.ListenableFuture;
import java.util.List;
@@ -31,6 +30,6 @@ public interface Dao {
T save(T t);
- ResultSet removeById(UUID id);
+ boolean removeById(UUID id);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java b/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java
index 27499bb1bd..3392ef7f3c 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/DaoUtil.java
@@ -15,15 +15,11 @@
*/
package org.thingsboard.server.dao;
-import java.util.ArrayList;
-import java.util.Collection;
-import java.util.Collections;
-import java.util.List;
-import java.util.UUID;
-
import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.dao.model.ToData;
+import java.util.*;
+
public abstract class DaoUtil {
private DaoUtil() {
@@ -34,7 +30,9 @@ public abstract class DaoUtil {
if (toDataList != null && !toDataList.isEmpty()) {
list = new ArrayList<>();
for (ToData object : toDataList) {
- list.add(object.toData());
+ if (object != null) {
+ list.add(object.toData());
+ }
}
}
return list;
diff --git a/dao/src/main/java/org/thingsboard/server/dao/EncryptionUtil.java b/dao/src/main/java/org/thingsboard/server/dao/EncryptionUtil.java
index 9a4e592e07..1d74d30a08 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/EncryptionUtil.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/EncryptionUtil.java
@@ -15,7 +15,6 @@
*/
package org.thingsboard.server.dao;
-import com.google.common.base.CharMatcher;
import lombok.extern.slf4j.Slf4j;
import org.bouncycastle.crypto.digests.SHA3Digest;
import org.bouncycastle.pqc.math.linearalgebra.ByteUtils;
diff --git a/dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java b/dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java
new file mode 100644
index 0000000000..5e4d5c17e3
--- /dev/null
+++ b/dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java
@@ -0,0 +1,38 @@
+/**
+ * Copyright © 2016-2017 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.server.dao;
+
+import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
+import org.springframework.boot.autoconfigure.domain.EntityScan;
+import org.springframework.context.annotation.ComponentScan;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.data.jpa.repository.config.EnableJpaRepositories;
+import org.springframework.transaction.annotation.EnableTransactionManagement;
+import org.thingsboard.server.dao.util.SqlDao;
+
+/**
+ * @author Valerii Sosliuk
+ */
+@Configuration
+@EnableAutoConfiguration
+@ComponentScan("org.thingsboard.server.dao.sql")
+@EnableJpaRepositories("org.thingsboard.server.dao.sql")
+@EntityScan("org.thingsboard.server.dao.model.sql")
+@EnableTransactionManagement
+@SqlDao
+public class JpaDaoConfig {
+
+}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/NoSqlDaoConfig.java b/dao/src/main/java/org/thingsboard/server/dao/NoSqlDaoConfig.java
new file mode 100644
index 0000000000..904d390657
--- /dev/null
+++ b/dao/src/main/java/org/thingsboard/server/dao/NoSqlDaoConfig.java
@@ -0,0 +1,33 @@
+/**
+ * Copyright © 2016-2017 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.server.dao;
+
+import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
+import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration;
+import org.springframework.boot.autoconfigure.jdbc.DataSourceTransactionManagerAutoConfiguration;
+import org.springframework.boot.autoconfigure.orm.jpa.HibernateJpaAutoConfiguration;
+import org.springframework.context.annotation.Configuration;
+import org.thingsboard.server.dao.util.NoSqlDao;
+
+@Configuration
+@EnableAutoConfiguration(
+ exclude = {
+ DataSourceAutoConfiguration.class,
+ DataSourceTransactionManagerAutoConfiguration.class,
+ HibernateJpaAutoConfiguration.class})
+@NoSqlDao
+public class NoSqlDaoConfig {
+}
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 289d5c9e59..0ab762024e 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
@@ -22,7 +22,6 @@ import org.thingsboard.server.common.data.alarm.AlarmQuery;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.dao.Dao;
-import org.thingsboard.server.dao.model.AlarmEntity;
import java.util.List;
import java.util.UUID;
@@ -30,13 +29,13 @@ import java.util.UUID;
/**
* Created by ashvayka on 11.05.17.
*/
-public interface AlarmDao extends Dao {
+public interface AlarmDao extends Dao {
ListenableFuture findLatestByOriginatorAndType(TenantId tenantId, EntityId originator, String type);
ListenableFuture findAlarmByIdAsync(UUID key);
- AlarmEntity save(Alarm alarm);
+ Alarm save(Alarm alarm);
ListenableFuture> findAlarms(AlarmQuery query);
}
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 785d029850..81f670a096 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
@@ -24,6 +24,7 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
+import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.alarm.*;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.page.TimePageData;
@@ -31,10 +32,8 @@ import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.dao.entity.AbstractEntityService;
-import org.thingsboard.server.dao.entity.BaseEntityService;
import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.dao.exception.DataValidationException;
-import org.thingsboard.server.dao.model.*;
import org.thingsboard.server.dao.relation.EntityRelationsQuery;
import org.thingsboard.server.dao.relation.EntitySearchDirection;
import org.thingsboard.server.dao.relation.RelationService;
@@ -53,8 +52,7 @@ import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.stream.Collectors;
-import static org.thingsboard.server.dao.DaoUtil.*;
-import static org.thingsboard.server.dao.service.Validator.*;
+import static org.thingsboard.server.dao.service.Validator.validateId;
@Service
@Slf4j
@@ -115,7 +113,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
private Alarm createAlarm(Alarm alarm) throws InterruptedException, ExecutionException {
log.debug("New Alarm : {}", alarm);
- Alarm saved = getData(alarmDao.save(new AlarmEntity(alarm)));
+ Alarm saved = alarmDao.save(alarm);
EntityRelationsQuery query = new EntityRelationsQuery();
query.setParameters(new RelationsSearchParameters(saved.getOriginator(), EntitySearchDirection.TO, Integer.MAX_VALUE));
List parentEntities = relationService.findByQuery(query).get().stream().map(r -> r.getFrom()).collect(Collectors.toList());
@@ -144,11 +142,11 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
private Alarm updateAlarm(Alarm oldAlarm, Alarm newAlarm) {
AlarmStatus oldStatus = oldAlarm.getStatus();
AlarmStatus newStatus = newAlarm.getStatus();
- AlarmEntity result = alarmDao.save(new AlarmEntity(merge(oldAlarm, newAlarm)));
+ Alarm result = alarmDao.save(merge(oldAlarm, newAlarm));
if (oldStatus != newStatus) {
updateRelations(oldAlarm, oldStatus, newStatus);
}
- return result.toData();
+ return result;
}
@Override
@@ -164,7 +162,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(new AlarmEntity(alarm));
+ alarmDao.save(alarm);
updateRelations(alarm, oldStatus, newStatus);
return true;
}
@@ -185,7 +183,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
AlarmStatus newStatus = oldStatus.isAck() ? AlarmStatus.CLEARED_ACK : AlarmStatus.CLEARED_UNACK;
alarm.setStatus(newStatus);
alarm.setClearTs(clearTime);
- alarmDao.save(new AlarmEntity(alarm));
+ alarmDao.save(alarm);
updateRelations(alarm, oldStatus, newStatus);
return true;
}
@@ -387,7 +385,7 @@ public class BaseAlarmService extends AbstractEntityService implements AlarmServ
if (alarm.getTenantId() == null) {
throw new DataValidationException("Alarm should be assigned to tenant!");
} else {
- TenantEntity tenant = tenantDao.findById(alarm.getTenantId().getId());
+ Tenant tenant = tenantDao.findById(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/AlarmDaoImpl.java b/dao/src/main/java/org/thingsboard/server/dao/alarm/CassandraAlarmDao.java
similarity index 89%
rename from dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDaoImpl.java
rename to dao/src/main/java/org/thingsboard/server/dao/alarm/CassandraAlarmDao.java
index 72fbae5e6f..83578eb29a 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/alarm/AlarmDaoImpl.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/alarm/CassandraAlarmDao.java
@@ -33,10 +33,10 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
-import org.thingsboard.server.dao.AbstractModelDao;
-import org.thingsboard.server.dao.AbstractSearchTimeDao;
-import org.thingsboard.server.dao.model.AlarmEntity;
+import org.thingsboard.server.dao.nosql.CassandraAbstractModelDao;
+import org.thingsboard.server.dao.util.NoSqlDao;
import org.thingsboard.server.dao.model.ModelConstants;
+import org.thingsboard.server.dao.model.nosql.AlarmEntity;
import org.thingsboard.server.dao.relation.RelationDao;
import java.util.ArrayList;
@@ -49,7 +49,8 @@ import static org.thingsboard.server.dao.model.ModelConstants.*;
@Component
@Slf4j
-public class AlarmDaoImpl extends AbstractModelDao implements AlarmDao {
+@NoSqlDao
+public class CassandraAlarmDao extends CassandraAbstractModelDao implements AlarmDao {
@Autowired
private RelationDao relationDao;
@@ -69,9 +70,9 @@ public class AlarmDaoImpl extends AbstractModelDao implements Alarm
}
@Override
- public AlarmEntity save(Alarm alarm) {
+ public Alarm save(Alarm alarm) {
log.debug("Save asset [{}] ", alarm);
- return save(new AlarmEntity(alarm));
+ return super.save(alarm);
}
@Override
@@ -84,16 +85,7 @@ public class AlarmDaoImpl extends AbstractModelDao implements Alarm
query.and(eq(ALARM_TYPE_PROPERTY, type));
query.limit(1);
query.orderBy(QueryBuilder.asc(ModelConstants.ALARM_TYPE_PROPERTY), QueryBuilder.desc(ModelConstants.ID_PROPERTY));
- return Futures.transform(findOneByStatementAsync(query), toDataFunction());
- }
-
- @Override
- public ListenableFuture findAlarmByIdAsync(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 Futures.transform(findOneByStatementAsync(query), toDataFunction());
+ return findOneByStatementAsync(query);
}
@Override
@@ -114,10 +106,19 @@ public class AlarmDaoImpl extends AbstractModelDao implements Alarm
List> alarmFutures = new ArrayList<>(input.size());
for (EntityRelation relation : input) {
alarmFutures.add(Futures.transform(
- findAlarmByIdAsync(relation.getTo().getId()), (Function)
- alarm1 -> new AlarmInfo(alarm1)));
+ findAlarmByIdAsync(relation.getTo().getId()),
+ (Function) AlarmInfo::new));
}
return Futures.successfulAsList(alarmFutures);
});
}
+
+ @Override
+ public ListenableFuture findAlarmByIdAsync(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);
+ }
}
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 c2ffd1fcc7..2b720ae04a 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
@@ -17,9 +17,9 @@ package org.thingsboard.server.dao.asset;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.asset.Asset;
+import org.thingsboard.server.common.data.asset.TenantAssetType;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.dao.Dao;
-import org.thingsboard.server.dao.model.AssetEntity;
import org.thingsboard.server.dao.model.TenantAssetTypeEntity;
import java.util.List;
@@ -30,7 +30,7 @@ import java.util.UUID;
* The Interface AssetDao.
*
*/
-public interface AssetDao extends Dao {
+public interface AssetDao extends Dao {
/**
* Save or update asset object
@@ -38,7 +38,7 @@ public interface AssetDao extends Dao {
* @param asset the asset object
* @return saved asset object
*/
- AssetEntity save(Asset asset);
+ Asset save(Asset asset);
/**
* Find assets by tenantId and page link.
@@ -47,7 +47,7 @@ public interface AssetDao extends Dao {
* @param pageLink the page link
* @return the list of asset objects
*/
- List findAssetsByTenantId(UUID tenantId, TextPageLink pageLink);
+ List findAssetsByTenantId(UUID tenantId, TextPageLink pageLink);
/**
* Find assets by tenantId, type and page link.
@@ -57,7 +57,7 @@ public interface AssetDao extends Dao {
* @param pageLink the page link
* @return the list of asset objects
*/
- List findAssetsByTenantIdAndType(UUID tenantId, String type, TextPageLink pageLink);
+ List findAssetsByTenantIdAndType(UUID tenantId, String type, TextPageLink pageLink);
/**
* Find assets by tenantId and assets Ids.
@@ -66,7 +66,7 @@ public interface AssetDao extends Dao {
* @param assetIds the asset Ids
* @return the list of asset objects
*/
- ListenableFuture> findAssetsByTenantIdAndIdsAsync(UUID tenantId, List assetIds);
+ ListenableFuture> findAssetsByTenantIdAndIdsAsync(UUID tenantId, List assetIds);
/**
* Find assets by tenantId, customerId and page link.
@@ -76,7 +76,7 @@ public interface AssetDao extends Dao {
* @param pageLink the page link
* @return the list of asset objects
*/
- List findAssetsByTenantIdAndCustomerId(UUID tenantId, UUID customerId, TextPageLink pageLink);
+ List findAssetsByTenantIdAndCustomerId(UUID tenantId, UUID customerId, TextPageLink pageLink);
/**
* Find assets by tenantId, customerId, type and page link.
@@ -87,7 +87,7 @@ public interface AssetDao extends Dao {
* @param pageLink the page link
* @return the list of asset objects
*/
- List findAssetsByTenantIdAndCustomerIdAndType(UUID tenantId, UUID customerId, String type, TextPageLink pageLink);
+ List findAssetsByTenantIdAndCustomerIdAndType(UUID tenantId, UUID customerId, String type, TextPageLink pageLink);
/**
* Find assets by tenantId, customerId and assets Ids.
@@ -97,7 +97,7 @@ public interface AssetDao extends Dao {
* @param assetIds the asset Ids
* @return the list of asset objects
*/
- ListenableFuture> findAssetsByTenantIdCustomerIdAndIdsAsync(UUID tenantId, UUID customerId, List assetIds);
+ ListenableFuture> findAssetsByTenantIdAndCustomerIdAndIdsAsync(UUID tenantId, UUID customerId, List assetIds);
/**
* Find assets by tenantId and asset name.
@@ -106,13 +106,13 @@ public interface AssetDao extends Dao {
* @param name the asset name
* @return the optional asset object
*/
- Optional findAssetsByTenantIdAndName(UUID tenantId, String name);
+ Optional findAssetsByTenantIdAndName(UUID tenantId, String name);
/**
* Find tenants asset types.
*
* @return the list of tenant asset type objects
*/
- ListenableFuture> findTenantAssetTypesAsync();
+ ListenableFuture> findTenantAssetTypesAsync();
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetSearchQuery.java b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetSearchQuery.java
index f3c69b23bc..440f0d2d22 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetSearchQuery.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/asset/AssetSearchQuery.java
@@ -23,7 +23,6 @@ import org.thingsboard.server.dao.relation.EntityRelationsQuery;
import org.thingsboard.server.dao.relation.EntityTypeFilter;
import javax.annotation.Nullable;
-import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
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 3a5f803046..7e9f129a10 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
@@ -24,7 +24,9 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
+import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.EntityType;
+import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.asset.TenantAssetType;
import org.thingsboard.server.common.data.id.AssetId;
@@ -37,7 +39,6 @@ import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.dao.customer.CustomerDao;
import org.thingsboard.server.dao.entity.AbstractEntityService;
import org.thingsboard.server.dao.exception.DataValidationException;
-import org.thingsboard.server.dao.model.*;
import org.thingsboard.server.dao.relation.EntitySearchDirection;
import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.dao.service.PaginatedRemover;
@@ -45,6 +46,7 @@ import org.thingsboard.server.dao.tenant.TenantDao;
import javax.annotation.Nullable;
import java.util.ArrayList;
+import java.util.Comparator;
import java.util.List;
import java.util.Optional;
import java.util.stream.Collectors;
@@ -70,35 +72,28 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
public Asset findAssetById(AssetId assetId) {
log.trace("Executing findAssetById [{}]", assetId);
validateId(assetId, "Incorrect assetId " + assetId);
- AssetEntity assetEntity = assetDao.findById(assetId.getId());
- return getData(assetEntity);
+ return assetDao.findById(assetId.getId());
}
@Override
public ListenableFuture findAssetByIdAsync(AssetId assetId) {
log.trace("Executing findAssetById [{}]", assetId);
validateId(assetId, "Incorrect assetId " + assetId);
- ListenableFuture assetEntity = assetDao.findByIdAsync(assetId.getId());
- return Futures.transform(assetEntity, (Function super AssetEntity, ? extends Asset>) input -> getData(input));
+ return assetDao.findByIdAsync(assetId.getId());
}
@Override
public Optional findAssetByTenantIdAndName(TenantId tenantId, String name) {
log.trace("Executing findAssetByTenantIdAndName [{}][{}]", tenantId, name);
validateId(tenantId, "Incorrect tenantId " + tenantId);
- Optional assetEntityOpt = assetDao.findAssetsByTenantIdAndName(tenantId.getId(), name);
- if (assetEntityOpt.isPresent()) {
- return Optional.of(getData(assetEntityOpt.get()));
- } else {
- return Optional.empty();
- }
+ return assetDao.findAssetsByTenantIdAndName(tenantId.getId(), name);
}
@Override
public Asset saveAsset(Asset asset) {
log.trace("Executing saveAsset [{}]", asset);
assetValidator.validate(asset);
- return getData(assetDao.save(asset));
+ return assetDao.save(asset);
}
@Override
@@ -128,8 +123,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
log.trace("Executing findAssetsByTenantId, tenantId [{}], pageLink [{}]", tenantId, pageLink);
validateId(tenantId, "Incorrect tenantId " + tenantId);
validatePageLink(pageLink, "Incorrect page link " + pageLink);
- List assetEntities = assetDao.findAssetsByTenantId(tenantId.getId(), pageLink);
- List assets = convertDataList(assetEntities);
+ List assets = assetDao.findAssetsByTenantId(tenantId.getId(), pageLink);
return new TextPageData<>(assets, pageLink);
}
@@ -139,8 +133,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
validateId(tenantId, "Incorrect tenantId " + tenantId);
validateString(type, "Incorrect type " + type);
validatePageLink(pageLink, "Incorrect page link " + pageLink);
- List assetEntities = assetDao.findAssetsByTenantIdAndType(tenantId.getId(), type, pageLink);
- List assets = convertDataList(assetEntities);
+ List assets = assetDao.findAssetsByTenantIdAndType(tenantId.getId(), type, pageLink);
return new TextPageData<>(assets, pageLink);
}
@@ -149,15 +142,14 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
log.trace("Executing findAssetsByTenantIdAndIdsAsync, tenantId [{}], assetIds [{}]", tenantId, assetIds);
validateId(tenantId, "Incorrect tenantId " + tenantId);
validateIds(assetIds, "Incorrect assetIds " + assetIds);
- ListenableFuture> assetEntities = assetDao.findAssetsByTenantIdAndIdsAsync(tenantId.getId(), toUUIDs(assetIds));
- return Futures.transform(assetEntities, (Function, List>) input -> convertDataList(input));
+ return assetDao.findAssetsByTenantIdAndIdsAsync(tenantId.getId(), toUUIDs(assetIds));
}
@Override
public void deleteAssetsByTenantId(TenantId tenantId) {
log.trace("Executing deleteAssetsByTenantId, tenantId [{}]", tenantId);
validateId(tenantId, "Incorrect tenantId " + tenantId);
- tenantAssetsRemover.removeEntitites(tenantId);
+ tenantAssetsRemover.removeEntities(tenantId);
}
@Override
@@ -166,9 +158,8 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
validateId(tenantId, "Incorrect tenantId " + tenantId);
validateId(customerId, "Incorrect customerId " + customerId);
validatePageLink(pageLink, "Incorrect page link " + pageLink);
- List assetEntities = assetDao.findAssetsByTenantIdAndCustomerId(tenantId.getId(), customerId.getId(), pageLink);
- List assets = convertDataList(assetEntities);
- return new TextPageData<>(assets, pageLink);
+ List assets = assetDao.findAssetsByTenantIdAndCustomerId(tenantId.getId(), customerId.getId(), pageLink);
+ return new TextPageData(assets, pageLink);
}
@Override
@@ -178,20 +169,17 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
validateId(customerId, "Incorrect customerId " + customerId);
validateString(type, "Incorrect type " + type);
validatePageLink(pageLink, "Incorrect page link " + pageLink);
- List assetEntities = assetDao.findAssetsByTenantIdAndCustomerIdAndType(tenantId.getId(), customerId.getId(), type, pageLink);
- List assets = convertDataList(assetEntities);
+ List assets = assetDao.findAssetsByTenantIdAndCustomerIdAndType(tenantId.getId(), customerId.getId(), type, pageLink);
return new TextPageData<>(assets, pageLink);
}
@Override
public ListenableFuture> findAssetsByTenantIdCustomerIdAndIdsAsync(TenantId tenantId, CustomerId customerId, List assetIds) {
- log.trace("Executing findAssetsByTenantIdCustomerIdAndIdsAsync, tenantId [{}], customerId [{}], assetIds [{}]", tenantId, customerId, assetIds);
+ log.trace("Executing findAssetsByTenantIdAndCustomerIdAndIdsAsync, tenantId [{}], customerId [{}], assetIds [{}]", tenantId, customerId, assetIds);
validateId(tenantId, "Incorrect tenantId " + tenantId);
validateId(customerId, "Incorrect customerId " + customerId);
validateIds(assetIds, "Incorrect assetIds " + assetIds);
- ListenableFuture> assetEntities = assetDao.findAssetsByTenantIdCustomerIdAndIdsAsync(tenantId.getId(),
- customerId.getId(), toUUIDs(assetIds));
- return Futures.transform(assetEntities, (Function, List>) input -> convertDataList(input));
+ return assetDao.findAssetsByTenantIdAndCustomerIdAndIdsAsync(tenantId.getId(), customerId.getId(), toUUIDs(assetIds));
}
@Override
@@ -199,7 +187,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
log.trace("Executing unassignCustomerAssets, tenantId [{}], customerId [{}]", tenantId, customerId);
validateId(tenantId, "Incorrect tenantId " + tenantId);
validateId(customerId, "Incorrect customerId " + customerId);
- new CustomerAssetsUnassigner(tenantId).removeEntitites(customerId);
+ new CustomerAssetsUnassigner(tenantId).removeEntities(customerId);
}
@Override
@@ -232,16 +220,16 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
public ListenableFuture> findAssetTypesByTenantId(TenantId tenantId) {
log.trace("Executing findAssetTypesByTenantId, tenantId [{}]", tenantId);
validateId(tenantId, "Incorrect tenantId " + tenantId);
- ListenableFuture> tenantAssetTypeEntities = assetDao.findTenantAssetTypesAsync();
+ ListenableFuture> tenantAssetTypeEntities = assetDao.findTenantAssetTypesAsync();
ListenableFuture> tenantAssetTypes = Futures.transform(tenantAssetTypeEntities,
- (Function, List>) assetTypeEntities -> {
+ (Function, List>) assetTypeEntities -> {
List assetTypes = new ArrayList<>();
- for (TenantAssetTypeEntity assetTypeEntity : assetTypeEntities) {
- if (assetTypeEntity.getTenantId().equals(tenantId.getId())) {
- assetTypes.add(assetTypeEntity.toTenantAssetType());
+ for (TenantAssetType assetType : assetTypeEntities) {
+ if (assetType.getTenantId().equals(tenantId)) {
+ assetTypes.add(assetType);
}
}
- assetTypes.sort((TenantAssetType o1, TenantAssetType o2) -> o1.getType().compareTo(o2.getType()));
+ assetTypes.sort(Comparator.comparing(TenantAssetType::getType));
return assetTypes;
});
return tenantAssetTypes;
@@ -263,7 +251,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
protected void validateUpdate(Asset asset) {
assetDao.findAssetsByTenantIdAndName(asset.getTenantId().getId(), asset.getName()).ifPresent(
d -> {
- if (!d.getId().equals(asset.getUuidId())) {
+ if (!d.getId().equals(asset.getId())) {
throw new DataValidationException("Asset with such name already exists!");
}
}
@@ -281,7 +269,7 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
if (asset.getTenantId() == null) {
throw new DataValidationException("Asset should be assigned to tenant!");
} else {
- TenantEntity tenant = tenantDao.findById(asset.getTenantId().getId());
+ Tenant tenant = tenantDao.findById(asset.getTenantId().getId());
if (tenant == null) {
throw new DataValidationException("Asset is referencing to non-existent tenant!");
}
@@ -289,32 +277,32 @@ 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)) {
- CustomerEntity customer = customerDao.findById(asset.getCustomerId().getId());
+ Customer customer = customerDao.findById(asset.getCustomerId().getId());
if (customer == null) {
throw new DataValidationException("Can't assign asset to non-existent customer!");
}
- if (!customer.getTenantId().equals(asset.getTenantId().getId())) {
+ if (!customer.getTenantId().equals(asset.getTenantId())) {
throw new DataValidationException("Can't assign asset to customer from different tenant!");
}
}
}
};
- private PaginatedRemover tenantAssetsRemover =
- new PaginatedRemover() {
+ private PaginatedRemover tenantAssetsRemover =
+ new PaginatedRemover() {
@Override
- protected List findEntities(TenantId id, TextPageLink pageLink) {
+ protected List findEntities(TenantId id, TextPageLink pageLink) {
return assetDao.findAssetsByTenantId(id.getId(), pageLink);
}
@Override
- protected void removeEntity(AssetEntity entity) {
- deleteAsset(new AssetId(entity.getId()));
+ protected void removeEntity(Asset entity) {
+ deleteAsset(new AssetId(entity.getId().getId()));
}
};
- class CustomerAssetsUnassigner extends PaginatedRemover {
+ class CustomerAssetsUnassigner extends PaginatedRemover {
private TenantId tenantId;
@@ -323,13 +311,13 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
}
@Override
- protected List findEntities(CustomerId id, TextPageLink pageLink) {
+ protected List findEntities(CustomerId id, TextPageLink pageLink) {
return assetDao.findAssetsByTenantIdAndCustomerId(tenantId.getId(), id.getId(), pageLink);
}
@Override
- protected void removeEntity(AssetEntity entity) {
- unassignAssetFromCustomer(new AssetId(entity.getId()));
+ protected void removeEntity(Asset entity) {
+ unassignAssetFromCustomer(new AssetId(entity.getId().getId()));
}
}
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetDaoImpl.java b/dao/src/main/java/org/thingsboard/server/dao/asset/CassandraAssetDao.java
similarity index 72%
rename from dao/src/main/java/org/thingsboard/server/dao/asset/AssetDaoImpl.java
rename to dao/src/main/java/org/thingsboard/server/dao/asset/CassandraAssetDao.java
index 5dba18b07e..c2025b9f7d 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/asset/AssetDaoImpl.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/asset/CassandraAssetDao.java
@@ -25,10 +25,13 @@ import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.asset.Asset;
+import org.thingsboard.server.common.data.asset.TenantAssetType;
import org.thingsboard.server.common.data.page.TextPageLink;
-import org.thingsboard.server.dao.AbstractSearchTextDao;
-import org.thingsboard.server.dao.model.AssetEntity;
+import org.thingsboard.server.dao.nosql.CassandraAbstractSearchTextDao;
+import org.thingsboard.server.dao.DaoUtil;
+import org.thingsboard.server.dao.util.NoSqlDao;
import org.thingsboard.server.dao.model.TenantAssetTypeEntity;
+import org.thingsboard.server.dao.model.nosql.AssetEntity;
import javax.annotation.Nullable;
import java.util.*;
@@ -38,7 +41,8 @@ import static org.thingsboard.server.dao.model.ModelConstants.*;
@Component
@Slf4j
-public class AssetDaoImpl extends AbstractSearchTextDao implements AssetDao {
+@NoSqlDao
+public class CassandraAssetDao extends CassandraAbstractSearchTextDao implements AssetDao {
@Override
protected Class getColumnFamilyClass() {
@@ -51,33 +55,26 @@ public class AssetDaoImpl extends AbstractSearchTextDao implements
}
@Override
- public AssetEntity save(Asset asset) {
- log.debug("Save asset [{}] ", asset);
- return save(new AssetEntity(asset));
- }
-
- @Override
- public List findAssetsByTenantId(UUID tenantId, TextPageLink pageLink) {
+ public List 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,
Collections.singletonList(eq(ASSET_TENANT_ID_PROPERTY, tenantId)), pageLink);
log.trace("Found assets [{}] by tenantId [{}] and pageLink [{}]", assetEntities, tenantId, pageLink);
- return assetEntities;
+ return DaoUtil.convertDataList(assetEntities);
}
@Override
- public List findAssetsByTenantIdAndType(UUID tenantId, String type, TextPageLink pageLink) {
+ public List 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,
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);
- return assetEntities;
+ return DaoUtil.convertDataList(assetEntities);
}
- @Override
- public ListenableFuture> findAssetsByTenantIdAndIdsAsync(UUID tenantId, List assetIds) {
+ public ListenableFuture> findAssetsByTenantIdAndIdsAsync(UUID tenantId, List assetIds) {
log.debug("Try to find assets by tenantId [{}] and asset Ids [{}]", tenantId, assetIds);
Select select = select().from(getColumnFamilyName());
Select.Where query = select.where();
@@ -87,7 +84,7 @@ public class AssetDaoImpl extends AbstractSearchTextDao implements
}
@Override
- public List findAssetsByTenantIdAndCustomerId(UUID tenantId, UUID customerId, TextPageLink pageLink) {
+ public List 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,
Arrays.asList(eq(ASSET_CUSTOMER_ID_PROPERTY, customerId),
@@ -95,11 +92,11 @@ public class AssetDaoImpl extends AbstractSearchTextDao implements
pageLink);
log.trace("Found assets [{}] by tenantId [{}], customerId [{}] and pageLink [{}]", assetEntities, tenantId, customerId, pageLink);
- return assetEntities;
+ return DaoUtil.convertDataList(assetEntities);
}
@Override
- public List findAssetsByTenantIdAndCustomerIdAndType(UUID tenantId, UUID customerId, String type, TextPageLink pageLink) {
+ public List 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,
Arrays.asList(eq(ASSET_TYPE_PROPERTY, type),
@@ -108,11 +105,11 @@ public class AssetDaoImpl extends AbstractSearchTextDao implements
pageLink);
log.trace("Found assets [{}] by tenantId [{}], customerId [{}], type [{}] and pageLink [{}]", assetEntities, tenantId, customerId, type, pageLink);
- return assetEntities;
+ return DaoUtil.convertDataList(assetEntities);
}
@Override
- public ListenableFuture> findAssetsByTenantIdCustomerIdAndIdsAsync(UUID tenantId, UUID customerId, List assetIds) {
+ public ListenableFuture> findAssetsByTenantIdAndCustomerIdAndIdsAsync(UUID tenantId, UUID customerId, List assetIds) {
log.debug("Try to find assets by tenantId [{}], customerId [{}] and asset Ids [{}]", tenantId, customerId, assetIds);
Select select = select().from(getColumnFamilyName());
Select.Where query = select.where();
@@ -123,16 +120,17 @@ public class AssetDaoImpl extends AbstractSearchTextDao implements
}
@Override
- public Optional findAssetsByTenantIdAndName(UUID tenantId, String assetName) {
+ public Optional findAssetsByTenantIdAndName(UUID tenantId, String assetName) {
Select select = select().from(ASSET_BY_TENANT_AND_NAME_VIEW_NAME);
Select.Where query = select.where();
query.and(eq(ASSET_TENANT_ID_PROPERTY, tenantId));
query.and(eq(ASSET_NAME_PROPERTY, assetName));
- return Optional.ofNullable(findOneByStatement(query));
+ AssetEntity assetEntity = (AssetEntity) findOneByStatement(query);
+ return Optional.ofNullable(DaoUtil.getData(assetEntity));
}
@Override
- public ListenableFuture> findTenantAssetTypesAsync() {
+ public ListenableFuture> findTenantAssetTypesAsync() {
Select statement = select().distinct().column(ASSET_TYPE_PROPERTY).column(ASSET_TENANT_ID_PROPERTY).from(ASSET_TYPES_BY_TENANT_VIEW_NAME);
statement.setConsistencyLevel(cluster.getDefaultReadConsistencyLevel());
ResultSetFuture resultSetFuture = getSession().executeAsync(statement);
@@ -148,7 +146,20 @@ public class AssetDaoImpl extends AbstractSearchTextDao implements
}
}
});
- return result;
+ return Futures.transform(result, new Function, List>() {
+ @Nullable
+ @Override
+ public List apply(@Nullable List entityList) {
+ List list = Collections.emptyList();
+ if (entityList != null && !entityList.isEmpty()) {
+ list = new ArrayList<>();
+ for (TenantAssetTypeEntity object : entityList) {
+ list.add(object.toTenantAssetType());
+ }
+ }
+ return list;
+ }
+ });
}
}
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 ae58d4d9c6..704210e4cc 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
@@ -15,8 +15,6 @@
*/
package org.thingsboard.server.dao.attributes;
-import com.datastax.driver.core.ResultSet;
-import com.datastax.driver.core.ResultSetFuture;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
@@ -36,7 +34,7 @@ public interface AttributesDao {
ListenableFuture> findAll(EntityId entityId, String attributeType);
- ResultSetFuture save(EntityId entityId, String attributeType, AttributeKvEntry attribute);
+ ListenableFuture save(EntityId entityId, String attributeType, AttributeKvEntry attribute);
- ListenableFuture> removeAll(EntityId entityId, String scope, List keys);
+ ListenableFuture> removeAll(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 6bf9fb2bd5..6090fa9349 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
@@ -15,12 +15,8 @@
*/
package org.thingsboard.server.dao.attributes;
-import com.datastax.driver.core.ResultSet;
-import com.datastax.driver.core.ResultSetFuture;
import com.google.common.util.concurrent.ListenableFuture;
-import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
-import org.thingsboard.server.common.data.id.UUIDBased;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import java.util.Collection;
@@ -38,7 +34,7 @@ public interface AttributesService {
ListenableFuture> findAll(EntityId entityId, String scope);
- ListenableFuture> save(EntityId entityId, String scope, List attributes);
+ ListenableFuture> save(EntityId entityId, String scope, List attributes);
- ListenableFuture> removeAll(EntityId entityId, String scope, List attributeKeys);
+ ListenableFuture> removeAll(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 43612419d0..bef2faa37f 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
@@ -20,11 +20,11 @@ import com.datastax.driver.core.ResultSetFuture;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.Futures;
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.kv.AttributeKvEntry;
import org.thingsboard.server.dao.exception.IncorrectParameterException;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.stereotype.Service;
import org.thingsboard.server.dao.service.Validator;
import java.util.Collection;
@@ -61,10 +61,10 @@ public class BaseAttributesService implements AttributesService {
}
@Override
- public ListenableFuture> save(EntityId entityId, String scope, List attributes) {
+ public ListenableFuture> save(EntityId entityId, String scope, List attributes) {
validate(entityId, scope);
attributes.forEach(attribute -> validate(attribute));
- List futures = Lists.newArrayListWithExpectedSize(attributes.size());
+ List> futures = Lists.newArrayListWithExpectedSize(attributes.size());
for (AttributeKvEntry attribute : attributes) {
futures.add(attributesDao.save(entityId, scope, attribute));
}
@@ -72,7 +72,7 @@ public class BaseAttributesService implements AttributesService {
}
@Override
- public ListenableFuture> removeAll(EntityId entityId, String scope, List keys) {
+ public ListenableFuture> removeAll(EntityId entityId, String scope, List keys) {
validate(entityId, scope);
return attributesDao.removeAll(entityId, scope, keys);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesDao.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java
similarity index 83%
rename from dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesDao.java
rename to dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java
index fd50f4d2ef..247463fbe4 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/attributes/BaseAttributesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/CassandraBaseAttributesDao.java
@@ -24,10 +24,12 @@ 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.dao.AbstractAsyncDao;
+import org.thingsboard.server.common.data.kv.AttributeKvEntry;
+import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
+import org.thingsboard.server.dao.nosql.CassandraAbstractAsyncDao;
+import org.thingsboard.server.dao.util.NoSqlDao;
import org.thingsboard.server.dao.model.ModelConstants;
-import org.thingsboard.server.common.data.kv.*;
-import org.thingsboard.server.dao.timeseries.BaseTimeseriesDao;
+import org.thingsboard.server.dao.timeseries.CassandraBaseTimeseriesDao;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
@@ -37,15 +39,17 @@ import java.util.List;
import java.util.Optional;
import java.util.stream.Collectors;
+import static com.datastax.driver.core.querybuilder.QueryBuilder.eq;
+import static com.datastax.driver.core.querybuilder.QueryBuilder.select;
import static org.thingsboard.server.dao.model.ModelConstants.*;
-import static com.datastax.driver.core.querybuilder.QueryBuilder.*;
/**
* @author Andrew Shvayka
*/
@Component
@Slf4j
-public class BaseAttributesDao extends AbstractAsyncDao implements AttributesDao {
+@NoSqlDao
+public class CassandraBaseAttributesDao extends CassandraAbstractAsyncDao implements AttributesDao {
private PreparedStatement saveStmt;
@@ -97,7 +101,7 @@ public class BaseAttributesDao extends AbstractAsyncDao implements AttributesDao
}
@Override
- public ResultSetFuture save(EntityId entityId, String attributeType, AttributeKvEntry attribute) {
+ public ListenableFuture save(EntityId entityId, String attributeType, AttributeKvEntry attribute) {
BoundStatement stmt = getSaveStmt().bind();
stmt.setString(0, entityId.getEntityType().name());
stmt.setUUID(1, entityId.getId());
@@ -120,23 +124,27 @@ public class BaseAttributesDao extends AbstractAsyncDao implements AttributesDao
} else {
stmt.setToNull(8);
}
- return executeAsyncWrite(stmt);
+ log.trace("Generated save stmt [{}] for entityId {} and attributeType {} and attribute", stmt, entityId, attributeType, attribute);
+ return getFuture(executeAsyncWrite(stmt), rs -> null);
}
@Override
- public ListenableFuture> removeAll(EntityId entityId, String attributeType, List keys) {
- List futures = keys.stream().map(key -> delete(entityId, attributeType, key)).collect(Collectors.toList());
+ public ListenableFuture> removeAll(EntityId entityId, String attributeType, List keys) {
+ List> futures = keys
+ .stream()
+ .map(key -> delete(entityId, attributeType, key))
+ .collect(Collectors.toList());
return Futures.allAsList(futures);
}
- private ResultSetFuture delete(EntityId entityId, String attributeType, String key) {
+ private ListenableFuture delete(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 getSession().executeAsync(delete);
+ return getFuture(getSession().executeAsync(delete), rs -> null);
}
private PreparedStatement getSaveStmt() {
@@ -161,7 +169,7 @@ public class BaseAttributesDao extends AbstractAsyncDao implements AttributesDao
AttributeKvEntry attributeEntry = null;
if (row != null) {
long lastUpdateTs = row.get(LAST_UPDATE_TS_COLUMN, Long.class);
- attributeEntry = new BaseAttributeKvEntry(BaseTimeseriesDao.toKvEntry(row, key), lastUpdateTs);
+ attributeEntry = new BaseAttributeKvEntry(CassandraBaseTimeseriesDao.toKvEntry(row, key), lastUpdateTs);
}
return attributeEntry;
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/cache/ServiceCacheConfiguration.java b/dao/src/main/java/org/thingsboard/server/dao/cache/ServiceCacheConfiguration.java
index 7c435bf779..35391785d0 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/cache/ServiceCacheConfiguration.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/cache/ServiceCacheConfiguration.java
@@ -57,16 +57,25 @@ public class ServiceCacheConfiguration {
Config config = new Config();
if (zkEnabled) {
- config.getNetworkConfig().getJoin().getMulticastConfig().setEnabled(false);
-
- config.setProperty(GroupProperty.DISCOVERY_SPI_ENABLED.getName(), Boolean.TRUE.toString());
- DiscoveryStrategyConfig discoveryStrategyConfig = new DiscoveryStrategyConfig(new ZookeeperDiscoveryStrategyFactory());
- discoveryStrategyConfig.addProperty(ZookeeperDiscoveryProperties.ZOOKEEPER_URL.key(), zkUrl);
- discoveryStrategyConfig.addProperty(ZookeeperDiscoveryProperties.ZOOKEEPER_PATH.key(), zkDir);
- discoveryStrategyConfig.addProperty(ZookeeperDiscoveryProperties.GROUP.key(), HAZELCAST_CLUSTER_NAME);
- config.getNetworkConfig().getJoin().getDiscoveryConfig().addDiscoveryStrategyConfig(discoveryStrategyConfig);
+ addZkConfig(config);
}
+ config.addMapConfig(createDeviceCredentialsCacheConfig());
+
+ return Hazelcast.newHazelcastInstance(config);
+ }
+
+ private void addZkConfig(Config config) {
+ config.getNetworkConfig().getJoin().getMulticastConfig().setEnabled(false);
+ config.setProperty(GroupProperty.DISCOVERY_SPI_ENABLED.getName(), Boolean.TRUE.toString());
+ DiscoveryStrategyConfig discoveryStrategyConfig = new DiscoveryStrategyConfig(new ZookeeperDiscoveryStrategyFactory());
+ discoveryStrategyConfig.addProperty(ZookeeperDiscoveryProperties.ZOOKEEPER_URL.key(), zkUrl);
+ discoveryStrategyConfig.addProperty(ZookeeperDiscoveryProperties.ZOOKEEPER_PATH.key(), zkDir);
+ discoveryStrategyConfig.addProperty(ZookeeperDiscoveryProperties.GROUP.key(), HAZELCAST_CLUSTER_NAME);
+ config.getNetworkConfig().getJoin().getDiscoveryConfig().addDiscoveryStrategyConfig(discoveryStrategyConfig);
+ }
+
+ private MapConfig createDeviceCredentialsCacheConfig() {
MapConfig deviceCredentialsCacheConfig = new MapConfig(CacheConstants.DEVICE_CREDENTIALS_CACHE);
deviceCredentialsCacheConfig.setTimeToLiveSeconds(cacheDeviceCredentialsTTL);
deviceCredentialsCacheConfig.setEvictionPolicy(EvictionPolicy.LRU);
@@ -75,9 +84,7 @@ public class ServiceCacheConfiguration {
cacheDeviceCredentialsMaxSizeSize,
MaxSizeConfig.MaxSizePolicy.valueOf(cacheDeviceCredentialsMaxSizePolicy))
);
- config.addMapConfig(deviceCredentialsCacheConfig);
-
- return Hazelcast.newHazelcastInstance(config);
+ return deviceCredentialsCacheConfig;
}
@Bean
diff --git a/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraCluster.java b/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraCluster.java
index 62e376245a..e180d140ad 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraCluster.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraCluster.java
@@ -16,13 +16,8 @@
package org.thingsboard.server.dao.cassandra;
-import com.datastax.driver.core.Cluster;
-import com.datastax.driver.core.ConsistencyLevel;
-import com.datastax.driver.core.HostDistance;
-import com.datastax.driver.core.PoolingOptions;
+import com.datastax.driver.core.*;
import com.datastax.driver.core.ProtocolOptions.Compression;
-import com.datastax.driver.core.Session;
-import com.datastax.driver.core.exceptions.NoHostAvailableException;
import com.datastax.driver.mapping.Mapper;
import com.datastax.driver.mapping.MappingManager;
import lombok.Data;
@@ -31,20 +26,19 @@ import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
-import org.thingsboard.server.dao.exception.DatabaseException;
+import org.thingsboard.server.dao.util.NoSqlDao;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
-import java.io.Closeable;
import java.net.InetSocketAddress;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
-import java.util.StringTokenizer;
+@Data
@Component
@Slf4j
-@Data
+@NoSqlDao
public class CassandraCluster {
private static final String COMMA = ",";
diff --git a/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraQueryOptions.java b/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraQueryOptions.java
index d5460b0081..1d29910c58 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraQueryOptions.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraQueryOptions.java
@@ -21,15 +21,14 @@ import lombok.Data;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Configuration;
import org.springframework.stereotype.Component;
-import org.springframework.util.StringUtils;
+import org.thingsboard.server.dao.util.NoSqlDao;
import javax.annotation.PostConstruct;
-import static org.apache.commons.lang3.StringUtils.isNotBlank;
-
@Component
@Configuration
@Data
+@NoSqlDao
public class CassandraQueryOptions {
@Value("${cassandra.query.default_fetch_size}")
diff --git a/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraSocketOptions.java b/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraSocketOptions.java
index c6f51d108e..5c8f196f05 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraSocketOptions.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/cassandra/CassandraSocketOptions.java
@@ -15,18 +15,19 @@
*/
package org.thingsboard.server.dao.cassandra;
+import com.datastax.driver.core.SocketOptions;
import lombok.Data;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Configuration;
import org.springframework.stereotype.Component;
-
-import com.datastax.driver.core.SocketOptions;
+import org.thingsboard.server.dao.util.NoSqlDao;
import javax.annotation.PostConstruct;
@Component
@Configuration
@Data
+@NoSqlDao
public class CassandraSocketOptions {
@Value("${cassandra.socket.connect_timeout}")
diff --git a/dao/src/main/java/org/thingsboard/server/dao/component/BaseComponentDescriptorService.java b/dao/src/main/java/org/thingsboard/server/dao/component/BaseComponentDescriptorService.java
index 22fa19d013..3a89e55489 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/component/BaseComponentDescriptorService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/component/BaseComponentDescriptorService.java
@@ -32,16 +32,12 @@ import org.thingsboard.server.common.data.plugin.ComponentScope;
import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.dao.exception.IncorrectParameterException;
-import org.thingsboard.server.dao.model.ComponentDescriptorEntity;
import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.dao.service.Validator;
import java.util.List;
import java.util.Optional;
-import static org.thingsboard.server.dao.DaoUtil.convertDataList;
-import static org.thingsboard.server.dao.DaoUtil.getData;
-
/**
* @author Andrew Shvayka
*/
@@ -55,39 +51,37 @@ public class BaseComponentDescriptorService implements ComponentDescriptorServic
@Override
public ComponentDescriptor saveComponent(ComponentDescriptor component) {
componentValidator.validate(component);
- Optional result = componentDescriptorDao.save(component);
+ Optional result = componentDescriptorDao.saveIfNotExist(component);
if (result.isPresent()) {
- return getData(result.get());
+ return result.get();
} else {
- return getData(componentDescriptorDao.findByClazz(component.getClazz()));
+ return componentDescriptorDao.findByClazz(component.getClazz());
}
}
@Override
public ComponentDescriptor findById(ComponentDescriptorId componentId) {
Validator.validateId(componentId, "Incorrect component id for search request.");
- return getData(componentDescriptorDao.findById(componentId));
+ return componentDescriptorDao.findById(componentId);
}
@Override
public ComponentDescriptor findByClazz(String clazz) {
Validator.validateString(clazz, "Incorrect clazz for search request.");
- return getData(componentDescriptorDao.findByClazz(clazz));
+ return componentDescriptorDao.findByClazz(clazz);
}
@Override
public TextPageData findByTypeAndPageLink(ComponentType type, TextPageLink pageLink) {
Validator.validatePageLink(pageLink, "Incorrect PageLink object for search plugin components request.");
- List pluginEntities = componentDescriptorDao.findByTypeAndPageLink(type, pageLink);
- List components = convertDataList(pluginEntities);
+ List components = componentDescriptorDao.findByTypeAndPageLink(type, pageLink);
return new TextPageData<>(components, pageLink);
}
@Override
public TextPageData findByScopeAndTypeAndPageLink(ComponentScope scope, ComponentType type, TextPageLink pageLink) {
Validator.validatePageLink(pageLink, "Incorrect PageLink object for search plugin components request.");
- List pluginEntities = componentDescriptorDao.findByScopeAndTypeAndPageLink(scope, type, pageLink);
- List components = convertDataList(pluginEntities);
+ List components = componentDescriptorDao.findByScopeAndTypeAndPageLink(scope, type, pageLink);
return new TextPageData<>(components, pageLink);
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/component/BaseComponentDescriptorDao.java b/dao/src/main/java/org/thingsboard/server/dao/component/CassandraBaseComponentDescriptorDao.java
similarity index 78%
rename from dao/src/main/java/org/thingsboard/server/dao/component/BaseComponentDescriptorDao.java
rename to dao/src/main/java/org/thingsboard/server/dao/component/CassandraBaseComponentDescriptorDao.java
index 7c43a765ae..d39b7f8bab 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/component/BaseComponentDescriptorDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/component/CassandraBaseComponentDescriptorDao.java
@@ -27,9 +27,11 @@ import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.plugin.ComponentDescriptor;
import org.thingsboard.server.common.data.plugin.ComponentScope;
import org.thingsboard.server.common.data.plugin.ComponentType;
-import org.thingsboard.server.dao.AbstractSearchTextDao;
+import org.thingsboard.server.dao.nosql.CassandraAbstractSearchTextDao;
+import org.thingsboard.server.dao.DaoUtil;
+import org.thingsboard.server.dao.util.NoSqlDao;
import org.thingsboard.server.dao.model.ModelConstants;
-import org.thingsboard.server.dao.model.ComponentDescriptorEntity;
+import org.thingsboard.server.dao.model.nosql.ComponentDescriptorEntity;
import java.util.Arrays;
import java.util.List;
@@ -44,7 +46,8 @@ import static com.datastax.driver.core.querybuilder.QueryBuilder.select;
*/
@Component
@Slf4j
-public class BaseComponentDescriptorDao extends AbstractSearchTextDao implements ComponentDescriptorDao {
+@NoSqlDao
+public class CassandraBaseComponentDescriptorDao extends CassandraAbstractSearchTextDao implements ComponentDescriptorDao {
@Override
protected Class getColumnFamilyClass() {
@@ -57,10 +60,10 @@ public class BaseComponentDescriptorDao extends AbstractSearchTextDao save(ComponentDescriptor component) {
+ public Optional saveIfNotExist(ComponentDescriptor component) {
ComponentDescriptorEntity entity = new ComponentDescriptorEntity(component);
log.debug("Save component entity [{}]", entity);
- Optional result = saveIfNotExist(entity);
+ Optional result = saveIfNotExist(entity);
if (log.isTraceEnabled()) {
log.trace("Saved result: [{}] for component entity [{}]", result.isPresent(), result.orElse(null));
} else {
@@ -70,19 +73,19 @@ public class BaseComponentDescriptorDao extends AbstractSearchTextDao findByTypeAndPageLink(ComponentType type, TextPageLink pageLink) {
+ public List findByTypeAndPageLink(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,
- Arrays.asList(eq(ModelConstants.COMPONENT_DESCRIPTOR_TYPE_PROPERTY, type.name())), pageLink);
+ Arrays.asList(eq(ModelConstants.COMPONENT_DESCRIPTOR_TYPE_PROPERTY, type)), pageLink);
if (log.isTraceEnabled()) {
log.trace("Search result: [{}]", Arrays.toString(entities.toArray()));
} else {
log.debug("Search result: [{}]", entities.size());
}
- return entities;
+ return DaoUtil.convertDataList(entities);
}
@Override
- public List findByScopeAndTypeAndPageLink(ComponentScope scope, ComponentType type, TextPageLink pageLink) {
+ public List findByScopeAndTypeAndPageLink(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,
- Arrays.asList(eq(ModelConstants.COMPONENT_DESCRIPTOR_TYPE_PROPERTY, type.name()),
+ Arrays.asList(eq(ModelConstants.COMPONENT_DESCRIPTOR_TYPE_PROPERTY, type),
eq(ModelConstants.COMPONENT_DESCRIPTOR_SCOPE_PROPERTY, scope.name())), pageLink);
if (log.isTraceEnabled()) {
log.trace("Search result: [{}]", Arrays.toString(entities.toArray()));
} else {
log.debug("Search result: [{}]", entities.size());
}
- return entities;
+ return DaoUtil.convertDataList(entities);
}
- public ResultSet removeById(UUID key) {
+ public boolean removeById(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 getSession().execute(delete);
+ return getSession().execute(delete).wasApplied();
}
@Override
public void deleteById(ComponentDescriptorId id) {
log.debug("Delete plugin meta-data entity by id [{}]", id);
- ResultSet resultSet = removeById(id.getId());
- log.debug("Delete result: [{}]", resultSet.wasApplied());
+ boolean result = removeById(id.getId());
+ log.debug("Delete result: [{}]", result);
}
@Override
@@ -144,7 +147,7 @@ public class BaseComponentDescriptorDao extends AbstractSearchTextDao saveIfNotExist(ComponentDescriptorEntity entity) {
+ private Optional saveIfNotExist(ComponentDescriptorEntity entity) {
if (entity.getId() == null) {
entity.setId(UUIDs.timeBased());
}
@@ -161,7 +164,7 @@ public class BaseComponentDescriptorDao extends AbstractSearchTextDao {
+public interface ComponentDescriptorDao extends Dao {
- Optional save(ComponentDescriptor component);
+ Optional saveIfNotExist(ComponentDescriptor component);
- ComponentDescriptorEntity findById(ComponentDescriptorId componentId);
+ ComponentDescriptor findById(ComponentDescriptorId componentId);
- ComponentDescriptorEntity findByClazz(String clazz);
+ ComponentDescriptor findByClazz(String clazz);
- List findByTypeAndPageLink(ComponentType type, TextPageLink pageLink);
+ List findByTypeAndPageLink(ComponentType type, TextPageLink pageLink);
- List findByScopeAndTypeAndPageLink(ComponentScope scope, ComponentType type, TextPageLink pageLink);
+ List findByScopeAndTypeAndPageLink(ComponentScope scope, ComponentType type, TextPageLink pageLink);
void deleteById(ComponentDescriptorId componentId);
diff --git a/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerDaoImpl.java b/dao/src/main/java/org/thingsboard/server/dao/customer/CassandraCustomerDao.java
similarity index 68%
rename from dao/src/main/java/org/thingsboard/server/dao/customer/CustomerDaoImpl.java
rename to dao/src/main/java/org/thingsboard/server/dao/customer/CassandraCustomerDao.java
index 25c111642a..cdb49f8b26 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerDaoImpl.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/customer/CassandraCustomerDao.java
@@ -15,31 +15,29 @@
*/
package org.thingsboard.server.dao.customer;
-import static com.datastax.driver.core.querybuilder.QueryBuilder.eq;
-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_TITLE_PROPERTY;
-import static org.thingsboard.server.dao.model.ModelConstants.CUSTOMER_TENANT_ID_PROPERTY;
-
-
-import java.util.Arrays;
-import java.util.List;
-import java.util.Optional;
-import java.util.UUID;
-
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.page.TextPageLink;
-import org.thingsboard.server.dao.AbstractSearchTextDao;
-import org.thingsboard.server.dao.model.CustomerEntity;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
+import org.thingsboard.server.dao.nosql.CassandraAbstractSearchTextDao;
+import org.thingsboard.server.dao.DaoUtil;
+import org.thingsboard.server.dao.util.NoSqlDao;
import org.thingsboard.server.dao.model.ModelConstants;
+import org.thingsboard.server.dao.model.nosql.CustomerEntity;
+
+import java.util.Arrays;
+import java.util.List;
+import java.util.Optional;
+import java.util.UUID;
+
+import static com.datastax.driver.core.querybuilder.QueryBuilder.eq;
+import static com.datastax.driver.core.querybuilder.QueryBuilder.select;
+import static org.thingsboard.server.dao.model.ModelConstants.*;
@Component
@Slf4j
-public class CustomerDaoImpl extends AbstractSearchTextDao implements CustomerDao {
+@NoSqlDao
+public class CassandraCustomerDao extends CassandraAbstractSearchTextDao implements CustomerDao {
@Override
protected Class getColumnFamilyClass() {
@@ -50,30 +48,26 @@ public class CustomerDaoImpl extends AbstractSearchTextDao imple
protected String getColumnFamilyName() {
return ModelConstants.CUSTOMER_COLUMN_FAMILY_NAME;
}
-
- @Override
- public CustomerEntity save(Customer customer) {
- log.debug("Save customer [{}] ", customer);
- return save(new CustomerEntity(customer));
- }
@Override
- public List findCustomersByTenantId(UUID tenantId, TextPageLink pageLink) {
+ public List 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,
Arrays.asList(eq(ModelConstants.CUSTOMER_TENANT_ID_PROPERTY, tenantId)),
pageLink);
log.trace("Found customers [{}] by tenantId [{}] and pageLink [{}]", customerEntities, tenantId, pageLink);
- return customerEntities;
+ return DaoUtil.convertDataList(customerEntities);
}
@Override
- public Optional findCustomersByTenantIdAndTitle(UUID tenantId, String title) {
+ public Optional findCustomersByTenantIdAndTitle(UUID tenantId, String title) {
Select select = select().from(CUSTOMER_BY_TENANT_AND_TITLE_VIEW_NAME);
Select.Where query = select.where();
query.and(eq(CUSTOMER_TENANT_ID_PROPERTY, tenantId));
query.and(eq(CUSTOMER_TITLE_PROPERTY, title));
- return Optional.ofNullable(findOneByStatement(query));
+ CustomerEntity customerEntity = findOneByStatement(query);
+ Customer customer = DaoUtil.getData(customerEntity);
+ return Optional.ofNullable(customer);
}
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerDao.java b/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerDao.java
index 5639044127..3c8f66da0a 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/customer/CustomerDao.java
@@ -15,20 +15,18 @@
*/
package org.thingsboard.server.dao.customer;
-import java.util.List;
import java.util.Optional;
-import java.util.UUID;
-
import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.dao.Dao;
-import org.thingsboard.server.dao.model.CustomerEntity;
-import org.thingsboard.server.dao.model.DeviceEntity;
+
+import java.util.List;
+import java.util.UUID;
/**
* The Interface CustomerDao.
*/
-public interface CustomerDao extends Dao {
+public interface CustomerDao extends Dao {
/**
* Save or update customer object
@@ -36,7 +34,7 @@ public interface CustomerDao extends Dao