diff --git a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java index b0f5f41c9a..856b17fc61 100644 --- a/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java +++ b/application/src/main/java/org/thingsboard/server/controller/TelemetryController.java @@ -284,7 +284,7 @@ public class TelemetryController extends BaseController { + MARKDOWN_CODE_BLOCK_END + "\n\n" + INVALID_ENTITY_ID_OR_ENTITY_TYPE_DESCRIPTION + TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH) @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')") - @GetMapping(value = "/{entityType}/{entityId}/values/timeseries", params = {"keys", "startTs", "endTs"}) + @GetMapping(value = "/{entityType}/{entityId}/values/timeseries", params = {"startTs", "endTs"}) public DeferredResult getTimeseries( @Parameter(description = ENTITY_TYPE_PARAM_DESCRIPTION, required = true, schema = @Schema(defaultValue = "DEVICE")) @PathVariable("entityType") String entityType, @Parameter(description = ENTITY_ID_PARAM_DESCRIPTION, required = true) @PathVariable("entityId") String entityIdStr, diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 67ad5be16b..e2bbea157e 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -341,6 +341,8 @@ cassandra: set_null_values_enabled: "${CASSANDRA_QUERY_SET_NULL_VALUES_ENABLED:true}" # log one of cassandra queries with specified frequency (0 - logging is disabled) print_queries_freq: "${CASSANDRA_QUERY_PRINT_FREQ:0}" + # Maximum total size in bytes of a Cassandra query result set across all pages. Default is 50MB. 0 means unlimited + max_result_set_size_in_bytes: "${CASSANDRA_QUERY_MAX_RESULT_SET_SIZE_IN_BYTES:52428800}" tenant_rate_limits: # Whether to print rate-limited tenant names when printing Cassandra query queue statistic print_tenant_names: "${CASSANDRA_QUERY_TENANT_RATE_LIMITS_PRINT_TENANT_NAMES:false}" diff --git a/application/src/test/java/org/thingsboard/server/controller/TelemetryControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/TelemetryControllerTest.java index dd51cdd440..ce6d291546 100644 --- a/application/src/test/java/org/thingsboard/server/controller/TelemetryControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/TelemetryControllerTest.java @@ -102,24 +102,29 @@ public class TelemetryControllerTest extends AbstractControllerTest { Assert.assertEquals(11L, thirdIntervalResult.get("value").asLong()); Assert.assertEquals(thirdIntervalTs, thirdIntervalResult.get("ts").asLong()); - result = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getId() + + ObjectNode resultByMonth = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getId() + "/values/timeseries?keys=t&startTs={startTs}&endTs={endTs}&agg={agg}&intervalType={intervalType}&timeZone={timeZone}", ObjectNode.class, startTs, endTs, "SUM", "MONTH", "Europe/Kyiv"); - Assert.assertNotNull(result); - Assert.assertNotNull(result.get("t")); - Assert.assertEquals(1, result.get("t").size()); + Assert.assertNotNull(resultByMonth); + Assert.assertNotNull(resultByMonth.get("t")); + Assert.assertEquals(1, resultByMonth.get("t").size()); - var monthResult = result.get("t").get(0); + var monthResult = resultByMonth.get("t").get(0); Assert.assertEquals(22L, monthResult.get("value").asLong()); Assert.assertEquals(middleOfTheInterval, monthResult.get("ts").asLong()); - // get all latest (without keys parameter) - ObjectNode allLatest = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getId() + + // check timeseries history with key parameter (instead of keys) has the same result as with keys parameter + ObjectNode resultByKey = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getId() + + "/values/timeseries?key=t&startTs={startTs}&endTs={endTs}&agg={agg}&intervalType={intervalType}&timeZone={timeZone}", + ObjectNode.class, startTs, endTs, "SUM", "WEEK_ISO", "Europe/Kyiv"); + Assert.assertEquals(result, resultByKey); + + // check timeseries history without keys and key results into empty object + ObjectNode resultWithoutKeyAndKeys = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getId() + "/values/timeseries?startTs={startTs}&endTs={endTs}&agg={agg}&intervalType={intervalType}&timeZone={timeZone}", ObjectNode.class, startTs, endTs, "SUM", "WEEK_ISO", "Europe/Kyiv"); - Assert.assertNotNull(allLatest); - Assert.assertNotNull(allLatest.get("t")); + Assert.assertTrue(resultWithoutKeyAndKeys.isEmpty()); } @Test diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/nosql/ResultSetSizeLimitExceededException.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/nosql/ResultSetSizeLimitExceededException.java new file mode 100644 index 0000000000..f20a8ebde0 --- /dev/null +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/nosql/ResultSetSizeLimitExceededException.java @@ -0,0 +1,32 @@ +/** + * Copyright © 2016-2026 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.nosql; + +import lombok.Getter; + +@Getter +public class ResultSetSizeLimitExceededException extends IllegalArgumentException { + + private final long limitBytes; + private final long actualBytes; + + public ResultSetSizeLimitExceededException(long limitBytes, long actualBytes) { + super("Result set size exceeds the maximum allowed limit. Please narrow your query"); + this.limitBytes = limitBytes; + this.actualBytes = actualBytes; + } + +} diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/nosql/TbResultSet.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/nosql/TbResultSet.java index 9dd5c6f4b1..a7a4aa5fa8 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/nosql/TbResultSet.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/nosql/TbResultSet.java @@ -34,6 +34,7 @@ import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletionStage; import java.util.concurrent.Executor; +import java.util.concurrent.atomic.AtomicLong; import java.util.function.Function; public class TbResultSet implements AsyncResultSet { @@ -89,9 +90,14 @@ public class TbResultSet implements AsyncResultSet { } public ListenableFuture> allRows(Executor executor) { + return allRows(executor, 0); + } + + public ListenableFuture> allRows(Executor executor, long maxResultSetSizeBytes) { List allRows = new ArrayList<>(); SettableFuture> resultFuture = SettableFuture.create(); - this.processRows(originalStatement, delegate, allRows, resultFuture, executor); + AtomicLong accumulatedBytes = new AtomicLong(0); + this.processRows(originalStatement, delegate, allRows, resultFuture, executor, maxResultSetSizeBytes, accumulatedBytes); return resultFuture; } @@ -99,7 +105,19 @@ public class TbResultSet implements AsyncResultSet { AsyncResultSet resultSet, List allRows, SettableFuture> resultFuture, - Executor executor) { + Executor executor, + long maxResultSetSizeBytes, + AtomicLong accumulatedBytes) { + if (maxResultSetSizeBytes > 0) { + int pageSizeInBytes = resultSet.getExecutionInfo().getResponseSizeInBytes(); + if (pageSizeInBytes > 0) { + accumulatedBytes.addAndGet(pageSizeInBytes); + } + if (accumulatedBytes.get() > maxResultSetSizeBytes) { + resultFuture.setException(new ResultSetSizeLimitExceededException(maxResultSetSizeBytes, accumulatedBytes.get())); + return; + } + } allRows.addAll(loadRows(resultSet)); if (resultSet.hasMorePages()) { ByteBuffer nextPagingState = resultSet.getExecutionInfo().getPagingState(); @@ -110,7 +128,7 @@ public class TbResultSet implements AsyncResultSet { @Override public void onSuccess(@Nullable TbResultSet result) { processRows(nextStatement, result, - allRows, resultFuture, executor); + allRows, resultFuture, executor, maxResultSetSizeBytes, accumulatedBytes); } @Override diff --git a/common/dao-api/src/test/java/org/thingsboard/server/dao/nosql/TbResultSetTest.java b/common/dao-api/src/test/java/org/thingsboard/server/dao/nosql/TbResultSetTest.java new file mode 100644 index 0000000000..01c1c2c506 --- /dev/null +++ b/common/dao-api/src/test/java/org/thingsboard/server/dao/nosql/TbResultSetTest.java @@ -0,0 +1,140 @@ +/** + * Copyright © 2016-2026 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.nosql; + +import com.datastax.oss.driver.api.core.cql.AsyncResultSet; +import com.datastax.oss.driver.api.core.cql.ColumnDefinitions; +import com.datastax.oss.driver.api.core.cql.ExecutionInfo; +import com.datastax.oss.driver.api.core.cql.Row; +import com.datastax.oss.driver.api.core.cql.Statement; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; +import com.google.common.util.concurrent.SettableFuture; +import org.junit.jupiter.api.Test; + +import java.nio.ByteBuffer; +import java.util.List; +import java.util.concurrent.ExecutionException; +import java.util.function.Function; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +class TbResultSetTest { + + @Test + void allRows_withinLimit_returnsAllRows() throws Exception { + Row row = mock(Row.class); + AsyncResultSet asyncResultSet = createMockResultSet(List.of(row), false, 1000); + Statement statement = mock(Statement.class); + + TbResultSet tbResultSet = new TbResultSet(statement, asyncResultSet, s -> null); + ListenableFuture> future = tbResultSet.allRows(MoreExecutors.directExecutor(), 5000); + + List result = future.get(); + assertThat(result).hasSize(1); + assertThat(result.get(0)).isSameAs(row); + } + + @Test + void allRows_exceedsLimitOnFirstPage_failsWithException() { + Row row = mock(Row.class); + AsyncResultSet asyncResultSet = createMockResultSet(List.of(row), false, 6000); + Statement statement = mock(Statement.class); + + TbResultSet tbResultSet = new TbResultSet(statement, asyncResultSet, s -> null); + ListenableFuture> future = tbResultSet.allRows(MoreExecutors.directExecutor(), 5000); + + assertThatThrownBy(future::get) + .isInstanceOf(ExecutionException.class) + .hasCauseInstanceOf(ResultSetSizeLimitExceededException.class); + } + + @Test + void allRows_exceedsLimitOnSecondPage_failsAfterSecondPage() { + Row row1 = mock(Row.class); + Row row2 = mock(Row.class); + Statement statement = mock(Statement.class); + doReturn(statement).when(statement).setPagingState((ByteBuffer) null); + + AsyncResultSet page2 = createMockResultSet(List.of(row2), false, 3000); + TbResultSet tbResultSetPage2 = new TbResultSet(statement, page2, s -> null); + SettableFuture page2Future = SettableFuture.create(); + page2Future.set(tbResultSetPage2); + TbResultSetFuture tbPage2Future = new TbResultSetFuture(page2Future); + + ExecutionInfo page1ExecInfo = mock(ExecutionInfo.class); + when(page1ExecInfo.getResponseSizeInBytes()).thenReturn(3000); + when(page1ExecInfo.getPagingState()).thenReturn(null); + + AsyncResultSet page1 = createMockResultSet(List.of(row1), true, 3000); + when(page1.getExecutionInfo()).thenReturn(page1ExecInfo); + + Function executeAsync = s -> tbPage2Future; + TbResultSet tbResultSet = new TbResultSet(statement, page1, executeAsync); + ListenableFuture> future = tbResultSet.allRows(MoreExecutors.directExecutor(), 5000); + + assertThatThrownBy(future::get) + .isInstanceOf(ExecutionException.class) + .hasCauseInstanceOf(ResultSetSizeLimitExceededException.class); + } + + @Test + void allRows_unlimitedWithZero_returnsAllRowsRegardlessOfSize() throws Exception { + Row row = mock(Row.class); + AsyncResultSet asyncResultSet = createMockResultSet(List.of(row), false, 999999); + Statement statement = mock(Statement.class); + + TbResultSet tbResultSet = new TbResultSet(statement, asyncResultSet, s -> null); + ListenableFuture> future = tbResultSet.allRows(MoreExecutors.directExecutor(), 0); + + List result = future.get(); + assertThat(result).hasSize(1); + } + + @Test + void allRows_noLimitOverload_returnsAllRows() throws Exception { + Row row = mock(Row.class); + AsyncResultSet asyncResultSet = createMockResultSet(List.of(row), false, 999999); + Statement statement = mock(Statement.class); + + TbResultSet tbResultSet = new TbResultSet(statement, asyncResultSet, s -> null); + ListenableFuture> future = tbResultSet.allRows(MoreExecutors.directExecutor()); + + List result = future.get(); + assertThat(result).hasSize(1); + } + + private AsyncResultSet createMockResultSet(List rows, boolean hasMorePages, int responseSizeInBytes) { + AsyncResultSet resultSet = mock(AsyncResultSet.class); + ExecutionInfo executionInfo = mock(ExecutionInfo.class); + ColumnDefinitions columnDefs = mock(ColumnDefinitions.class); + + when(executionInfo.getResponseSizeInBytes()).thenReturn(responseSizeInBytes); + when(executionInfo.getPagingState()).thenReturn(null); + when(resultSet.getExecutionInfo()).thenReturn(executionInfo); + when(resultSet.getColumnDefinitions()).thenReturn(columnDefs); + when(resultSet.currentPage()).thenReturn(rows); + when(resultSet.hasMorePages()).thenReturn(hasMorePages); + when(resultSet.remaining()).thenReturn(rows.size()); + + return resultSet; + } + +} diff --git a/common/version-control/src/main/java/org/thingsboard/server/service/sync/DefaultGitSyncService.java b/common/version-control/src/main/java/org/thingsboard/server/service/sync/DefaultGitSyncService.java index 2abf025cda..3f405e145b 100644 --- a/common/version-control/src/main/java/org/thingsboard/server/service/sync/DefaultGitSyncService.java +++ b/common/version-control/src/main/java/org/thingsboard/server/service/sync/DefaultGitSyncService.java @@ -18,6 +18,7 @@ package org.thingsboard.server.service.sync; import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; +import org.eclipse.jgit.revwalk.RevCommit; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardExecutors; @@ -47,6 +48,8 @@ public class DefaultGitSyncService implements GitSyncService { private final Map repositories = new ConcurrentHashMap<>(); private final Map updateListeners = new ConcurrentHashMap<>(); + private RevCommit lastCommit; + @Override public void registerSync(String key, String repoUri, String branch, long fetchFrequencyMs, Runnable onUpdate) { RepositorySettings settings = new RepositorySettings(); @@ -84,7 +87,7 @@ public class DefaultGitSyncService implements GitSyncService { @Override public List listFiles(String key, String path, int depth, FileType type) { GitRepository repository = getRepository(key); - return repository.listFilesAtCommit(getBranchRef(repository), path, depth).stream() + return repository.listFilesAtCommit(lastCommit, path, depth).stream() .filter(file -> type == null || file.type() == type) .toList(); } @@ -93,7 +96,7 @@ public class DefaultGitSyncService implements GitSyncService { @Override public byte[] getFileContent(String key, String path) { GitRepository repository = getRepository(key); - return repository.getFileContentAtCommit(path, getBranchRef(repository)); + return repository.getFileContentAtCommit(path, lastCommit); } @Override @@ -137,6 +140,15 @@ public class DefaultGitSyncService implements GitSyncService { } private void onUpdate(String key) { + GitRepository repository = getRepository(key); + String branchRef = getBranchRef(repository); + try { + lastCommit = repository.resolveCommit(branchRef); + } catch (Throwable e) { + log.error("[{}] Failed to resolve commit for ref {}", key, branchRef, e); + return; + } + Runnable listener = updateListeners.get(key); if (listener != null) { log.debug("[{}] Handling repository update", key); diff --git a/common/version-control/src/main/java/org/thingsboard/server/service/sync/vc/GitRepository.java b/common/version-control/src/main/java/org/thingsboard/server/service/sync/vc/GitRepository.java index 434d52a4ca..3602dd55ca 100644 --- a/common/version-control/src/main/java/org/thingsboard/server/service/sync/vc/GitRepository.java +++ b/common/version-control/src/main/java/org/thingsboard/server/service/sync/vc/GitRepository.java @@ -277,13 +277,17 @@ public class GitRepository { return listFilesAtCommit(commitId, path, -1).stream().map(RepoFile::path).toList(); } - @SneakyThrows public List listFilesAtCommit(String commitId, String path, int depth) { - log.debug("Executing listFilesAtCommit [{}][{}][{}]", settings.getRepositoryUri(), commitId, path); + RevCommit commit = resolveCommit(commitId); + return listFilesAtCommit(commit, path, depth); + } + + @SneakyThrows + public List listFilesAtCommit(RevCommit commit, String path, int depth) { + log.debug("Executing listFilesAtCommit [{}][{}][{}]", settings.getRepositoryUri(), commit, path); List files = new ArrayList<>(); - RevCommit revCommit = resolveCommit(commitId); try (TreeWalk treeWalk = new TreeWalk(git.getRepository())) { - treeWalk.reset(revCommit.getTree().getId()); + treeWalk.reset(commit.getTree().getId()); if (StringUtils.isNotEmpty(path)) { treeWalk.setFilter(PathFilter.create(path)); } @@ -301,11 +305,15 @@ public class GitRepository { return files; } - @SneakyThrows public byte[] getFileContentAtCommit(String file, String commitId) { - log.debug("Executing getFileContentAtCommit [{}][{}][{}]", settings.getRepositoryUri(), commitId, file); - RevCommit revCommit = resolveCommit(commitId); - try (TreeWalk treeWalk = TreeWalk.forPath(git.getRepository(), file, revCommit.getTree())) { + RevCommit commit = resolveCommit(commitId); + return getFileContentAtCommit(file, commit); + } + + @SneakyThrows + public byte[] getFileContentAtCommit(String file, RevCommit commit) { + log.debug("Executing getFileContentAtCommit [{}][{}][{}]", settings.getRepositoryUri(), commit, file); + try (TreeWalk treeWalk = TreeWalk.forPath(git.getRepository(), file, commit.getTree())) { if (treeWalk == null) { throw new IllegalArgumentException("File not found"); } @@ -373,9 +381,9 @@ public class GitRepository { for (RemoteRefUpdate update : pushResult.getRemoteUpdates()) { RemoteRefUpdate.Status status = update.getStatus(); if (status == REJECTED_NONFASTFORWARD || status == REJECTED_NODELETE || - status == REJECTED_REMOTE_CHANGED || status == REJECTED_OTHER_REASON) { + status == REJECTED_REMOTE_CHANGED || status == REJECTED_OTHER_REASON) { throw new RuntimeException("Remote repository answered with error: " + - Optional.ofNullable(update.getMessage()).orElseGet(status::name)); + Optional.ofNullable(update.getMessage()).orElseGet(status::name)); } } }); @@ -450,7 +458,8 @@ public class GitRepository { revCommit.getFullMessage(), revCommit.getAuthorIdent().getName(), revCommit.getAuthorIdent().getEmailAddress()); } - private RevCommit resolveCommit(String id) throws IOException { + @SneakyThrows + public RevCommit resolveCommit(String id) { return git.getRepository().parseCommit(resolve(id)); } @@ -481,8 +490,8 @@ public class GitRepository { private static final Function> revCommitComparatorFunction = pageLink -> { SortOrder sortOrder = pageLink.getSortOrder(); if (sortOrder != null - && sortOrder.getProperty().equals("timestamp") - && SortOrder.Direction.ASC.equals(sortOrder.getDirection())) { + && sortOrder.getProperty().equals("timestamp") + && SortOrder.Direction.ASC.equals(sortOrder.getDirection())) { return Comparator.comparingInt(RevCommit::getCommitTime); } return null; diff --git a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraAbstractAsyncDao.java b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraAbstractAsyncDao.java index 900958dcba..3d2ebfd1fb 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraAbstractAsyncDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/nosql/CassandraAbstractAsyncDao.java @@ -37,6 +37,9 @@ public abstract class CassandraAbstractAsyncDao extends CassandraAbstractDao { @Value("${cassandra.query.result_processing_threads:50}") private int threadPoolSize; + @Value("${cassandra.query.max_result_set_size_in_bytes:52428800}") + protected long maxResultSetSizeBytes; + @PostConstruct public void startExecutor() { readResultsProcessingExecutor = ThingsBoardExecutors.newWorkStealingPool(threadPoolSize, "cassandra-callback"); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index e68e315b4b..a421cd6842 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -52,6 +52,7 @@ import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntryAggWrapper; import org.thingsboard.server.common.data.kv.TsKvQuery; import org.thingsboard.server.dao.model.ModelConstants; +import org.thingsboard.server.dao.nosql.ResultSetSizeLimitExceededException; import org.thingsboard.server.dao.nosql.TbResultSet; import org.thingsboard.server.dao.nosql.TbResultSetFuture; import org.thingsboard.server.dao.sqlts.AggregationTimeseriesDao; @@ -256,7 +257,8 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD @Override public void onFailure(Throwable t) { - log.error("[{}][{}] Failed to fetch partitions for interval {}-{}", entityId.getEntityType().name(), entityId.getId(), minPartition, maxPartition, t); + log.error("[{}][{}][{}] Failed to fetch partitions for interval {}-{}", tenantId, entityId.getEntityType(), entityId.getId(), minPartition, maxPartition, t); + resultFuture.setException(t); } }, readResultsProcessingExecutor); return resultFuture; @@ -330,7 +332,8 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD @Override public void onFailure(Throwable t) { - log.error("[{}][{}] Failed to fetch partitions for interval {}-{}", entityId.getEntityType().name(), entityId.getId(), toPartitionTs(query.getStartTs()), toPartitionTs(query.getEndTs()), t); + log.error("[{}][{}][{}] Failed to fetch partitions for interval {}-{}", tenantId, entityId.getEntityType(), entityId.getId(), toPartitionTs(query.getStartTs()), toPartitionTs(query.getEndTs()), t); + resultFuture.setException(t); } }, readResultsProcessingExecutor); @@ -372,8 +375,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD cursor.addData(convertResultToTsKvEntryList(Collections.emptyList())); findAllAsyncSequentiallyWithLimit(tenantId, cursor, resultFuture); } else { - Futures.addCallback(result.allRows(readResultsProcessingExecutor), new FutureCallback>() { - + Futures.addCallback(result.allRows(readResultsProcessingExecutor, maxResultSetSizeBytes), new FutureCallback>() { @Override public void onSuccess(@Nullable List result) { cursor.addData(convertResultToTsKvEntryList(result == null ? Collections.emptyList() : result)); @@ -382,7 +384,13 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD @Override public void onFailure(Throwable t) { - log.error("[{}][{}] Failed to fetch data for query {}-{}", stmt, t); + if (t instanceof ResultSetSizeLimitExceededException e) { + log.warn("[{}][{}][{}] Result set size limit exceeded for key [{}], query [{}]: {} bytes, limit {} bytes", + tenantId, cursor.getEntityType(), cursor.getEntityId(), cursor.getKey(), stmt.getPreparedStatement().getQuery(), e.getActualBytes(), e.getLimitBytes()); + } else { + log.error("[{}][{}][{}] Failed to fetch data for key [{}], query [{}]", tenantId, cursor.getEntityType(), cursor.getEntityId(), cursor.getKey(), stmt.getPreparedStatement().getQuery(), t); + } + resultFuture.setException(t); } }, readResultsProcessingExecutor); @@ -392,7 +400,8 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD @Override public void onFailure(Throwable t) { - log.error("[{}][{}] Failed to fetch data for query {}-{}", stmt, t); + log.error("[{}][{}][{}] Failed to fetch data for key [{}], query [{}]", tenantId, cursor.getEntityType(), cursor.getEntityId(), cursor.getKey(), stmt.getPreparedStatement().getQuery(), t); + resultFuture.setException(t); } }, readResultsProcessingExecutor); } @@ -425,7 +434,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD } if (!isUseTsKeyValuePartitioningOnRead()) { final long estimatedPartitionCount = estimatePartitionCount(minPartition, maxPartition); - if (estimatedPartitionCount <= useTsKeyValuePartitioningOnReadMaxEstimatedPartitionCount) { + if (estimatedPartitionCount <= useTsKeyValuePartitioningOnReadMaxEstimatedPartitionCount) { return Futures.immediateFuture(calculatePartitions(minPartition, maxPartition, (int) estimatedPartitionCount)); } } @@ -446,7 +455,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD } List calculatePartitions(long minPartition, long maxPartition) { - return calculatePartitions(minPartition, maxPartition, 0); + return calculatePartitions(minPartition, maxPartition, 0); } List calculatePartitions(long minPartition, long maxPartition, int estimatedPartitionCount) { @@ -533,6 +542,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD public void onFailure(Throwable t) { } + } private long computeTtl(long ttl) { @@ -569,7 +579,8 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD @Override public void onFailure(Throwable t) { - log.error("[{}][{}] Failed to delete data for query {}-{}", stmt, t); + log.error("[{}][{}][{}] Failed to delete data for key [{}], query [{}]", tenantId, cursor.getEntityType(), cursor.getEntityId(), cursor.getKey(), stmt.getPreparedStatement().getQuery(), t); + resultFuture.setException(t); } }, readResultsProcessingExecutor); } @@ -581,12 +592,12 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD try { if (deleteStmt == null) { deleteStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_CF + - " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.TS_COLUMN + " >= ? " - + "AND " + ModelConstants.TS_COLUMN + " < ?"); + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.TS_COLUMN + " >= ? " + + "AND " + ModelConstants.TS_COLUMN + " < ?"); } } finally { stmtCreationLock.unlock(); @@ -661,13 +672,13 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD private String getPreparedStatementQuery(DataType type) { return INSERT_INTO + ModelConstants.TS_KV_CF + - "(" + ModelConstants.ENTITY_TYPE_COLUMN + - "," + ModelConstants.ENTITY_ID_COLUMN + - "," + ModelConstants.KEY_COLUMN + - "," + ModelConstants.PARTITION_COLUMN + - "," + ModelConstants.TS_COLUMN + - "," + getColumnName(type) + ")" + - " VALUES(?, ?, ?, ?, ?, ?)"; + "(" + ModelConstants.ENTITY_TYPE_COLUMN + + "," + ModelConstants.ENTITY_ID_COLUMN + + "," + ModelConstants.KEY_COLUMN + + "," + ModelConstants.PARTITION_COLUMN + + "," + ModelConstants.TS_COLUMN + + "," + getColumnName(type) + ")" + + " VALUES(?, ?, ?, ?, ?, ?)"; } private String getPreparedStatementQueryWithTtl(DataType type) { @@ -680,11 +691,11 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD try { if (partitionInsertStmt == null) { partitionInsertStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF + - "(" + ModelConstants.ENTITY_TYPE_COLUMN + - "," + ModelConstants.ENTITY_ID_COLUMN + - "," + ModelConstants.PARTITION_COLUMN + - "," + ModelConstants.KEY_COLUMN + ")" + - " VALUES(?, ?, ?, ?)"); + "(" + ModelConstants.ENTITY_TYPE_COLUMN + + "," + ModelConstants.ENTITY_ID_COLUMN + + "," + ModelConstants.PARTITION_COLUMN + + "," + ModelConstants.KEY_COLUMN + ")" + + " VALUES(?, ?, ?, ?)"); } } finally { stmtCreationLock.unlock(); @@ -699,11 +710,11 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD try { if (partitionInsertTtlStmt == null) { partitionInsertTtlStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF + - "(" + ModelConstants.ENTITY_TYPE_COLUMN + - "," + ModelConstants.ENTITY_ID_COLUMN + - "," + ModelConstants.PARTITION_COLUMN + - "," + ModelConstants.KEY_COLUMN + ")" + - " VALUES(?, ?, ?, ?) USING TTL ?"); + "(" + ModelConstants.ENTITY_TYPE_COLUMN + + "," + ModelConstants.ENTITY_ID_COLUMN + + "," + ModelConstants.PARTITION_COLUMN + + "," + ModelConstants.KEY_COLUMN + ")" + + " VALUES(?, ?, ?, ?) USING TTL ?"); } } finally { stmtCreationLock.unlock(); @@ -809,16 +820,17 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD fetchStmts[type.ordinal()] = fetchStmts[Aggregation.SUM.ordinal()]; } else { fetchStmts[type.ordinal()] = prepare(SELECT_PREFIX + - String.join(", ", ModelConstants.getFetchColumnNames(type)) + " FROM " + ModelConstants.TS_KV_CF - + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.TS_COLUMN + " >= ? " - + "AND " + ModelConstants.TS_COLUMN + " < ?" - + (type == Aggregation.NONE ? " ORDER BY " + ModelConstants.TS_COLUMN + " " + orderBy + " LIMIT ?" : "")); + String.join(", ", ModelConstants.getFetchColumnNames(type)) + " FROM " + ModelConstants.TS_KV_CF + + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.TS_COLUMN + " >= ? " + + "AND " + ModelConstants.TS_COLUMN + " < ?" + + (type == Aggregation.NONE ? " ORDER BY " + ModelConstants.TS_COLUMN + " " + orderBy + " LIMIT ?" : "")); } } return fetchStmts; } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java index b3cc8496a0..e9bf3db134 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java @@ -194,7 +194,7 @@ public class CassandraBaseTimeseriesLatestDao extends AbstractCassandraBaseTimes } private ListenableFuture> convertAsyncResultSetToTsKvEntryList(TbResultSet rs) { - return Futures.transform(rs.allRows(readResultsProcessingExecutor), + return Futures.transform(rs.allRows(readResultsProcessingExecutor, maxResultSetSizeBytes), rows -> this.convertResultToTsKvEntryList(rows), readResultsProcessingExecutor); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/SimpleListenableFuture.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/SimpleListenableFuture.java index 01581026d6..ded8c74b98 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/SimpleListenableFuture.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/SimpleListenableFuture.java @@ -26,4 +26,8 @@ public class SimpleListenableFuture extends AbstractFuture { return super.set(value); } + public boolean setException(Throwable throwable) { + return super.setException(throwable); + } + } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java index 495ffc131b..aa43826734 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/nosql/TimeseriesServiceNoSqlTest.java @@ -16,7 +16,10 @@ package org.thingsboard.server.dao.service.timeseries.nosql; import com.datastax.oss.driver.api.core.uuid.Uuids; +import org.apache.commons.lang3.exception.ExceptionUtils; import org.junit.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.test.util.ReflectionTestUtils; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; @@ -27,9 +30,12 @@ import org.thingsboard.server.common.data.kv.JsonDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.dao.nosql.ResultSetSizeLimitExceededException; import org.thingsboard.server.dao.service.DaoNoSqlTest; import org.thingsboard.server.dao.service.timeseries.BaseTimeseriesServiceTest; +import org.thingsboard.server.dao.timeseries.CassandraBaseTimeseriesDao; +import java.util.ArrayList; import java.util.Collections; import java.util.List; import java.util.concurrent.ExecutionException; @@ -37,6 +43,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; @@ -44,6 +51,9 @@ import static org.junit.Assert.assertTrue; @DaoNoSqlTest public class TimeseriesServiceNoSqlTest extends BaseTimeseriesServiceTest { + @Autowired + private CassandraBaseTimeseriesDao cassandraBaseTimeseriesDao; + @Test public void shouldSaveEntryOfEachTypeWithTtl() throws ExecutionException, InterruptedException, TimeoutException { long ttlInSec = TimeUnit.SECONDS.toSeconds(3); @@ -94,4 +104,36 @@ public class TimeseriesServiceNoSqlTest extends BaseTimeseriesServiceTest { double expectedValue = (doubleValue + longValue)/ 2; assertThat(listWithAgg.get(0).getDoubleValue().get()).isEqualTo(expectedValue); } + + @Test + public void testResultSetSizeLimitExceeded() throws Exception { + DeviceId deviceId = new DeviceId(Uuids.timeBased()); + + String value = "x".repeat(500); + List entries = new ArrayList<>(); + for (int i = 0; i < 5; i++) { + entries.add(new BasicTsKvEntry(TimeUnit.MINUTES.toMillis(i + 1), new StringDataEntry("bigKey", value))); + } + tsService.save(tenantId, deviceId, entries, 0).get(MAX_TIMEOUT, TimeUnit.SECONDS); + + long originalLimit = (long) ReflectionTestUtils.getField(cassandraBaseTimeseriesDao, "maxResultSetSizeBytes"); + try { + // Set a very low limit to trigger the exception + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "maxResultSetSizeBytes", 1024L); + + assertThatThrownBy(() -> tsService.findAll(tenantId, deviceId, Collections.singletonList( + new BaseReadTsKvQuery("bigKey", 0L, TimeUnit.MINUTES.toMillis(6), 1000, 10, Aggregation.NONE) + )).get(MAX_TIMEOUT, TimeUnit.SECONDS)) + .isInstanceOf(ExecutionException.class) + .satisfies(e -> assertThat(ExceptionUtils.getRootCause(e)).isInstanceOf(ResultSetSizeLimitExceededException.class)); + } finally { + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "maxResultSetSizeBytes", originalLimit); + } + + // Verify query succeeds with original limit restored + List result = tsService.findAll(tenantId, deviceId, Collections.singletonList( + new BaseReadTsKvQuery("bigKey", 0L, TimeUnit.MINUTES.toMillis(6), 1000, 10, Aggregation.NONE) + )).get(MAX_TIMEOUT, TimeUnit.SECONDS); + assertThat(result).hasSize(5); + } } diff --git a/rest-client/src/main/resources/logback.xml b/rest-client/src/main/resources/logback.xml deleted file mode 100644 index f8c9b51bd0..0000000000 --- a/rest-client/src/main/resources/logback.xml +++ /dev/null @@ -1,34 +0,0 @@ - - - - - - - - %d{ISO8601} [%thread] %-5level %logger{36} - %msg%n - - - - - - - - - - diff --git a/ui-ngx/src/app/modules/home/components/widget/lib/settings/common/color-settings-panel.component.ts b/ui-ngx/src/app/modules/home/components/widget/lib/settings/common/color-settings-panel.component.ts index ecbb03b580..0c0e504d74 100644 --- a/ui-ngx/src/app/modules/home/components/widget/lib/settings/common/color-settings-panel.component.ts +++ b/ui-ngx/src/app/modules/home/components/widget/lib/settings/common/color-settings-panel.component.ts @@ -34,6 +34,7 @@ import { IAliasController } from '@core/api/widget-api.models'; import { coerceBoolean } from '@shared/decorators/coercion'; import { DataKeysCallbacks } from '@home/components/widget/lib/settings/common/key/data-keys.component.models'; import { Datasource } from '@shared/models/widget.models'; +import { takeUntilDestroyed } from '@angular/core/rxjs-interop'; @Component({ selector: 'tb-color-settings-panel', @@ -107,6 +108,12 @@ export class ColorSettingsPanelComponent extends PageComponent implements OnInit colorFunction: [this.colorSettings?.colorFunction, []] } ); + this.colorSettingsFormGroup.get('type').valueChanges.pipe( + takeUntilDestroyed(this.destroyRef) + ).subscribe(() => { + this.updateValidators(); + setTimeout(() => {this.popover?.updatePosition();}, 0); + }); this.updateValidators(); } diff --git a/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.module.ts b/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.module.ts index 078234dbda..50d847609f 100644 --- a/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.module.ts +++ b/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.module.ts @@ -15,36 +15,47 @@ /// import { Overlay, OverlayContainer, OverlayModule } from '@angular/cdk/overlay'; -import { NgModule } from '@angular/core'; +import { inject, Injector, NgModule } from '@angular/core'; import { DEFAULT_DIALOG_CONFIG, Dialog, DialogConfig, DialogModule } from '@angular/cdk/dialog'; import { MatDialogModule } from '@angular/material/dialog'; import { DynamicDialog, DynamicMatDialog } from './dynamic-dialog'; import { DynamicOverlay } from './dynamic-overlay'; -import { DynamicOverlayContainer } from './dynamic-overlay-container'; +import { DynamicOverlayContainer, PARENT_OVERLAY_CONTAINER } from './dynamic-overlay-container'; export const DYNAMIC_MAT_DIALOG_PROVIDERS = [ - DynamicOverlayContainer, - { provide: OverlayContainer, useExisting: DynamicOverlayContainer }, - DynamicOverlay, - { provide: Overlay, useExisting: DynamicOverlay }, - DynamicDialog, - { provide: Dialog, useExisting: DynamicDialog }, - DynamicMatDialog, { - provide: DEFAULT_DIALOG_CONFIG, - useValue: { - ...new DialogConfig() + provide: DynamicMatDialog, + useFactory: () => { + const parentInjector = inject(Injector); + const parentOverlayContainer = parentInjector.get(OverlayContainer); + + const customInjector = Injector.create({ + providers: [ + { provide: PARENT_OVERLAY_CONTAINER, useValue: parentOverlayContainer }, + DynamicOverlayContainer, + { provide: OverlayContainer, useExisting: DynamicOverlayContainer }, + DynamicOverlay, + { provide: Overlay, useExisting: DynamicOverlay }, + DynamicDialog, + { provide: Dialog, useExisting: DynamicDialog }, + DynamicMatDialog, + { provide: DEFAULT_DIALOG_CONFIG, useValue: new DialogConfig() } + ], + parent: parentInjector + }); + + return customInjector.get(DynamicMatDialog); } } ]; -@NgModule( { +@NgModule({ imports: [ OverlayModule, DialogModule, MatDialogModule ], providers: DYNAMIC_MAT_DIALOG_PROVIDERS -} ) +}) export class DynamicMatDialogModule { } diff --git a/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.ts b/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.ts index f717ccb0a7..df707ad49d 100644 --- a/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.ts +++ b/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-dialog.ts @@ -34,20 +34,13 @@ export class DynamicMatDialog extends MatDialog { config.containerElement.style.transform = 'translateZ(0)'; this._customOverlay.setContainerElement(config.containerElement); } - const ref = super.open(component, config); - if (config?.containerElement) { - ref.afterClosed().subscribe( - { - next: () => { - this._customOverlay.setContainerElement(null); - }, - error: () => { - this._customOverlay.setContainerElement(null); - } - } - ); + try { + return super.open(component, config); + } finally { + if (config?.containerElement) { + this._customOverlay.setContainerElement(null); + } } - return ref; } } diff --git a/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-overlay-container.ts b/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-overlay-container.ts index 83feda56fc..99a84f5a6a 100644 --- a/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-overlay-container.ts +++ b/ui-ngx/src/app/shared/components/dialog/dynamic/dynamic-overlay-container.ts @@ -15,13 +15,21 @@ /// import { OverlayContainer } from "@angular/cdk/overlay"; -import { Injectable } from "@angular/core"; +import { inject, Injectable, InjectionToken } from "@angular/core"; + +export const PARENT_OVERLAY_CONTAINER = new InjectionToken('PARENT_OVERLAY_CONTAINER'); @Injectable() export class DynamicOverlayContainer extends OverlayContainer { - public setContainerElement( containerElement:HTMLElement ):void { + private _globalContainer = inject(PARENT_OVERLAY_CONTAINER); + private _customElement: HTMLElement | null = null; + + public override getContainerElement(): HTMLElement { + return this._customElement || this._globalContainer.getContainerElement(); + } - this._containerElement = containerElement; + setContainerElement(element: HTMLElement | null): void { + this._customElement = element; } -} +} \ No newline at end of file