diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java index f0f2e1eaa7..71df830f8d 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java @@ -34,6 +34,7 @@ import org.springframework.web.socket.CloseStatus; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.dao.nosql.ResultSetSizeLimitExceededException; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult; @@ -242,7 +243,10 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc @Override public void onFailure(Throwable t) { - log.warn("[{}][{}] Failed to process command", finalCtx.getSessionId(), finalCtx.getCmdId()); + log.warn("[{}][{}] Failed to process command", finalCtx.getSessionId(), finalCtx.getCmdId(), t); + if (t instanceof ResultSetSizeLimitExceededException) { + finalCtx.sendWsMsg(new EntityDataUpdate(finalCtx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), t.getMessage())); + } } }, wsCallBackExecutor); } @@ -258,7 +262,18 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc handleLatestCmd(ctx, cmd.getLatestCmd()); } if (cmd.getTsCmd() != null) { - handleTimeSeriesCmd(ctx, cmd.getTsCmd()); + Futures.addCallback(handleTimeSeriesCmd(ctx, cmd.getTsCmd()), new FutureCallback<>() { + @Override + public void onSuccess(TbEntityDataSubCtx result) {} + + @Override + public void onFailure(Throwable t) { + log.warn("[{}][{}] Failed to process timeseries command", ctx.getSessionId(), ctx.getCmdId(), t); + if (t instanceof ResultSetSizeLimitExceededException) { + ctx.sendWsMsg(new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), t.getMessage())); + } + } + }, wsCallBackExecutor); } } else { checkAndSendInitialData(ctx); diff --git a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java index 87ba0ec3e8..aa7ffaeaf5 100644 --- a/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java @@ -19,13 +19,16 @@ import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.ArrayNode; import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; import lombok.extern.slf4j.Slf4j; import org.checkerframework.checker.nullness.qual.Nullable; import org.junit.After; import org.junit.Assert; import org.junit.Before; import org.junit.Test; +import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.mock.mockito.SpyBean; import org.springframework.test.context.TestPropertySource; import org.testcontainers.shaded.org.apache.commons.lang3.RandomStringUtils; import org.thingsboard.common.util.JacksonUtil; @@ -60,7 +63,9 @@ import org.thingsboard.server.common.data.query.NumericFilterPredicate; import org.thingsboard.server.common.data.query.SingleEntityFilter; import org.thingsboard.server.common.data.query.TsValue; import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.dao.nosql.ResultSetSizeLimitExceededException; import org.thingsboard.server.dao.service.DaoSqlTest; +import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.service.subscription.SubscriptionErrorCode; import org.thingsboard.server.service.subscription.TbAttributeSubscriptionScope; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; @@ -95,6 +100,9 @@ public class WebsocketApiTest extends AbstractControllerTest { @Autowired private TelemetrySubscriptionService tsService; + @SpyBean + private TimeseriesService timeseriesService; + Device device; DeviceTypeFilter dtf; @@ -965,6 +973,44 @@ public class WebsocketApiTest extends AbstractControllerTest { } + @Test + public void testHistoryCmdSendsWsErrorOnResultSetSizeLimitExceeded() throws Exception { + ResultSetSizeLimitExceededException exception = new ResultSetSizeLimitExceededException(100L, 200L); + Mockito.doReturn(Futures.immediateFailedFuture(exception)) + .when(timeseriesService).findAllByQueries(Mockito.any(), Mockito.any(), Mockito.any()); + + List keys = List.of("temperature"); + long now = System.currentTimeMillis(); + + // Register for 2 messages: initial entity page data + error + getWsClient().registerWaitForUpdate(2); + getWsClient().sendHistoryCmd(keys, now, TimeUnit.HOURS.toMillis(1), dtf); + getWsClient().waitForUpdate(); + + EntityDataUpdate errorUpdate = JacksonUtil.fromString(getWsClient().getLastMsg(), EntityDataUpdate.class); + assertThat(errorUpdate.getErrorCode()).isEqualTo(SubscriptionErrorCode.INTERNAL_ERROR.getCode()); + assertThat(errorUpdate.getErrorMsg()).isEqualTo(exception.getMessage()); + } + + @Test + public void testTimeSeriesCmdSendsWsErrorOnResultSetSizeLimitExceeded() throws Exception { + ResultSetSizeLimitExceededException exception = new ResultSetSizeLimitExceededException(100L, 200L); + Mockito.doReturn(Futures.immediateFailedFuture(exception)) + .when(timeseriesService).findAllByQueries(Mockito.any(), Mockito.any(), Mockito.any()); + + List keys = List.of("temperature"); + long now = System.currentTimeMillis(); + + // Register for 2 messages: initial entity page data + error + getWsClient().registerWaitForUpdate(2); + getWsClient().subscribeTsUpdate(keys, now, TimeUnit.HOURS.toMillis(1), dtf); + getWsClient().waitForUpdate(); + + EntityDataUpdate errorUpdate = JacksonUtil.fromString(getWsClient().getLastMsg(), EntityDataUpdate.class); + assertThat(errorUpdate.getErrorCode()).isEqualTo(SubscriptionErrorCode.INTERNAL_ERROR.getCode()); + assertThat(errorUpdate.getErrorMsg()).isEqualTo(exception.getMessage()); + } + private void sendTelemetry(Device device, List tsData) throws InterruptedException { CountDownLatch latch = new CountDownLatch(1); tsService.saveTimeseries(TimeseriesSaveRequest.builder()