Browse Source

sending ws error when the telemetry queries exceed limit

pull/15160/head
dashevchenko 7 months ago
parent
commit
30b046fa30
  1. 19
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  2. 46
      application/src/test/java/org/thingsboard/server/controller/WebsocketApiTest.java

19
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);

46
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<String> 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<String> 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<TsKvEntry> tsData) throws InterruptedException {
CountDownLatch latch = new CountDownLatch(1);
tsService.saveTimeseries(TimeseriesSaveRequest.builder()

Loading…
Cancel
Save