diff --git a/application/pom.xml b/application/pom.xml index 0d9982f281..b653465050 100644 --- a/application/pom.xml +++ b/application/pom.xml @@ -144,21 +144,6 @@ org.eclipse.paho org.eclipse.paho.mqttv5.client - - org.cassandraunit - cassandra-unit - - - org.slf4j - slf4j-log4j12 - - - org.hibernate - hibernate-validator - - - test - org.thingsboard ui-ngx @@ -329,6 +314,11 @@ spring-test-dbunit test + + org.testcontainers + cassandra + test + org.testcontainers postgresql diff --git a/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java b/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java index 47355d6c44..9a6ffb0e32 100644 --- a/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java +++ b/application/src/main/java/org/thingsboard/server/config/WebSocketConfiguration.java @@ -15,6 +15,8 @@ */ package org.thingsboard.server.config; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.http.HttpStatus; @@ -40,11 +42,15 @@ import java.util.Map; @Configuration @TbCoreComponent @EnableWebSocket +@RequiredArgsConstructor +@Slf4j public class WebSocketConfiguration implements WebSocketConfigurer { public static final String WS_PLUGIN_PREFIX = "/api/ws/plugins/"; private static final String WS_PLUGIN_MAPPING = WS_PLUGIN_PREFIX + "**"; + private final WebSocketHandler wsHandler; + @Bean public ServletServerContainerFactoryBean createWebSocketContainer() { ServletServerContainerFactoryBean container = new ServletServerContainerFactoryBean(); @@ -55,7 +61,11 @@ public class WebSocketConfiguration implements WebSocketConfigurer { @Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { - registry.addHandler(wsHandler(), WS_PLUGIN_MAPPING).setAllowedOriginPatterns("*") + if (!(wsHandler instanceof TbWebSocketHandler)) { + log.error("TbWebSocketHandler expected but [{}] provided", wsHandler); + throw new RuntimeException("TbWebSocketHandler expected but " + wsHandler + " provided"); + } + registry.addHandler(wsHandler, WS_PLUGIN_MAPPING).setAllowedOriginPatterns("*") .addInterceptors(new HttpSessionHandshakeInterceptor(), new HandshakeInterceptor() { @Override @@ -82,11 +92,6 @@ public class WebSocketConfiguration implements WebSocketConfigurer { }); } - @Bean - public WebSocketHandler wsHandler() { - return new TbWebSocketHandler(); - } - protected SecurityUser getCurrentUser() throws ThingsboardException { Authentication authentication = SecurityContextHolder.getContext().getAuthentication(); if (authentication != null && authentication.getPrincipal() instanceof SecurityUser) { diff --git a/application/src/main/java/org/thingsboard/server/controller/AuthController.java b/application/src/main/java/org/thingsboard/server/controller/AuthController.java index c73ac1ae66..113c292380 100644 --- a/application/src/main/java/org/thingsboard/server/controller/AuthController.java +++ b/application/src/main/java/org/thingsboard/server/controller/AuthController.java @@ -20,6 +20,7 @@ import io.swagger.annotations.ApiOperation; import io.swagger.annotations.ApiParam; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; import org.springframework.context.ApplicationEventPublisher; import org.springframework.http.HttpHeaders; import org.springframework.http.HttpStatus; @@ -41,11 +42,13 @@ import org.thingsboard.server.common.data.edge.EdgeEventActionType; import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.security.UserCredentials; import org.thingsboard.server.common.data.security.event.UserCredentialsInvalidationEvent; import org.thingsboard.server.common.data.security.event.UserSessionInvalidationEvent; import org.thingsboard.server.common.data.security.model.SecuritySettings; import org.thingsboard.server.common.data.security.model.UserPasswordPolicy; +import org.thingsboard.server.common.msg.tools.TbRateLimits; import org.thingsboard.server.dao.audit.AuditLogService; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.security.auth.rest.RestAuthenticationDetails; @@ -62,6 +65,8 @@ import org.thingsboard.server.service.security.system.SystemSecurityService; import javax.servlet.http.HttpServletRequest; import java.net.URI; import java.net.URISyntaxException; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; @RestController @TbCoreComponent @@ -69,6 +74,10 @@ import java.net.URISyntaxException; @Slf4j @RequiredArgsConstructor public class AuthController extends BaseController { + + @Value("${server.rest.rate_limits.reset_password_per_user:5:3600}") + private String defaultLimitsConfiguration; + private final ConcurrentMap resetPasswordRateLimits = new ConcurrentHashMap<>(); private final BCryptPasswordEncoder passwordEncoder; private final JwtTokenFactory tokenFactory; private final MailService mailService; @@ -211,7 +220,12 @@ public class AuthController extends BaseController { HttpStatus responseStatus; String resetURI = "/login/resetPassword"; UserCredentials userCredentials = userService.findUserCredentialsByResetToken(TenantId.SYS_TENANT_ID, resetToken); + if (userCredentials != null) { + TbRateLimits tbRateLimits = getTbRateLimits(userCredentials.getUserId()); + if (!tbRateLimits.tryConsume()) { + return ResponseEntity.status(HttpStatus.TOO_MANY_REQUESTS).build(); + } try { URI location = new URI(resetURI + "?resetToken=" + resetToken); headers.setLocation(location); @@ -323,4 +337,9 @@ public class AuthController extends BaseController { throw handleException(e); } } + + private TbRateLimits getTbRateLimits(UserId userId) { + return resetPasswordRateLimits.computeIfAbsent(userId, + key -> new TbRateLimits(defaultLimitsConfiguration, true)); + } } diff --git a/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java b/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java index c689ae3593..f43607b24c 100644 --- a/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java +++ b/application/src/main/java/org/thingsboard/server/controller/plugin/TbWebSocketHandler.java @@ -67,8 +67,8 @@ import static org.thingsboard.server.service.telemetry.DefaultTelemetryWebSocket @Slf4j public class TbWebSocketHandler extends TextWebSocketHandler implements TelemetryWebSocketMsgEndpoint { - private static final ConcurrentMap internalSessionMap = new ConcurrentHashMap<>(); - private static final ConcurrentMap externalSessionMap = new ConcurrentHashMap<>(); + private final ConcurrentMap internalSessionMap = new ConcurrentHashMap<>(); + private final ConcurrentMap externalSessionMap = new ConcurrentHashMap<>(); @Autowired @@ -82,13 +82,13 @@ public class TbWebSocketHandler extends TextWebSocketHandler implements Telemetr @Value("${server.ws.ping_timeout:30000}") private long pingTimeout; - private ConcurrentMap blacklistedSessions = new ConcurrentHashMap<>(); - private ConcurrentMap perSessionUpdateLimits = new ConcurrentHashMap<>(); + private final ConcurrentMap blacklistedSessions = new ConcurrentHashMap<>(); + private final ConcurrentMap perSessionUpdateLimits = new ConcurrentHashMap<>(); - private ConcurrentMap> tenantSessionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> customerSessionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> regularUserSessionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> publicUserSessionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> tenantSessionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> customerSessionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> regularUserSessionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> publicUserSessionsMap = new ConcurrentHashMap<>(); @Override public void handleTextMessage(WebSocketSession session, TextMessage message) { diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java index a2a0d0fd8e..ffde6192ec 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetryWebSocketService.java @@ -16,7 +16,6 @@ package org.thingsboard.server.service.telemetry; import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.ObjectMapper; import com.google.common.base.Function; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; @@ -27,6 +26,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.springframework.web.socket.CloseStatus; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.DataConstants; @@ -95,6 +95,8 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; import java.util.function.Consumer; import java.util.stream.Collectors; @@ -112,7 +114,6 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi private static final Aggregation DEFAULT_AGGREGATION = Aggregation.NONE; private static final int UNKNOWN_SUBSCRIPTION_ID = 0; private static final String PROCESSING_MSG = "[{}] Processing: {}"; - private static final ObjectMapper jsonMapper = new ObjectMapper(); private static final String FAILED_TO_FETCH_DATA = "Failed to fetch data!"; private static final String FAILED_TO_FETCH_ATTRIBUTES = "Failed to fetch attributes!"; private static final String SESSION_META_DATA_NOT_FOUND = "Session meta-data not found!"; @@ -147,10 +148,10 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi @Value("${server.ws.ping_timeout:30000}") private long pingTimeout; - private ConcurrentMap> tenantSubscriptionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> customerSubscriptionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> regularUserSubscriptionsMap = new ConcurrentHashMap<>(); - private ConcurrentMap> publicUserSubscriptionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> tenantSubscriptionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> customerSubscriptionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> regularUserSubscriptionsMap = new ConcurrentHashMap<>(); + private final ConcurrentMap> publicUserSubscriptionsMap = new ConcurrentHashMap<>(); private ExecutorService executor; private String serviceId; @@ -204,7 +205,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi } try { - TelemetryPluginCmdsWrapper cmdsWrapper = jsonMapper.readValue(msg, TelemetryPluginCmdsWrapper.class); + TelemetryPluginCmdsWrapper cmdsWrapper = JacksonUtil.OBJECT_MAPPER.readValue(msg, TelemetryPluginCmdsWrapper.class); if (cmdsWrapper != null) { if (cmdsWrapper.getAttrSubCmds() != null) { cmdsWrapper.getAttrSubCmds().forEach(cmd -> { @@ -450,7 +451,6 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi @Override public void onSuccess(List data) { List attributesData = data.stream().map(d -> new BasicTsKvEntry(d.getLastUpdateTs(), d)).collect(Collectors.toList()); - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); Map subState = new HashMap<>(keys.size()); keys.forEach(key -> subState.put(key, 0L)); @@ -458,6 +458,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi TbAttributeSubscriptionScope scope = StringUtils.isEmpty(cmd.getScope()) ? TbAttributeSubscriptionScope.ANY_SCOPE : TbAttributeSubscriptionScope.valueOf(cmd.getScope()); + Lock subLock = new ReentrantLock(); TbAttributeSubscription sub = TbAttributeSubscription.builder() .serviceId(serviceId) .sessionId(sessionId) @@ -467,9 +468,24 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .allKeys(false) .keyStates(subState) .scope(scope) - .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) + .updateConsumer((sessionId, update) -> { + subLock.lock(); + try { + sendWsMsg(sessionId, update); + } finally { + subLock.unlock(); + } + }) .build(); - oldSubService.addSubscription(sub); + + subLock.lock(); + try{ + oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); + } finally { + subLock.unlock(); + } + } @Override @@ -550,13 +566,13 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi @Override public void onSuccess(List data) { List attributesData = data.stream().map(d -> new BasicTsKvEntry(d.getLastUpdateTs(), d)).collect(Collectors.toList()); - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); Map subState = new HashMap<>(attributesData.size()); attributesData.forEach(v -> subState.put(v.getKey(), v.getTs())); TbAttributeSubscriptionScope scope = StringUtils.isEmpty(cmd.getScope()) ? TbAttributeSubscriptionScope.ANY_SCOPE : TbAttributeSubscriptionScope.valueOf(cmd.getScope()); + Lock subLock = new ReentrantLock(); TbAttributeSubscription sub = TbAttributeSubscription.builder() .serviceId(serviceId) .sessionId(sessionId) @@ -565,9 +581,24 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi .entityId(entityId) .allKeys(true) .keyStates(subState) - .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) - .scope(scope).build(); - oldSubService.addSubscription(sub); + .updateConsumer((sessionId, update) -> { + subLock.lock(); + try { + sendWsMsg(sessionId, update); + } finally { + subLock.unlock(); + } + }) + .scope(scope) + .build(); + + subLock.lock(); + try { + oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); + } finally { + subLock.unlock(); + } } @Override @@ -636,20 +667,34 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi FutureCallback> callback = new FutureCallback>() { @Override public void onSuccess(List data) { - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); Map subState = new HashMap<>(data.size()); data.forEach(v -> subState.put(v.getKey(), v.getTs())); + Lock subLock = new ReentrantLock(); TbTimeseriesSubscription sub = TbTimeseriesSubscription.builder() .serviceId(serviceId) .sessionId(sessionId) .subscriptionId(cmd.getCmdId()) .tenantId(sessionRef.getSecurityCtx().getTenantId()) .entityId(entityId) - .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) + .updateConsumer((sessionId, update) -> { + subLock.lock(); + try { + sendWsMsg(sessionId, update); + } finally { + subLock.unlock(); + } + }) .allKeys(true) .keyStates(subState).build(); - oldSubService.addSubscription(sub); + + subLock.lock(); + try { + oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); + } finally { + subLock.unlock(); + } } @Override @@ -673,21 +718,35 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi return new FutureCallback<>() { @Override public void onSuccess(List data) { - sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); Map subState = new HashMap<>(keys.size()); keys.forEach(key -> subState.put(key, startTs)); data.forEach(v -> subState.put(v.getKey(), v.getTs())); + Lock subLock = new ReentrantLock(); TbTimeseriesSubscription sub = TbTimeseriesSubscription.builder() .serviceId(serviceId) .sessionId(sessionId) .subscriptionId(cmd.getCmdId()) .tenantId(sessionRef.getSecurityCtx().getTenantId()) .entityId(entityId) - .updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) + .updateConsumer((sessionId, update) -> { + subLock.lock(); + try { + sendWsMsg(sessionId, update); + } finally { + subLock.unlock(); + } + }) .allKeys(false) .keyStates(subState).build(); - oldSubService.addSubscription(sub); + + subLock.lock(); + try{ + oldSubService.addSubscription(sub); + sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); + } finally { + subLock.unlock(); + } } @Override @@ -793,7 +852,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi private void sendWsMsg(TelemetryWebSocketSessionRef sessionRef, int cmdId, Object update) { try { - String msg = jsonMapper.writeValueAsString(update); + String msg = JacksonUtil.OBJECT_MAPPER.writeValueAsString(update); executor.submit(() -> { try { msgEndpoint.send(sessionRef, cmdId, msg); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 61bef18bcd..10c4c9edac 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -73,6 +73,8 @@ server: min_timeout: "${MIN_SERVER_SIDE_RPC_TIMEOUT:5000}" # Default value of the server side RPC timeout. default_timeout: "${DEFAULT_SERVER_SIDE_RPC_TIMEOUT:10000}" + rate_limits: + reset_password_per_user: "${RESET_PASSWORD_PER_USER_RATE_LIMIT_CONFIGURATION:5:3600}" # Application info app: @@ -1210,3 +1212,4 @@ management: exposure: # Expose metrics endpoint (use value 'prometheus' to enable prometheus metrics). include: '${METRICS_ENDPOINTS_EXPOSE:info}' + 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 06cee7b71d..32d1f65765 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractControllerTest.java @@ -52,7 +52,7 @@ public abstract class AbstractControllerTest extends AbstractNotifyEntityTest { @LocalServerPort protected int wsPort; - private TbTestWebSocketClient wsClient; // lazy + private volatile TbTestWebSocketClient wsClient; // lazy public TbTestWebSocketClient getWsClient() { if (wsClient == null) { diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java index aa4febc303..4afb12e638 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java @@ -211,17 +211,6 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest { this.mappingJackson2HttpMessageConverter); } - @BeforeClass - public static void beforeWebTestClass() throws Exception { - - } - - @AfterClass - public static void afterWebTestClass() throws Exception { - Mockito.clearAllCaches(); - } - - @Before public void setupWebTest() throws Exception { log.debug("Executing web test setup"); diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java index 09fb5c505d..ce354749f3 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java @@ -548,7 +548,7 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { SingleEntityFilter entityFilter = new SingleEntityFilter(); entityFilter.setSingleEntity(tenantId); - assertThatNoException().isThrownBy(() -> { + assertThatNoException().as("subscribeForAttributes").isThrownBy(() -> { JsonNode update = getWsClient().subscribeForAttributes(tenantId, TbAttributeSubscriptionScope.SERVER_SCOPE.name(), List.of("attr")); assertThat(update.get("errorMsg").isNull()).isTrue(); assertThat(update.get("errorCode").asInt()).isEqualTo(SubscriptionErrorCode.NO_ERROR.getCode()); @@ -560,7 +560,7 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { new BaseAttributeKvEntry(System.currentTimeMillis(), new StringDataEntry("attr", expectedAttrValue)) )); JsonNode update = JacksonUtil.toJsonNode(getWsClient().waitForUpdate()); - assertThat(update).isNotNull(); + assertThat(update).as("waitForUpdate").isNotNull(); assertThat(update.get("data").get("attr").get(0).get(1).asText()).isEqualTo(expectedAttrValue); } @@ -569,15 +569,17 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { tsService.saveAndNotify(device.getTenantId(), null, device.getId(), tsData, 0, new FutureCallback() { @Override public void onSuccess(@Nullable Void result) { + log.debug("sendTelemetry callback onSuccess"); latch.countDown(); } @Override public void onFailure(Throwable t) { + log.error("Failed to send telemetry", t); latch.countDown(); } }); - latch.await(3, TimeUnit.SECONDS); + assertThat(latch.await(TIMEOUT, TimeUnit.SECONDS)).as("await sendTelemetry callback"); } private void sendAttributes(Device device, TbAttributeSubscriptionScope scope, List attrData) throws InterruptedException { @@ -589,14 +591,16 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { tsService.saveAndNotify(tenantId, entityId, scope.name(), attrData, new FutureCallback() { @Override public void onSuccess(@Nullable Void result) { + log.debug("sendAttributes callback onSuccess"); latch.countDown(); } @Override public void onFailure(Throwable t) { + log.error("Failed to sendAttributes", t); latch.countDown(); } }); - latch.await(3, TimeUnit.SECONDS); + assertThat(latch.await(TIMEOUT, TimeUnit.SECONDS)).as("await sendAttributes callback").isTrue(); } } diff --git a/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java b/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java index db26e9f6df..1f9ac4ea69 100644 --- a/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java +++ b/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java @@ -47,6 +47,7 @@ import java.util.concurrent.TimeUnit; @Slf4j public class TbTestWebSocketClient extends WebSocketClient { + private static final long TIMEOUT = TimeUnit.SECONDS.toMillis(30); private volatile String lastMsg; private volatile CountDownLatch reply; private volatile CountDownLatch update; @@ -87,12 +88,14 @@ public class TbTestWebSocketClient extends WebSocketClient { } public void registerWaitForUpdate(int count) { + log.debug("registerWaitForUpdate [{}]", count); lastMsg = null; update = new CountDownLatch(count); } @Override public void send(String text) throws NotYetConnectedException { + log.debug("send [{}]", text); reply = new CountDownLatch(1); super.send(text); } @@ -110,21 +113,31 @@ public class TbTestWebSocketClient extends WebSocketClient { } public String waitForUpdate() { - return waitForUpdate(TimeUnit.SECONDS.toMillis(3)); + return waitForUpdate(TIMEOUT); } public String waitForUpdate(long ms) { + log.debug("waitForUpdate [{}]", ms); try { - update.await(ms, TimeUnit.MILLISECONDS); + if (!update.await(ms, TimeUnit.MILLISECONDS)) { + log.warn("Failed to await update (waiting time [{}]ms elapsed)", ms, new RuntimeException("stacktrace")); + } } catch (InterruptedException e) { - log.warn("Failed to await reply", e); + log.warn("Failed to await update", e); } return lastMsg; } public String waitForReply() { + return waitForReply(TIMEOUT); + } + + public String waitForReply(long ms) { + log.debug("waitForReply [{}]", ms); try { - reply.await(3, TimeUnit.SECONDS); + if (!reply.await(ms, TimeUnit.MILLISECONDS)) { + log.warn("Failed to await reply (waiting time [{}]ms elapsed)", ms, new RuntimeException("stacktrace")); + } } catch (InterruptedException e) { log.warn("Failed to await reply", e); } diff --git a/application/src/test/java/org/thingsboard/server/transport/TransportNoSqlTestSuite.java b/application/src/test/java/org/thingsboard/server/transport/TransportNoSqlTestSuite.java index 41cc4c2b76..ab2f6b9bca 100644 --- a/application/src/test/java/org/thingsboard/server/transport/TransportNoSqlTestSuite.java +++ b/application/src/test/java/org/thingsboard/server/transport/TransportNoSqlTestSuite.java @@ -15,30 +15,14 @@ */ package org.thingsboard.server.transport; -import org.cassandraunit.dataset.cql.ClassPathCQLDataSet; -import org.junit.BeforeClass; -import org.junit.ClassRule; import org.junit.extensions.cpsuite.ClasspathSuite; import org.junit.runner.RunWith; -import org.thingsboard.server.dao.CustomCassandraCQLUnit; -import org.thingsboard.server.queue.memory.InMemoryStorage; - -import java.util.Arrays; +import org.thingsboard.server.dao.AbstractNoSqlContainer; @RunWith(ClasspathSuite.class) @ClasspathSuite.ClassnameFilters({ "org.thingsboard.server.transport.*.telemetry.timeseries.nosql.*Test", }) -public class TransportNoSqlTestSuite { - - @ClassRule - public static CustomCassandraCQLUnit cassandraUnit = - new CustomCassandraCQLUnit( - Arrays.asList( - new ClassPathCQLDataSet("cassandra/schema-keyspace.cql", false, false), - new ClassPathCQLDataSet("cassandra/schema-ts.cql", false, false), - new ClassPathCQLDataSet("cassandra/schema-ts-latest.cql", false, false) - ), - "cassandra-test.yaml", 30000l); +public class TransportNoSqlTestSuite extends AbstractNoSqlContainer { } diff --git a/application/src/test/resources/application-test.properties b/application/src/test/resources/application-test.properties index a4bedd9527..ad86ff736b 100644 --- a/application/src/test/resources/application-test.properties +++ b/application/src/test/resources/application-test.properties @@ -67,4 +67,3 @@ sql.ttl.audit_logs.ttl=2592000 sql.edge_events.partition_size=168 sql.ttl.edge_events.edge_event_ttl=2592000 -spring.test.context.cache.maxSize=1 \ No newline at end of file diff --git a/application/src/test/resources/logback-test.xml b/application/src/test/resources/logback-test.xml index c79a7d3905..23e6d8c2f6 100644 --- a/application/src/test/resources/logback-test.xml +++ b/application/src/test/resources/logback-test.xml @@ -13,10 +13,11 @@ - + + diff --git a/common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java b/common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java index a9f262d6ca..c5af954392 100644 --- a/common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java +++ b/common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java @@ -108,6 +108,10 @@ public abstract class RedisTbTransactionalCache keys) { + //Redis expects at least 1 key to delete. Otherwise - ERR wrong number of arguments for 'del' command + if (keys.isEmpty()) { + return; + } try (var connection = connectionFactory.getConnection()) { connection.del(keys.stream().map(this::getRawKey).toArray(byte[][]::new)); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java b/common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java index 6a2b0a58d6..3b38aa57c1 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/StringUtils.java @@ -18,9 +18,14 @@ package org.thingsboard.server.common.data; import com.google.common.base.Splitter; import org.apache.commons.lang3.RandomStringUtils; +import java.security.SecureRandom; +import java.util.Base64; + import static org.apache.commons.lang3.StringUtils.repeat; public class StringUtils { + public static final SecureRandom RANDOM = new SecureRandom(); + public static final String EMPTY = ""; public static final int INDEX_NOT_FOUND = -1; @@ -180,4 +185,11 @@ public class StringUtils { return RandomStringUtils.randomAlphabetic(count); } + public static String generateSafeToken(int length) { + byte[] bytes = new byte[length]; + RANDOM.nextBytes(bytes); + Base64.Encoder encoder = Base64.getUrlEncoder().withoutPadding(); + return encoder.encodeToString(bytes); + } + } diff --git a/dao/pom.xml b/dao/pom.xml index 30dc1adcf0..4fc8a3c08d 100644 --- a/dao/pom.xml +++ b/dao/pom.xml @@ -162,26 +162,6 @@ com.google.guava guava - - org.cassandraunit - cassandra-unit - - - org.slf4j - slf4j-log4j12 - - - org.hibernate - hibernate-validator - - - test - - - org.apache.cassandra - cassandra-thrift - test - com.google.protobuf protobuf-java @@ -211,6 +191,11 @@ spring-test test + + org.testcontainers + cassandra + test + org.testcontainers postgresql diff --git a/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java b/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java index 584ae84531..017fe4315c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/attributes/CachedAttributesService.java @@ -66,7 +66,7 @@ public class CachedAttributesService implements AttributesService { private final TbTransactionalCache cache; private ListeningExecutorService cacheExecutor; - @Value("${cache.type}") + @Value("${cache.type:caffeine}") private String cacheType; public CachedAttributesService(AttributesDao attributesDao, diff --git a/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java index c8c5175e5f..b4ad62c0fc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/user/UserServiceImpl.java @@ -29,7 +29,6 @@ import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.EntityId; @@ -40,7 +39,6 @@ import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.security.UserCredentials; -import org.thingsboard.server.common.data.security.UserSettings; import org.thingsboard.server.common.data.security.event.UserCredentialsInvalidationEvent; import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.exception.IncorrectParameterException; @@ -51,6 +49,7 @@ import java.util.HashMap; import java.util.Map; import java.util.Optional; +import static org.thingsboard.server.common.data.StringUtils.generateSafeToken; import static org.thingsboard.server.dao.service.Validator.validateId; import static org.thingsboard.server.dao.service.Validator.validatePageLink; import static org.thingsboard.server.dao.service.Validator.validateString; @@ -126,7 +125,7 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic if (user.getId() == null) { UserCredentials userCredentials = new UserCredentials(); userCredentials.setEnabled(false); - userCredentials.setActivateToken(StringUtils.randomAlphanumeric(DEFAULT_TOKEN_LENGTH)); + userCredentials.setActivateToken(generateSafeToken(DEFAULT_TOKEN_LENGTH)); userCredentials.setUserId(new UserId(savedUser.getUuidId())); saveUserCredentialsAndPasswordHistory(user.getTenantId(), userCredentials); } @@ -192,7 +191,7 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic if (!userCredentials.isEnabled()) { throw new DisabledException(String.format("User credentials not enabled [%s]", email)); } - userCredentials.setResetToken(StringUtils.randomAlphanumeric(DEFAULT_TOKEN_LENGTH)); + userCredentials.setResetToken(generateSafeToken(DEFAULT_TOKEN_LENGTH)); return saveUserCredentials(tenantId, userCredentials); } @@ -202,7 +201,7 @@ public class UserServiceImpl extends AbstractEntityService implements UserServic if (!userCredentials.isEnabled()) { throw new IncorrectParameterException("Unable to reset password for inactive user"); } - userCredentials.setResetToken(StringUtils.randomAlphanumeric(DEFAULT_TOKEN_LENGTH)); + userCredentials.setResetToken(generateSafeToken(DEFAULT_TOKEN_LENGTH)); return saveUserCredentials(tenantId, userCredentials); } diff --git a/dao/src/main/resources/cassandra/schema-ts-latest.cql b/dao/src/main/resources/cassandra/schema-ts-latest.cql index e3c72ab6da..6b76c6187c 100644 --- a/dao/src/main/resources/cassandra/schema-ts-latest.cql +++ b/dao/src/main/resources/cassandra/schema-ts-latest.cql @@ -15,7 +15,7 @@ -- CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_latest_cf ( - entity_type text, // (DEVICE, CUSTOMER, TENANT) + entity_type text, -- (DEVICE, CUSTOMER, TENANT) entity_id timeuuid, key text, ts bigint, diff --git a/dao/src/main/resources/cassandra/schema-ts.cql b/dao/src/main/resources/cassandra/schema-ts.cql index fb60ed9e78..ae8c11a703 100644 --- a/dao/src/main/resources/cassandra/schema-ts.cql +++ b/dao/src/main/resources/cassandra/schema-ts.cql @@ -15,7 +15,7 @@ -- CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_cf ( - entity_type text, // (DEVICE, CUSTOMER, TENANT) + entity_type text, -- (DEVICE, CUSTOMER, TENANT) entity_id timeuuid, key text, partition bigint, @@ -29,7 +29,7 @@ CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_cf ( ); CREATE TABLE IF NOT EXISTS thingsboard.ts_kv_partitions_cf ( - entity_type text, // (DEVICE, CUSTOMER, TENANT) + entity_type text, -- (DEVICE, CUSTOMER, TENANT) entity_id timeuuid, key text, partition bigint, diff --git a/dao/src/test/java/org/apache/cassandra/io/sstable/Descriptor.java b/dao/src/test/java/org/apache/cassandra/io/sstable/Descriptor.java deleted file mode 100644 index 69a164f729..0000000000 --- a/dao/src/test/java/org/apache/cassandra/io/sstable/Descriptor.java +++ /dev/null @@ -1,366 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you 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.apache.cassandra.io.sstable; - -import java.io.File; -import java.io.IOError; -import java.io.IOException; -import java.util.*; -import java.util.regex.Pattern; - -import com.google.common.annotations.VisibleForTesting; -import com.google.common.base.CharMatcher; -import com.google.common.base.Objects; - -import org.apache.cassandra.db.Directories; -import org.apache.cassandra.io.sstable.format.SSTableFormat; -import org.apache.cassandra.io.sstable.format.Version; -import org.apache.cassandra.io.sstable.metadata.IMetadataSerializer; -import org.apache.cassandra.io.sstable.metadata.LegacyMetadataSerializer; -import org.apache.cassandra.io.sstable.metadata.MetadataSerializer; -import org.apache.cassandra.utils.Pair; - -import static org.apache.cassandra.io.sstable.Component.separator; - -/** - * A SSTable is described by the keyspace and column family it contains data - * for, a generation (where higher generations contain more recent data) and - * an alphabetic version string. - * - * A descriptor can be marked as temporary, which influences generated filenames. - */ -public class Descriptor -{ - public static String TMP_EXT = ".tmp"; - - /** canonicalized path to the directory where SSTable resides */ - public final File directory; - /** version has the following format: [a-z]+ */ - public final Version version; - public final String ksname; - public final String cfname; - public final int generation; - public final SSTableFormat.Type formatType; - /** digest component - might be {@code null} for old, legacy sstables */ - public final Component digestComponent; - private final int hashCode; - - /** - * A descriptor that assumes CURRENT_VERSION. - */ - @VisibleForTesting - public Descriptor(File directory, String ksname, String cfname, int generation) - { - this(SSTableFormat.Type.current().info.getLatestVersion(), directory, ksname, cfname, generation, SSTableFormat.Type.current(), null); - } - - /** - * Constructor for sstable writers only. - */ - public Descriptor(File directory, String ksname, String cfname, int generation, SSTableFormat.Type formatType) - { - this(formatType.info.getLatestVersion(), directory, ksname, cfname, generation, formatType, Component.digestFor(formatType.info.getLatestVersion().uncompressedChecksumType())); - } - - @VisibleForTesting - public Descriptor(String version, File directory, String ksname, String cfname, int generation, SSTableFormat.Type formatType) - { - this(formatType.info.getVersion(version), directory, ksname, cfname, generation, formatType, Component.digestFor(formatType.info.getLatestVersion().uncompressedChecksumType())); - } - - public Descriptor(Version version, File directory, String ksname, String cfname, int generation, SSTableFormat.Type formatType, Component digestComponent) - { - assert version != null && directory != null && ksname != null && cfname != null && formatType.info.getLatestVersion().getClass().equals(version.getClass()); - this.version = version; - try - { - this.directory = directory.getCanonicalFile(); - } - catch (IOException e) - { - throw new IOError(e); - } - this.ksname = ksname; - this.cfname = cfname; - this.generation = generation; - this.formatType = formatType; - this.digestComponent = digestComponent; - - hashCode = Objects.hashCode(version, this.directory, generation, ksname, cfname, formatType); - } - - public Descriptor withGeneration(int newGeneration) - { - return new Descriptor(version, directory, ksname, cfname, newGeneration, formatType, digestComponent); - } - - public Descriptor withFormatType(SSTableFormat.Type newType) - { - return new Descriptor(newType.info.getLatestVersion(), directory, ksname, cfname, generation, newType, digestComponent); - } - - public Descriptor withDigestComponent(Component newDigestComponent) - { - return new Descriptor(version, directory, ksname, cfname, generation, formatType, newDigestComponent); - } - - public String tmpFilenameFor(Component component) - { - return filenameFor(component) + TMP_EXT; - } - - public String filenameFor(Component component) - { - return baseFilename() + separator + component.name(); - } - - public String baseFilename() - { - StringBuilder buff = new StringBuilder(); - buff.append(directory).append(File.separatorChar); - appendFileName(buff); - return buff.toString(); - } - - private void appendFileName(StringBuilder buff) - { - if (!version.hasNewFileName()) - { - buff.append(ksname).append(separator); - buff.append(cfname).append(separator); - } - buff.append(version).append(separator); - buff.append(generation); - if (formatType != SSTableFormat.Type.LEGACY) - buff.append(separator).append(formatType.name); - } - - public String relativeFilenameFor(Component component) - { - final StringBuilder buff = new StringBuilder(); - appendFileName(buff); - buff.append(separator).append(component.name()); - return buff.toString(); - } - - public SSTableFormat getFormat() - { - return formatType.info; - } - - /** Return any temporary files found in the directory */ - public List getTemporaryFiles() - { - List ret = new ArrayList<>(); - File[] tmpFiles = directory.listFiles((dir, name) -> - name.endsWith(Descriptor.TMP_EXT)); - - for (File tmpFile : tmpFiles) - ret.add(tmpFile); - - return ret; - } - - /** - * Files obsoleted by CASSANDRA-7066 : temporary files and compactions_in_progress. We support - * versions 2.1 (ka) and 2.2 (la). - * Temporary files have tmp- or tmplink- at the beginning for 2.2 sstables or after ks-cf- for 2.1 sstables - */ - - private final static String LEGACY_COMP_IN_PROG_REGEX_STR = "^compactions_in_progress(\\-[\\d,a-f]{32})?$"; - private final static Pattern LEGACY_COMP_IN_PROG_REGEX = Pattern.compile(LEGACY_COMP_IN_PROG_REGEX_STR); - private final static String LEGACY_TMP_REGEX_STR = "^((.*)\\-(.*)\\-)?tmp(link)?\\-((?:l|k).)\\-(\\d)*\\-(.*)$"; - private final static Pattern LEGACY_TMP_REGEX = Pattern.compile(LEGACY_TMP_REGEX_STR); - - public static boolean isLegacyFile(File file) - { - if (file.isDirectory()) - return file.getParentFile() != null && - file.getParentFile().getName().equalsIgnoreCase("system") && - LEGACY_COMP_IN_PROG_REGEX.matcher(file.getName()).matches(); - else - return LEGACY_TMP_REGEX.matcher(file.getName()).matches(); - } - - public static boolean isValidFile(String fileName) - { - return fileName.endsWith(".db") && !LEGACY_TMP_REGEX.matcher(fileName).matches(); - } - - /** - * @see #fromFilename(File directory, String name) - * @param filename The SSTable filename - * @return Descriptor of the SSTable initialized from filename - */ - public static Descriptor fromFilename(String filename) - { - return fromFilename(filename, false); - } - - public static Descriptor fromFilename(String filename, SSTableFormat.Type formatType) - { - return fromFilename(filename).withFormatType(formatType); - } - - public static Descriptor fromFilename(String filename, boolean skipComponent) - { - File file = new File(filename).getAbsoluteFile(); - return fromFilename(file.getParentFile(), file.getName(), skipComponent).left; - } - - public static Pair fromFilename(File directory, String name) - { - return fromFilename(directory, name, false); - } - - /** - * Filename of the form is vary by version: - * - *
    - *
  • <ksname>-<cfname>-(tmp-)?<version>-<gen>-<component> for cassandra 2.0 and before
  • - *
  • (<tmp marker>-)?<version>-<gen>-<component> for cassandra 3.0 and later
  • - *
- * - * If this is for SSTable of secondary index, directory should ends with index name for 2.1+. - * - * @param directory The directory of the SSTable files - * @param name The name of the SSTable file - * @param skipComponent true if the name param should not be parsed for a component tag - * - * @return A Descriptor for the SSTable, and the Component remainder. - */ - @SuppressWarnings("deprecation") - public static Pair fromFilename(File directory, String name, boolean skipComponent) - { - File parentDirectory = directory != null ? directory : new File("."); - - // tokenize the filename - StringTokenizer st = new StringTokenizer(name, String.valueOf(separator)); - String nexttok; - - // read tokens backwards to determine version - Deque tokenStack = new ArrayDeque<>(); - while (st.hasMoreTokens()) - { - tokenStack.push(st.nextToken()); - } - - // component suffix - String component = skipComponent ? null : tokenStack.pop(); - - nexttok = tokenStack.pop(); - // generation OR format type - SSTableFormat.Type fmt = SSTableFormat.Type.LEGACY; - if (!CharMatcher.digit().matchesAllOf(nexttok)) - { - fmt = SSTableFormat.Type.validate(nexttok); - nexttok = tokenStack.pop(); - } - - // generation - int generation = Integer.parseInt(nexttok); - - // version - nexttok = tokenStack.pop(); - - if (!Version.validate(nexttok)) - throw new UnsupportedOperationException("SSTable " + name + " is too old to open. Upgrade to 2.0 first, and run upgradesstables"); - - Version version = fmt.info.getVersion(nexttok); - - // ks/cf names - String ksname, cfname; - if (version.hasNewFileName()) - { - // for 2.1+ read ks and cf names from directory - File cfDirectory = parentDirectory; - // check if this is secondary index - String indexName = ""; - if (cfDirectory.getName().startsWith(Directories.SECONDARY_INDEX_NAME_SEPARATOR)) - { - indexName = cfDirectory.getName(); - cfDirectory = cfDirectory.getParentFile(); - } - if (cfDirectory.getName().equals(Directories.BACKUPS_SUBDIR)) - { - cfDirectory = cfDirectory.getParentFile(); - } - else if (cfDirectory.getParentFile().getName().equals(Directories.SNAPSHOT_SUBDIR)) - { - cfDirectory = cfDirectory.getParentFile().getParentFile(); - } - cfname = cfDirectory.getName().split("-")[0] + indexName; - ksname = cfDirectory.getParentFile().getName(); - } - else - { - cfname = tokenStack.pop(); - ksname = tokenStack.pop(); - } - assert tokenStack.isEmpty() : "Invalid file name " + name + " in " + directory; - - return Pair.create(new Descriptor(version, parentDirectory, ksname, cfname, generation, fmt, - // _assume_ version from version - Component.digestFor(version.uncompressedChecksumType())), - component); - } - - @SuppressWarnings("deprecation") - public IMetadataSerializer getMetadataSerializer() - { - if (version.hasNewStatsFile()) - return new MetadataSerializer(); - else - return new LegacyMetadataSerializer(); - } - - /** - * @return true if the current Cassandra version can read the given sstable version - */ - public boolean isCompatible() - { - return version.isCompatible(); - } - - @Override - public String toString() - { - return baseFilename(); - } - - @Override - public boolean equals(Object o) - { - if (o == this) - return true; - if (!(o instanceof Descriptor)) - return false; - Descriptor that = (Descriptor)o; - return that.directory.equals(this.directory) - && that.generation == this.generation - && that.ksname.equals(this.ksname) - && that.cfname.equals(this.cfname) - && that.formatType == this.formatType; - } - - @Override - public int hashCode() - { - return hashCode; - } -} diff --git a/dao/src/test/java/org/apache/cassandra/io/sstable/format/SSTableFormat.java b/dao/src/test/java/org/apache/cassandra/io/sstable/format/SSTableFormat.java deleted file mode 100644 index 350d27591f..0000000000 --- a/dao/src/test/java/org/apache/cassandra/io/sstable/format/SSTableFormat.java +++ /dev/null @@ -1,86 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you 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.apache.cassandra.io.sstable.format; - -import com.google.common.base.CharMatcher; -import org.apache.cassandra.config.CFMetaData; -import org.apache.cassandra.db.RowIndexEntry; -import org.apache.cassandra.db.SerializationHeader; -import org.apache.cassandra.io.sstable.format.big.BigFormat; - -/** - * Provides the accessors to data on disk. - */ -public interface SSTableFormat -{ - static boolean enableSSTableDevelopmentTestMode = Boolean.getBoolean("cassandra.test.sstableformatdevelopment"); - - - Version getLatestVersion(); - Version getVersion(String version); - - SSTableWriter.Factory getWriterFactory(); - SSTableReader.Factory getReaderFactory(); - - RowIndexEntry.IndexSerializer getIndexSerializer(CFMetaData cfm, Version version, SerializationHeader header); - - public static enum Type - { - //Used internally to refer to files with no - //format flag in the filename - LEGACY("big", BigFormat.instance), - - //The original sstable format - BIG("big", BigFormat.instance); - - public final SSTableFormat info; - public final String name; - - public static Type current() - { - return BIG; - } - - @SuppressWarnings("deprecation") - private Type(String name, SSTableFormat info) - { - //Since format comes right after generation - //we disallow formats with numeric names - // We have removed this check for compatibility with the embedded cassandra used for tests. - assert !CharMatcher.digit().matchesAllOf(name); - - this.name = name; - this.info = info; - } - - public static Type validate(String name) - { - for (Type valid : Type.values()) - { - //This is used internally for old sstables - if (valid == LEGACY) - continue; - - if (valid.name.equalsIgnoreCase(name)) - return valid; - } - - throw new IllegalArgumentException("No Type constant " + name); - } - } -} diff --git a/dao/src/test/java/org/apache/cassandra/io/util/FileUtils.java b/dao/src/test/java/org/apache/cassandra/io/util/FileUtils.java deleted file mode 100644 index 2609a2c3d0..0000000000 --- a/dao/src/test/java/org/apache/cassandra/io/util/FileUtils.java +++ /dev/null @@ -1,760 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you 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.apache.cassandra.io.util; - -import java.io.*; -import java.nio.ByteBuffer; -import java.nio.channels.FileChannel; -import java.nio.charset.Charset; -import java.nio.charset.StandardCharsets; -import java.nio.file.*; -import java.nio.file.attribute.BasicFileAttributes; -import java.nio.file.attribute.FileAttributeView; -import java.nio.file.attribute.FileStoreAttributeView; -import java.text.DecimalFormat; -import java.util.Arrays; -import java.util.Collections; -import java.util.List; -import java.util.Optional; -import java.util.concurrent.atomic.AtomicReference; -import java.util.function.Consumer; -import java.util.function.Predicate; -import java.util.stream.StreamSupport; - -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import org.apache.cassandra.concurrent.ScheduledExecutors; -import org.apache.cassandra.io.FSError; -import org.apache.cassandra.io.FSErrorHandler; -import org.apache.cassandra.io.FSReadError; -import org.apache.cassandra.io.FSWriteError; -import org.apache.cassandra.io.sstable.CorruptSSTableException; -import org.apache.cassandra.utils.JVMStabilityInspector; - -import static com.google.common.base.Throwables.throwIfUnchecked; -import static org.apache.cassandra.utils.Throwables.maybeFail; -import static org.apache.cassandra.utils.Throwables.merge; - -public final class FileUtils -{ - public static final Charset CHARSET = StandardCharsets.UTF_8; - - private static final Logger logger = LoggerFactory.getLogger(FileUtils.class); - public static final long ONE_KB = 1024; - public static final long ONE_MB = 1024 * ONE_KB; - public static final long ONE_GB = 1024 * ONE_MB; - public static final long ONE_TB = 1024 * ONE_GB; - - private static final DecimalFormat df = new DecimalFormat("#.##"); - public static final boolean isCleanerAvailable = false; - private static final AtomicReference> fsErrorHandler = new AtomicReference<>(Optional.empty()); - - public static void createHardLink(String from, String to) - { - createHardLink(new File(from), new File(to)); - } - - public static void createHardLink(File from, File to) - { - if (to.exists()) - throw new RuntimeException("Tried to create duplicate hard link to " + to); - if (!from.exists()) - throw new RuntimeException("Tried to hard link to file that does not exist " + from); - - try - { - Files.createLink(to.toPath(), from.toPath()); - } - catch (IOException e) - { - throw new FSWriteError(e, to); - } - } - - public static File createTempFile(String prefix, String suffix, File directory) - { - try - { - return File.createTempFile(prefix, suffix, directory); - } - catch (IOException e) - { - throw new FSWriteError(e, directory); - } - } - - public static File createTempFile(String prefix, String suffix) - { - return createTempFile(prefix, suffix, new File(System.getProperty("java.io.tmpdir"))); - } - - public static Throwable deleteWithConfirm(String filePath, boolean expect, Throwable accumulate) - { - return deleteWithConfirm(new File(filePath), expect, accumulate); - } - - public static Throwable deleteWithConfirm(File file, boolean expect, Throwable accumulate) - { - boolean exists = file.exists(); - assert exists || !expect : "attempted to delete non-existing file " + file.getName(); - try - { - if (exists) - Files.delete(file.toPath()); - } - catch (Throwable t) - { - try - { - throw new FSWriteError(t, file); - } - catch (Throwable t2) - { - accumulate = merge(accumulate, t2); - } - } - return accumulate; - } - - public static void deleteWithConfirm(String file) - { - deleteWithConfirm(new File(file)); - } - - public static void deleteWithConfirm(File file) - { - maybeFail(deleteWithConfirm(file, true, null)); - } - - public static void renameWithOutConfirm(String from, String to) - { - try - { - atomicMoveWithFallback(new File(from).toPath(), new File(to).toPath()); - } - catch (IOException e) - { - if (logger.isTraceEnabled()) - logger.trace("Could not move file "+from+" to "+to, e); - } - } - - public static void renameWithConfirm(String from, String to) - { - renameWithConfirm(new File(from), new File(to)); - } - - public static void renameWithConfirm(File from, File to) - { - assert from.exists(); - if (logger.isTraceEnabled()) - logger.trace("Renaming {} to {}", from.getPath(), to.getPath()); - // this is not FSWE because usually when we see it it's because we didn't close the file before renaming it, - // and Windows is picky about that. - try - { - atomicMoveWithFallback(from.toPath(), to.toPath()); - } - catch (IOException e) - { - throw new RuntimeException(String.format("Failed to rename %s to %s", from.getPath(), to.getPath()), e); - } - } - - /** - * Move a file atomically, if it fails, it falls back to a non-atomic operation - * @param from - * @param to - * @throws IOException - */ - private static void atomicMoveWithFallback(Path from, Path to) throws IOException - { - try - { - Files.move(from, to, StandardCopyOption.REPLACE_EXISTING, StandardCopyOption.ATOMIC_MOVE); - } - catch (AtomicMoveNotSupportedException e) - { - logger.trace("Could not do an atomic move", e); - Files.move(from, to, StandardCopyOption.REPLACE_EXISTING); - } - - } - public static void truncate(String path, long size) - { - try(FileChannel channel = FileChannel.open(Paths.get(path), StandardOpenOption.READ, StandardOpenOption.WRITE)) - { - channel.truncate(size); - } - catch (IOException e) - { - throw new RuntimeException(e); - } - } - - public static void closeQuietly(Closeable c) - { - try - { - if (c != null) - c.close(); - } - catch (Exception e) - { - logger.warn("Failed closing {}", c, e); - } - } - - public static void closeQuietly(AutoCloseable c) - { - try - { - if (c != null) - c.close(); - } - catch (Exception e) - { - logger.warn("Failed closing {}", c, e); - } - } - - public static void close(Closeable... cs) throws IOException - { - close(Arrays.asList(cs)); - } - - public static void close(Iterable cs) throws IOException - { - Throwable e = null; - for (Closeable c : cs) - { - try - { - if (c != null) - c.close(); - } - catch (Throwable ex) - { - if (e == null) e = ex; - else e.addSuppressed(ex); - logger.warn("Failed closing stream {}", c, ex); - } - } - maybeFail(e, IOException.class); - } - - public static void closeQuietly(Iterable cs) - { - for (AutoCloseable c : cs) - { - try - { - if (c != null) - c.close(); - } - catch (Exception ex) - { - logger.warn("Failed closing {}", c, ex); - } - } - } - - public static String getCanonicalPath(String filename) - { - try - { - return new File(filename).getCanonicalPath(); - } - catch (IOException e) - { - throw new FSReadError(e, filename); - } - } - - public static String getCanonicalPath(File file) - { - try - { - return file.getCanonicalPath(); - } - catch (IOException e) - { - throw new FSReadError(e, file); - } - } - - /** Return true if file is contained in folder */ - public static boolean isContained(File folder, File file) - { - Path folderPath = Paths.get(getCanonicalPath(folder)); - Path filePath = Paths.get(getCanonicalPath(file)); - - return filePath.startsWith(folderPath); - } - - /** Convert absolute path into a path relative to the base path */ - public static String getRelativePath(String basePath, String path) - { - try - { - return Paths.get(basePath).relativize(Paths.get(path)).toString(); - } - catch(Exception ex) - { - String absDataPath = FileUtils.getCanonicalPath(basePath); - return Paths.get(absDataPath).relativize(Paths.get(path)).toString(); - } - } - - public static void clean(ByteBuffer buffer) - { - if (buffer == null) - return; - } - - public static void createDirectory(String directory) - { - createDirectory(new File(directory)); - } - - public static void createDirectory(File directory) - { - if (!directory.exists()) - { - if (!directory.mkdirs()) - throw new FSWriteError(new IOException("Failed to mkdirs " + directory), directory); - } - } - - public static boolean delete(String file) - { - File f = new File(file); - return f.delete(); - } - - public static void delete(File... files) - { - if (files == null) - { - // CASSANDRA-13389: some callers use Files.listFiles() which, on error, silently returns null - logger.debug("Received null list of files to delete"); - return; - } - - for ( File file : files ) - { - file.delete(); - } - } - - public static void deleteAsync(final String file) - { - Runnable runnable = new Runnable() - { - public void run() - { - deleteWithConfirm(new File(file)); - } - }; - ScheduledExecutors.nonPeriodicTasks.execute(runnable); - } - - public static void visitDirectory(Path dir, Predicate filter, Consumer consumer) - { - try (DirectoryStream stream = Files.newDirectoryStream(dir)) - { - StreamSupport.stream(stream.spliterator(), false) - .map(Path::toFile) - // stream directories are weakly consistent so we always check if the file still exists - .filter(f -> f.exists() && (filter == null || filter.test(f))) - .forEach(consumer); - } - catch (IOException|DirectoryIteratorException ex) - { - logger.error("Failed to list files in {} with exception: {}", dir, ex.getMessage(), ex); - } - } - - public static String stringifyFileSize(double value) - { - double d; - if ( value >= ONE_TB ) - { - d = value / ONE_TB; - String val = df.format(d); - return val + " TiB"; - } - else if ( value >= ONE_GB ) - { - d = value / ONE_GB; - String val = df.format(d); - return val + " GiB"; - } - else if ( value >= ONE_MB ) - { - d = value / ONE_MB; - String val = df.format(d); - return val + " MiB"; - } - else if ( value >= ONE_KB ) - { - d = value / ONE_KB; - String val = df.format(d); - return val + " KiB"; - } - else - { - String val = df.format(value); - return val + " bytes"; - } - } - - /** - * Deletes all files and subdirectories under "dir". - * @param dir Directory to be deleted - * @throws FSWriteError if any part of the tree cannot be deleted - */ - public static void deleteRecursive(File dir) - { - if (dir.isDirectory()) - { - String[] children = dir.list(); - for (String child : children) - deleteRecursive(new File(dir, child)); - } - - // The directory is now empty so now it can be smoked - deleteWithConfirm(dir); - } - - /** - * Schedules deletion of all file and subdirectories under "dir" on JVM shutdown. - * @param dir Directory to be deleted - */ - public static void deleteRecursiveOnExit(File dir) - { - if (dir.isDirectory()) - { - String[] children = dir.list(); - for (String child : children) - deleteRecursiveOnExit(new File(dir, child)); - } - - logger.trace("Scheduling deferred deletion of file: {}", dir); - dir.deleteOnExit(); - } - - public static void handleCorruptSSTable(CorruptSSTableException e) - { - fsErrorHandler.get().ifPresent(handler -> handler.handleCorruptSSTable(e)); - } - - public static void handleFSError(FSError e) - { - fsErrorHandler.get().ifPresent(handler -> handler.handleFSError(e)); - } - - /** - * handleFSErrorAndPropagate will invoke the disk failure policy error handler, - * which may or may not stop the daemon or transports. However, if we don't exit, - * we still want to propagate the exception to the caller in case they have custom - * exception handling - * - * @param e A filesystem error - */ - public static void handleFSErrorAndPropagate(FSError e) - { - JVMStabilityInspector.inspectThrowable(e); - throwIfUnchecked(e); - throw new RuntimeException(e); - } - - /** - * Get the size of a directory in bytes - * @param folder The directory for which we need size. - * @return The size of the directory - */ - public static long folderSize(File folder) - { - final long [] sizeArr = {0L}; - try - { - Files.walkFileTree(folder.toPath(), new SimpleFileVisitor() - { - @Override - public FileVisitResult visitFile(Path file, BasicFileAttributes attrs) - { - sizeArr[0] += attrs.size(); - return FileVisitResult.CONTINUE; - } - }); - } - catch (IOException e) - { - logger.error("Error while getting {} folder size. {}", folder, e); - } - return sizeArr[0]; - } - - public static void copyTo(DataInput in, OutputStream out, int length) throws IOException - { - byte[] buffer = new byte[64 * 1024]; - int copiedBytes = 0; - - while (copiedBytes + buffer.length < length) - { - in.readFully(buffer); - out.write(buffer); - copiedBytes += buffer.length; - } - - if (copiedBytes < length) - { - int left = length - copiedBytes; - in.readFully(buffer, 0, left); - out.write(buffer, 0, left); - } - } - - public static boolean isSubDirectory(File parent, File child) throws IOException - { - parent = parent.getCanonicalFile(); - child = child.getCanonicalFile(); - - File toCheck = child; - while (toCheck != null) - { - if (parent.equals(toCheck)) - return true; - toCheck = toCheck.getParentFile(); - } - return false; - } - - public static void append(File file, String ... lines) - { - if (file.exists()) - write(file, Arrays.asList(lines), StandardOpenOption.APPEND); - else - write(file, Arrays.asList(lines), StandardOpenOption.CREATE); - } - - public static void appendAndSync(File file, String ... lines) - { - if (file.exists()) - write(file, Arrays.asList(lines), StandardOpenOption.APPEND, StandardOpenOption.SYNC); - else - write(file, Arrays.asList(lines), StandardOpenOption.CREATE, StandardOpenOption.SYNC); - } - - public static void replace(File file, String ... lines) - { - write(file, Arrays.asList(lines), StandardOpenOption.TRUNCATE_EXISTING); - } - - public static void write(File file, List lines, StandardOpenOption ... options) - { - try - { - Files.write(file.toPath(), - lines, - CHARSET, - options); - } - catch (IOException ex) - { - throw new RuntimeException(ex); - } - } - - public static List readLines(File file) - { - try - { - return Files.readAllLines(file.toPath(), CHARSET); - } - catch (IOException ex) - { - if (ex instanceof NoSuchFileException) - return Collections.emptyList(); - - throw new RuntimeException(ex); - } - } - - public static void setFSErrorHandler(FSErrorHandler handler) - { - fsErrorHandler.getAndSet(Optional.ofNullable(handler)); - } - - /** - * Returns the size of the specified partition. - *

This method handles large file system by returning {@code Long.MAX_VALUE} if the size overflow. - * See JDK-8179320 for more information.

- * - * @param file the partition - * @return the size, in bytes, of the partition or {@code 0L} if the abstract pathname does not name a partition - */ - public static long getTotalSpace(File file) - { - return handleLargeFileSystem(file.getTotalSpace()); - } - - /** - * Returns the number of unallocated bytes on the specified partition. - *

This method handles large file system by returning {@code Long.MAX_VALUE} if the number of unallocated bytes - * overflow. See JDK-8179320 for more information

- * - * @param file the partition - * @return the number of unallocated bytes on the partition or {@code 0L} - * if the abstract pathname does not name a partition. - */ - public static long getFreeSpace(File file) - { - return handleLargeFileSystem(file.getFreeSpace()); - } - - /** - * Returns the number of available bytes on the specified partition. - *

This method handles large file system by returning {@code Long.MAX_VALUE} if the number of available bytes - * overflow. See JDK-8179320 for more information

- * - * @param file the partition - * @return the number of available bytes on the partition or {@code 0L} - * if the abstract pathname does not name a partition. - */ - public static long getUsableSpace(File file) - { - return handleLargeFileSystem(file.getUsableSpace()); - } - - /** - * Returns the {@link FileStore} representing the file store where a file - * is located. This {@link FileStore} handles large file system by returning {@code Long.MAX_VALUE} - * from {@code FileStore#getTotalSpace()}, {@code FileStore#getUnallocatedSpace()} and {@code FileStore#getUsableSpace()} - * it the value is bigger than {@code Long.MAX_VALUE}. See JDK-8162520 - * for more information. - * - * @param path the path to the file - * @return the file store where the file is stored - */ - public static FileStore getFileStore(Path path) throws IOException - { - return new SafeFileStore(Files.getFileStore(path)); - } - - /** - * Handle large file system by returning {@code Long.MAX_VALUE} when the size overflows. - * @param size returned by the Java's FileStore methods - * @return the size or {@code Long.MAX_VALUE} if the size was bigger than {@code Long.MAX_VALUE} - */ - private static long handleLargeFileSystem(long size) - { - return size < 0 ? Long.MAX_VALUE : size; - } - - /** - * Private constructor as the class contains only static methods. - */ - private FileUtils() - { - } - - /** - * FileStore decorator used to safely handle large file system. - * - *

Java's FileStore methods (getTotalSpace/getUnallocatedSpace/getUsableSpace) are limited to reporting bytes as - * signed long (2^63-1), if the filesystem is any bigger, then the size overflows. {@code SafeFileStore} will - * return {@code Long.MAX_VALUE} if the size overflow.

- * - * @see https://bugs.openjdk.java.net/browse/JDK-8162520. - */ - private static final class SafeFileStore extends FileStore - { - /** - * The decorated {@code FileStore} - */ - private final FileStore fileStore; - - public SafeFileStore(FileStore fileStore) - { - this.fileStore = fileStore; - } - - @Override - public String name() - { - return fileStore.name(); - } - - @Override - public String type() - { - return fileStore.type(); - } - - @Override - public boolean isReadOnly() - { - return fileStore.isReadOnly(); - } - - @Override - public long getTotalSpace() throws IOException - { - return handleLargeFileSystem(fileStore.getTotalSpace()); - } - - @Override - public long getUsableSpace() throws IOException - { - return handleLargeFileSystem(fileStore.getUsableSpace()); - } - - @Override - public long getUnallocatedSpace() throws IOException - { - return handleLargeFileSystem(fileStore.getUnallocatedSpace()); - } - - @Override - public boolean supportsFileAttributeView(Class type) - { - return fileStore.supportsFileAttributeView(type); - } - - @Override - public boolean supportsFileAttributeView(String name) - { - return fileStore.supportsFileAttributeView(name); - } - - @Override - public V getFileStoreAttributeView(Class type) - { - return fileStore.getFileStoreAttributeView(type); - } - - @Override - public Object getAttribute(String attribute) throws IOException - { - return fileStore.getAttribute(attribute); - } - } -} diff --git a/dao/src/test/java/org/thingsboard/server/dao/AbstractNoSqlContainer.java b/dao/src/test/java/org/thingsboard/server/dao/AbstractNoSqlContainer.java new file mode 100644 index 0000000000..3f7c2c779a --- /dev/null +++ b/dao/src/test/java/org/thingsboard/server/dao/AbstractNoSqlContainer.java @@ -0,0 +1,96 @@ +/** + * Copyright © 2016-2023 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 com.github.dockerjava.api.command.InspectContainerResponse; +import lombok.extern.slf4j.Slf4j; +import org.apache.commons.io.IOUtils; +import org.junit.ClassRule; +import org.junit.rules.ExternalResource; +import org.testcontainers.containers.CassandraContainer; +import org.testcontainers.containers.delegate.CassandraDatabaseDelegate; +import org.testcontainers.delegate.DatabaseDelegate; +import org.testcontainers.ext.ScriptUtils; + +import javax.script.ScriptException; +import java.io.IOException; +import java.net.URL; +import java.nio.charset.StandardCharsets; +import java.util.List; + +@Slf4j +public abstract class AbstractNoSqlContainer { + + public static final List INIT_SCRIPTS = List.of( + "cassandra/schema-keyspace.cql", + "cassandra/schema-ts.cql", + "cassandra/schema-ts-latest.cql" + ); + + @ClassRule(order = 0) + public static final CassandraContainer cassandra = (CassandraContainer) new CassandraContainer("cassandra:4.1") { + @Override + protected void containerIsStarted(InspectContainerResponse containerInfo) { + super.containerIsStarted(containerInfo); + DatabaseDelegate db = new CassandraDatabaseDelegate(this); + INIT_SCRIPTS.forEach(script -> runInitScriptIfRequired(db, script)); + } + + private void runInitScriptIfRequired(DatabaseDelegate db, String initScriptPath) { + logger().info("Init script [{}]", initScriptPath); + if (initScriptPath != null) { + try { + URL resource = Thread.currentThread().getContextClassLoader().getResource(initScriptPath); + if (resource == null) { + logger().warn("Could not load classpath init script: {}", initScriptPath); + throw new ScriptUtils.ScriptLoadException("Could not load classpath init script: " + initScriptPath + ". Resource not found."); + } + String cql = IOUtils.toString(resource, StandardCharsets.UTF_8); + ScriptUtils.executeDatabaseScript(db, initScriptPath, cql); + } catch (IOException e) { + logger().warn("Could not load classpath init script: {}", initScriptPath); + throw new ScriptUtils.ScriptLoadException("Could not load classpath init script: " + initScriptPath, e); + } catch (ScriptException e) { + logger().error("Error while executing init script: {}", initScriptPath, e); + throw new ScriptUtils.UncategorizedScriptException("Error while executing init script: " + initScriptPath, e); + } + } + } + } + .withEnv("HEAP_NEWSIZE", "64M") + .withEnv("MAX_HEAP_SIZE", "512M") + .withEnv("CASSANDRA_CLUSTER_NAME", "ThingsBoard Cluster"); + + @ClassRule(order = 1) + public static ExternalResource resource = new ExternalResource() { + @Override + protected void before() throws Throwable { + cassandra.start(); + String cassandraUrl = String.format("%s:%s", cassandra.getHost(), cassandra.getMappedPort(9042)); + log.debug("Cassandra url [{}]", cassandraUrl); + System.setProperty("cassandra.url", cassandraUrl); + } + + @Override + protected void after() { + cassandra.stop(); + List.of("cassandra.url") + .forEach(System.getProperties()::remove); + } + }; + +} diff --git a/dao/src/test/java/org/thingsboard/server/dao/AbstractRedisContainer.java b/dao/src/test/java/org/thingsboard/server/dao/AbstractRedisContainer.java new file mode 100644 index 0000000000..e9a6cb8641 --- /dev/null +++ b/dao/src/test/java/org/thingsboard/server/dao/AbstractRedisContainer.java @@ -0,0 +1,51 @@ +/** + * Copyright © 2016-2023 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 lombok.extern.slf4j.Slf4j; +import org.junit.ClassRule; +import org.junit.rules.ExternalResource; +import org.testcontainers.containers.GenericContainer; + +import java.util.List; + +@Slf4j +public class AbstractRedisContainer { + + @ClassRule(order = 0) + public static GenericContainer redis = new GenericContainer("redis:7.0") + .withExposedPorts(6379); + + @ClassRule(order = 1) + public static ExternalResource resource = new ExternalResource() { + @Override + protected void before() throws Throwable { + redis.start(); + System.setProperty("cache.type", "redis"); + System.setProperty("redis.connection.type", "standalone"); + System.setProperty("redis.standalone.host", redis.getHost()); + System.setProperty("redis.standalone.port", String.valueOf(redis.getMappedPort(6379))); + } + + @Override + protected void after() { + redis.stop(); + List.of("cache.type", "redis.connection.type", "redis.standalone.host", "redis.standalone.port") + .forEach(System.getProperties()::remove); + } + }; + +} diff --git a/dao/src/test/java/org/thingsboard/server/dao/CustomCassandraCQLUnit.java b/dao/src/test/java/org/thingsboard/server/dao/CustomCassandraCQLUnit.java deleted file mode 100644 index 7bab5e1138..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/CustomCassandraCQLUnit.java +++ /dev/null @@ -1,88 +0,0 @@ -/** - * Copyright © 2016-2023 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 com.datastax.oss.driver.api.core.CqlSession; -import org.cassandraunit.BaseCassandraUnit; -import org.cassandraunit.CQLDataLoader; -import org.cassandraunit.dataset.CQLDataSet; -import org.cassandraunit.utils.EmbeddedCassandraServerHelper; - -import java.util.List; - -public class CustomCassandraCQLUnit extends BaseCassandraUnit { - protected List dataSets; - - public CqlSession session; - - public CustomCassandraCQLUnit(List dataSets) { - this.dataSets = dataSets; - } - - public CustomCassandraCQLUnit(List dataSets, int readTimeoutMillis) { - this.dataSets = dataSets; - this.readTimeoutMillis = readTimeoutMillis; - } - - public CustomCassandraCQLUnit(List dataSets, String configurationFileName) { - this(dataSets); - this.configurationFileName = configurationFileName; - } - - public CustomCassandraCQLUnit(List dataSets, String configurationFileName, int readTimeoutMillis) { - this(dataSets); - this.configurationFileName = configurationFileName; - this.readTimeoutMillis = readTimeoutMillis; - } - - public CustomCassandraCQLUnit(List dataSets, String configurationFileName, long startUpTimeoutMillis) { - super(startUpTimeoutMillis); - this.dataSets = dataSets; - this.configurationFileName = configurationFileName; - } - - public CustomCassandraCQLUnit(List dataSets, String configurationFileName, long startUpTimeoutMillis, int readTimeoutMillis) { - super(startUpTimeoutMillis); - this.dataSets = dataSets; - this.configurationFileName = configurationFileName; - this.readTimeoutMillis = readTimeoutMillis; - } - - @Override - protected void load() { - session = EmbeddedCassandraServerHelper.getSession(); - CQLDataLoader dataLoader = new CQLDataLoader(session); - dataSets.forEach(dataLoader::load); - session = dataLoader.getSession(); - System.setSecurityManager(null); - } - - @Override - protected void after() { - super.after(); - try (CqlSession s = session) { - session = null; - } - System.setSecurityManager(null); - } - - // Getters for those who do not like to directly access fields - - public CqlSession getSession() { - return session; - } - -} diff --git a/dao/src/test/java/org/thingsboard/server/dao/NoSqlDaoServiceTestSuite.java b/dao/src/test/java/org/thingsboard/server/dao/NoSqlDaoServiceTestSuite.java index 1cc88f5734..9f82b8735a 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/NoSqlDaoServiceTestSuite.java +++ b/dao/src/test/java/org/thingsboard/server/dao/NoSqlDaoServiceTestSuite.java @@ -15,28 +15,14 @@ */ package org.thingsboard.server.dao; -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 java.util.Arrays; - @RunWith(ClasspathSuite.class) @ClassnameFilters({ "org.thingsboard.server.dao.service.*.nosql.*ServiceNoSqlTest", }) -public class NoSqlDaoServiceTestSuite { - - @ClassRule - public static CustomCassandraCQLUnit cassandraUnit = - new CustomCassandraCQLUnit( - Arrays.asList( - new ClassPathCQLDataSet("cassandra/schema-keyspace.cql", false, false), - new ClassPathCQLDataSet("cassandra/schema-ts.cql", false, false), - new ClassPathCQLDataSet("cassandra/schema-ts-latest.cql", false, false) - ), - "cassandra-test.yaml", 30000L); +public class NoSqlDaoServiceTestSuite extends AbstractNoSqlContainer { } diff --git a/dao/src/test/java/org/thingsboard/server/dao/RedisSqlTestSuite.java b/dao/src/test/java/org/thingsboard/server/dao/RedisSqlTestSuite.java index 2fa63f2387..61f77ee1b7 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/RedisSqlTestSuite.java +++ b/dao/src/test/java/org/thingsboard/server/dao/RedisSqlTestSuite.java @@ -15,37 +15,15 @@ */ package org.thingsboard.server.dao; -import org.junit.ClassRule; import org.junit.extensions.cpsuite.ClasspathSuite; import org.junit.extensions.cpsuite.ClasspathSuite.ClassnameFilters; import org.junit.runner.RunWith; -import org.springframework.context.ApplicationContextInitializer; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.test.context.ContextConfiguration; -import org.springframework.test.context.support.TestPropertySourceUtils; -import org.testcontainers.containers.GenericContainer; -@ContextConfiguration(initializers = RedisSqlTestSuite.class) @RunWith(ClasspathSuite.class) @ClassnameFilters( //All the same tests using redis instead of caffeine. "org.thingsboard.server.dao.service.*ServiceSqlTest" ) -public class RedisSqlTestSuite implements ApplicationContextInitializer { - - @ClassRule - public static GenericContainer redis = new GenericContainer("redis:4.0").withExposedPorts(6379); - - @Override - public void initialize(ConfigurableApplicationContext applicationContext) { - TestPropertySourceUtils.addInlinedPropertiesToEnvironment( - applicationContext, "cache.type=redis"); - TestPropertySourceUtils.addInlinedPropertiesToEnvironment( - applicationContext, "redis.connection.type=standalone"); - TestPropertySourceUtils.addInlinedPropertiesToEnvironment( - applicationContext, "redis.standalone.host=localhost"); - TestPropertySourceUtils.addInlinedPropertiesToEnvironment( - applicationContext, "redis.standalone.port=" + redis.getMappedPort(6379)); - } +public class RedisSqlTestSuite extends AbstractRedisContainer { } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseDeviceCredentialsCacheTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseDeviceCredentialsCacheTest.java index 9415da99ea..7fc9f7b218 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseDeviceCredentialsCacheTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseDeviceCredentialsCacheTest.java @@ -43,8 +43,8 @@ import static org.mockito.Mockito.when; public abstract class BaseDeviceCredentialsCacheTest extends AbstractServiceTest { - private static final String CREDENTIALS_ID_1 = StringUtils.randomAlphanumeric(20); - private static final String CREDENTIALS_ID_2 = StringUtils.randomAlphanumeric(20); + private final String CREDENTIALS_ID_1 = StringUtils.randomAlphanumeric(20); + private final String CREDENTIALS_ID_2 = StringUtils.randomAlphanumeric(20); @Autowired private DeviceCredentialsService deviceCredentialsService; diff --git a/dao/src/test/resources/application-test.properties b/dao/src/test/resources/application-test.properties index 15131011d8..a2a21abc11 100644 --- a/dao/src/test/resources/application-test.properties +++ b/dao/src/test/resources/application-test.properties @@ -7,10 +7,9 @@ updates.enabled=false audit-log.enabled=true audit-log.sink.type=none -cache.type=caffeine +#cache.type=caffeine # will be injected redis by RedisContainer or will be default (caffeine) cache.maximumPoolSize=16 cache.attributes.enabled=true -#cache.type=redis cache.specs.relations.timeToLiveInMinutes=1440 cache.specs.relations.maxSize=100000 diff --git a/dao/src/test/resources/cassandra-test.properties b/dao/src/test/resources/cassandra-test.properties index 5765153a6f..5876582ee6 100644 --- a/dao/src/test/resources/cassandra-test.properties +++ b/dao/src/test/resources/cassandra-test.properties @@ -2,7 +2,7 @@ cassandra.cluster_name=Thingsboard Cluster cassandra.keyspace_name=thingsboard -cassandra.url=127.0.0.1:9142 +#cassandra.url=127.0.0.1:9142 # will be injected by NoSqlContainer cassandra.local_datacenter=datacenter1 diff --git a/dao/src/test/resources/logback.xml b/dao/src/test/resources/logback.xml index f74ba574b5..61397ec6f1 100644 --- a/dao/src/test/resources/logback.xml +++ b/dao/src/test/resources/logback.xml @@ -8,9 +8,7 @@ - - - + diff --git a/msa/black-box-tests/README.md b/msa/black-box-tests/README.md index ae0e8b79a0..5b63ba712a 100644 --- a/msa/black-box-tests/README.md +++ b/msa/black-box-tests/README.md @@ -18,7 +18,7 @@ As result, in REPOSITORY column, next images should be present: thingsboard/tb-web-ui thingsboard/tb-js-executor -- Run the black box tests in the [msa/black-box-tests](../black-box-tests) directory with Redis standalone: +- Run the black box tests (without ui tests) in the [msa/black-box-tests](../black-box-tests) directory with Redis standalone: mvn clean install -DblackBoxTests.skip=false @@ -34,11 +34,23 @@ As result, in REPOSITORY column, next images should be present: mvn clean install -DblackBoxTests.skip=false -DrunLocal=true -- To run ui smoke tests in the [msa/black-box-tests](../black-box-tests) directory specifying suite name: +- To run only ui tests in the [msa/black-box-tests](../black-box-tests) directory: mvn clean install -DblackBoxTests.skip=false -Dsuite=uiTests -- To run all tests in the [msa/black-box-tests](../black-box-tests) directory specifying suite name: +- To run only ui smoke rule chains tests in the [msa/black-box-tests](../black-box-tests) directory: + + mvn clean install -DblackBoxTests.skip=false -Dsuite=smokesRuleChain + +- To run only ui smoke customers tests in the [msa/black-box-tests](../black-box-tests) directory: + + mvn clean install -DblackBoxTests.skip=false -Dsuite=smokesCustomer + +- To run only ui smoke profiles tests in the [msa/black-box-tests](../black-box-tests) directory: + + mvn clean install -DblackBoxTests.skip=false -Dsuite=smokesPrifiles + +- To run all tests (black-box and ui) in the [msa/black-box-tests](../black-box-tests) directory: mvn clean install -DblackBoxTests.skip=false -Dsuite=all diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ui/pages/OtherPageElements.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ui/pages/OtherPageElements.java index 3b18df801f..99e737dba1 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ui/pages/OtherPageElements.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ui/pages/OtherPageElements.java @@ -42,7 +42,7 @@ public class OtherPageElements extends AbstractBasePage { private static final String MARKS_CHECKBOX = "//mat-row[contains (@class,'mat-selected')]//mat-checkbox[contains(@class, 'checked')]"; private static final String SELECT_ALL_CHECKBOX = "//thead//mat-checkbox"; private static final String ALL_ENTITY = "//mat-row[@class='mat-row cdk-row mat-row-select ng-star-inserted']"; - private static final String EDIT_PENCIL_BTN = "//mat-icon[contains(text(),'edit')]/ancestor::button"; + private static final String EDIT_PENCIL_BTN = "//tb-details-panel//mat-icon[contains(text(),'edit')]/ancestor::button"; private static final String NAME_FIELD_EDIT_VIEW = "//input[@formcontrolname='name']"; private static final String HEADER_NAME_VIEW = "//header//div[@class='tb-details-title']/span"; private static final String DONE_BTN_EDIT_VIEW = "//mat-icon[contains(text(),'done')]/ancestor::button"; diff --git a/pom.xml b/pom.xml index 7d6b59f10b..77aeb81498 100755 --- a/pom.xml +++ b/pom.xml @@ -126,7 +126,6 @@ 2.8.5 4.1.0 - 4.3.1.0 2.7.2 1.5.2 5.8.2 @@ -659,7 +658,7 @@ ${surefire.version} - --illegal-access=permit -Xss384k -XX:+UseStringDeduplication -XX:MaxGCPauseMillis=20 + --illegal-access=permit -XX:+UseStringDeduplication -XX:MaxGCPauseMillis=20 @@ -1603,26 +1602,6 @@
- - org.cassandraunit - cassandra-unit - ${cassandra-unit.version} - test - - - junit - junit - - - org.hamcrest - hamcrest-core - - - org.hamcrest - hamcrest-library - - - org.apache.cassandra cassandra-all @@ -1746,6 +1725,12 @@ bcpkix-jdk15on ${bouncycastle.version} + + org.testcontainers + cassandra + ${testcontainers.version} + test + org.testcontainers postgresql diff --git a/rule-engine/rule-engine-components/pom.xml b/rule-engine/rule-engine-components/pom.xml index 56621162f7..e85f0094b5 100644 --- a/rule-engine/rule-engine-components/pom.xml +++ b/rule-engine/rule-engine-components/pom.xml @@ -150,22 +150,6 @@ mockserver-client-java test - - - org.cassandraunit - cassandra-unit - - - org.slf4j - slf4j-log4j12 - - - org.hibernate - hibernate-validator - - - test - com.jayway.jsonpath json-path