From 4a20d153c573ae31d95df5e0cd34048933011014 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Mon, 16 May 2022 16:35:07 +0300 Subject: [PATCH 1/4] UI: Use 'Maximum entities per datasource' parameter from widget configuration instead of hardcoded value 1024 --- ui-ngx/src/app/core/api/alias-controller.ts | 10 ++++++---- ui-ngx/src/app/core/api/widget-api.models.ts | 4 +++- ui-ngx/src/app/core/api/widget-subscription.ts | 8 +++++--- ui-ngx/src/app/core/http/entity.service.ts | 6 ++++-- .../home/components/widget/lib/maps/map-widget2.ts | 8 ++++++-- .../components/widget/lib/maps/providers/image-map.ts | 1 + .../components/widget/lib/markdown-widget.component.ts | 6 ++++-- .../components/widget/widget-config.component.html | 8 ++++++++ .../home/components/widget/widget-config.component.ts | 2 ++ .../modules/home/components/widget/widget.component.ts | 3 ++- ui-ngx/src/app/shared/models/widget.models.ts | 1 + ui-ngx/src/assets/locale/locale.constant-en_US.json | 1 + 12 files changed, 43 insertions(+), 15 deletions(-) diff --git a/ui-ngx/src/app/core/api/alias-controller.ts b/ui-ngx/src/app/core/api/alias-controller.ts index 14a26388a5..7ad7d7c715 100644 --- a/ui-ngx/src/app/core/api/alias-controller.ts +++ b/ui-ngx/src/app/core/api/alias-controller.ts @@ -17,14 +17,15 @@ import { AliasInfo, IAliasController, StateControllerHolder, StateEntityInfo } from '@core/api/widget-api.models'; import { forkJoin, Observable, of, ReplaySubject, Subject } from 'rxjs'; import { Datasource, DatasourceType, datasourceTypeTranslationMap } from '@app/shared/models/widget.models'; -import { deepClone, isEqual } from '@core/utils'; +import { deepClone, isDefinedAndNotNull, isEqual } from '@core/utils'; import { EntityService } from '@core/http/entity.service'; import { UtilsService } from '@core/services/utils.service'; import { AliasFilterType, EntityAliases, SingleEntityFilter } from '@shared/models/alias.models'; import { EntityInfo } from '@shared/models/entity.models'; import { map, mergeMap } from 'rxjs/operators'; import { - defaultEntityDataPageLink, Filter, FilterInfo, filterInfoToKeyFilters, Filters, KeyFilter, singleEntityDataPageLink, + createDefaultEntityDataPageLink, + Filter, FilterInfo, filterInfoToKeyFilters, Filters, KeyFilter, singleEntityDataPageLink, updateDatasourceFromEntityInfo } from '@shared/models/query/query.models'; import { TranslateService } from '@ngx-translate/core'; @@ -322,7 +323,7 @@ export class AliasController implements IAliasController { ); } - resolveDatasources(datasources: Array, singleEntity?: boolean): Observable> { + resolveDatasources(datasources: Array, singleEntity?: boolean, pageSize = 1024): Observable> { if (!datasources || !datasources.length) { return of([]); } @@ -360,7 +361,8 @@ export class AliasController implements IAliasController { if (singleEntity) { datasource.pageLink = deepClone(singleEntityDataPageLink); } else if (!datasource.pageLink) { - datasource.pageLink = deepClone(defaultEntityDataPageLink); + pageSize = isDefinedAndNotNull(pageSize) && pageSize > 0 ? pageSize : 1024; + datasource.pageLink = createDefaultEntityDataPageLink(pageSize); } } }); diff --git a/ui-ngx/src/app/core/api/widget-api.models.ts b/ui-ngx/src/app/core/api/widget-api.models.ts index 6ae258663c..6a413f4a95 100644 --- a/ui-ngx/src/app/core/api/widget-api.models.ts +++ b/ui-ngx/src/app/core/api/widget-api.models.ts @@ -119,7 +119,7 @@ export interface IAliasController { getEntityAliasId(aliasName: string): string; getInstantAliasInfo(aliasId: string): AliasInfo; resolveSingleEntityInfo(aliasId: string): Observable; - resolveDatasources(datasources: Array, singleEntity?: boolean): Observable>; + resolveDatasources(datasources: Array, singleEntity?: boolean, pageSize?: number): Observable>; resolveAlarmSource(alarmSource: Datasource): Observable; getEntityAliases(): EntityAliases; getFilters(): Filters; @@ -184,6 +184,7 @@ export interface SubscriptionInfo { deviceName?: string; deviceNamePrefix?: string; deviceIds?: Array; + pageSize?: number; } export class WidgetSubscriptionContext { @@ -242,6 +243,7 @@ export interface WidgetSubscriptionOptions { datasourcesOptional?: boolean; hasDataPageLink?: boolean; singleEntity?: boolean; + pageSize?: number; warnOnPageDataOverflow?: boolean; ignoreDataUpdateOnIntervalTick?: boolean; targetDeviceAliasIds?: Array; diff --git a/ui-ngx/src/app/core/api/widget-subscription.ts b/ui-ngx/src/app/core/api/widget-subscription.ts index bdbddffdcb..635e98d89c 100644 --- a/ui-ngx/src/app/core/api/widget-subscription.ts +++ b/ui-ngx/src/app/core/api/widget-subscription.ts @@ -97,6 +97,7 @@ export class WidgetSubscription implements IWidgetSubscription { hasDataPageLink: boolean; singleEntity: boolean; + pageSize: number; warnOnPageDataOverflow: boolean; ignoreDataUpdateOnIntervalTick: boolean; @@ -229,6 +230,7 @@ export class WidgetSubscription implements IWidgetSubscription { this.entityDataListeners = []; this.hasDataPageLink = options.hasDataPageLink; this.singleEntity = options.singleEntity; + this.pageSize = options.pageSize; this.warnOnPageDataOverflow = options.warnOnPageDataOverflow; this.ignoreDataUpdateOnIntervalTick = options.ignoreDataUpdateOnIntervalTick; this.datasourcePages = []; @@ -387,7 +389,7 @@ export class WidgetSubscription implements IWidgetSubscription { } ); } else { - this.ctx.aliasController.resolveDatasources(this.configuredDatasources, this.singleEntity).subscribe( + this.ctx.aliasController.resolveDatasources(this.configuredDatasources, this.singleEntity, this.pageSize).subscribe( (datasources) => { this.configuredDatasources = datasources; this.prepareDataSubscriptions().subscribe( @@ -1132,7 +1134,7 @@ export class WidgetSubscription implements IWidgetSubscription { } ); } else { - this.ctx.aliasController.resolveDatasources(this.configuredDatasources, this.singleEntity).subscribe( + this.ctx.aliasController.resolveDatasources(this.configuredDatasources, this.singleEntity, this.pageSize).subscribe( (datasources) => { this.configuredDatasources = datasources; this.prepareDataSubscriptions().subscribe( @@ -1271,7 +1273,7 @@ export class WidgetSubscription implements IWidgetSubscription { totalPages: pageData.totalPages }; if (datasource.type === DatasourceType.entity && - pageData.hasNext && pageLink.pageSize > 1) { + pageData.hasNext && !this.singleEntity) { if (this.warnOnPageDataOverflow) { const message = this.ctx.translate.instant('widget.data-overflow', {count: pageData.data.length, total: pageData.totalElements}); diff --git a/ui-ngx/src/app/core/http/entity.service.ts b/ui-ngx/src/app/core/http/entity.service.ts index b71f5a95ef..1b27367ece 100644 --- a/ui-ngx/src/app/core/http/entity.service.ts +++ b/ui-ngx/src/app/core/http/entity.service.ts @@ -1319,7 +1319,8 @@ export class EntityService { pageLink = deepClone(singleEntityDataPageLink); } else { nameFilter = subscriptionInfo.entityNamePrefix; - pageLink = deepClone(defaultEntityDataPageLink); + const pageSize = isDefinedAndNotNull(subscriptionInfo.pageSize) && subscriptionInfo.pageSize > 0 ? subscriptionInfo.pageSize : 1024; + pageLink = createDefaultEntityDataPageLink(pageSize); } datasource.entityFilter = { type: AliasFilterType.entityName, @@ -1333,7 +1334,8 @@ export class EntityService { entityType: subscriptionInfo.entityType, entityList: subscriptionInfo.entityIds }; - datasource.pageLink = deepClone(defaultEntityDataPageLink); + const pageSize = isDefinedAndNotNull(subscriptionInfo.pageSize) && subscriptionInfo.pageSize > 0 ? subscriptionInfo.pageSize : 1024; + datasource.pageLink = createDefaultEntityDataPageLink(pageSize); } } diff --git a/ui-ngx/src/app/modules/home/components/widget/lib/maps/map-widget2.ts b/ui-ngx/src/app/modules/home/components/widget/lib/maps/map-widget2.ts index d603848018..bccc8f394f 100644 --- a/ui-ngx/src/app/modules/home/components/widget/lib/maps/map-widget2.ts +++ b/ui-ngx/src/app/modules/home/components/widget/lib/maps/map-widget2.ts @@ -45,7 +45,7 @@ import { TranslateService } from '@ngx-translate/core'; import { UtilsService } from '@core/services/utils.service'; import { EntityDataPageLink } from '@shared/models/query/query.models'; import { providerClass } from '@home/components/widget/lib/maps/providers'; -import { isDefined, parseFunction } from '@core/utils'; +import { isDefined, isDefinedAndNotNull, parseFunction } from '@core/utils'; import L from 'leaflet'; import { forkJoin, Observable, of } from 'rxjs'; import { AttributeService } from '@core/http/attribute.service'; @@ -87,9 +87,13 @@ export class MapWidgetController implements MapWidgetInterface { this.map.saveMarkerLocation = this.setMarkerLocation.bind(this); this.map.savePolygonLocation = this.savePolygonLocation.bind(this); this.map.saveLocation = this.saveLocation.bind(this); + let pageSize = this.settings.mapPageSize; + if (isDefinedAndNotNull(this.ctx.widgetConfig.pageSize)) { + pageSize = Math.max(pageSize, this.ctx.widgetConfig.pageSize); + } this.pageLink = { page: 0, - pageSize: this.settings.mapPageSize, + pageSize, textSearch: null, dynamic: true }; diff --git a/ui-ngx/src/app/modules/home/components/widget/lib/maps/providers/image-map.ts b/ui-ngx/src/app/modules/home/components/widget/lib/maps/providers/image-map.ts index 411258d3af..d72c2798f4 100644 --- a/ui-ngx/src/app/modules/home/components/widget/lib/maps/providers/image-map.ts +++ b/ui-ngx/src/app/modules/home/components/widget/lib/maps/providers/image-map.ts @@ -88,6 +88,7 @@ export class ImageMap extends LeafletMap { const imageUrlSubscriptionOptions: WidgetSubscriptionOptions = { datasources, hasDataPageLink: true, + singleEntity: true, useDashboardTimewindow: false, type: widgetType.latest, callbacks: { diff --git a/ui-ngx/src/app/modules/home/components/widget/lib/markdown-widget.component.ts b/ui-ngx/src/app/modules/home/components/widget/lib/markdown-widget.component.ts index 54988775c3..daf2e40221 100644 --- a/ui-ngx/src/app/modules/home/components/widget/lib/markdown-widget.component.ts +++ b/ui-ngx/src/app/modules/home/components/widget/lib/markdown-widget.component.ts @@ -26,7 +26,7 @@ import { fillDataPattern, flatFormattedData, formattedDataFormDatasourceData, - hashCode, + hashCode, isDefinedAndNotNull, isNotEmptyStr, parseFunction, processDataPattern, safeExecute @@ -83,9 +83,11 @@ export class MarkdownWidgetComponent extends PageComponent implements OnInit { cssParser.cssPreviewNamespace = this.markdownClass; cssParser.createStyleElement(this.markdownClass, cssString); } + const pageSize = isDefinedAndNotNull(this.ctx.widgetConfig.pageSize) && + this.ctx.widgetConfig.pageSize > 0 ? this.ctx.widgetConfig.pageSize : 16384; const pageLink: EntityDataPageLink = { page: 0, - pageSize: 16384, + pageSize, textSearch: null, dynamic: true }; diff --git a/ui-ngx/src/app/modules/home/components/widget/widget-config.component.html b/ui-ngx/src/app/modules/home/components/widget/widget-config.component.html index b3df4de3be..b3ef13bbae 100644 --- a/ui-ngx/src/app/modules/home/components/widget/widget-config.component.html +++ b/ui-ngx/src/app/modules/home/components/widget/widget-config.component.html @@ -317,6 +317,14 @@ widget-config.data-settings +
+ + widget-config.data-page-size + + +
widget-config.units diff --git a/ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts b/ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts index 20ec1d6340..db593e64f2 100644 --- a/ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts +++ b/ui-ngx/src/app/modules/home/components/widget/widget-config.component.ts @@ -213,6 +213,7 @@ export class WidgetConfigComponent extends PageComponent implements OnInit, Cont widgetStyle: [null, []], widgetCss: [null, []], titleStyle: [null, []], + pageSize: [1024, [Validators.min(1), Validators.pattern(/^\d*$/)]], units: [null, []], decimals: [null, [Validators.min(0), Validators.max(15), Validators.pattern(/^\d*$/)]], noDataDisplayMessage: [null, []], @@ -420,6 +421,7 @@ export class WidgetConfigComponent extends PageComponent implements OnInit, Cont fontSize: '16px', fontWeight: 400 }, + pageSize: isDefined(config.pageSize) ? config.pageSize : 1024, units: config.units, decimals: config.decimals, noDataDisplayMessage: isDefined(config.noDataDisplayMessage) ? config.noDataDisplayMessage : '', diff --git a/ui-ngx/src/app/modules/home/components/widget/widget.component.ts b/ui-ngx/src/app/modules/home/components/widget/widget.component.ts index da8fa8cbd0..11f168a3b0 100644 --- a/ui-ngx/src/app/modules/home/components/widget/widget.component.ts +++ b/ui-ngx/src/app/modules/home/components/widget/widget.component.ts @@ -977,7 +977,8 @@ export class WidgetComponent extends PageComponent implements OnInit, AfterViewI ignoreDataUpdateOnIntervalTick: this.typeParameters.ignoreDataUpdateOnIntervalTick, comparisonEnabled: comparisonSettings.comparisonEnabled, timeForComparison: comparisonSettings.timeForComparison, - comparisonCustomIntervalValue: comparisonSettings.comparisonCustomIntervalValue + comparisonCustomIntervalValue: comparisonSettings.comparisonCustomIntervalValue, + pageSize: this.widget.config.pageSize }; if (this.widget.type === widgetType.alarm) { options.alarmSource = deepClone(this.widget.config.alarmSource); diff --git a/ui-ngx/src/app/shared/models/widget.models.ts b/ui-ngx/src/app/shared/models/widget.models.ts index 3a7ef4426f..12167a2f18 100644 --- a/ui-ngx/src/app/shared/models/widget.models.ts +++ b/ui-ngx/src/app/shared/models/widget.models.ts @@ -554,6 +554,7 @@ export interface WidgetConfig { units?: string; decimals?: number; noDataDisplayMessage?: string; + pageSize?: number; actions?: {[actionSourceId: string]: Array}; settings?: WidgetSettings; alarmSource?: Datasource; diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 904f5e2fa0..cf576aedbf 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -3246,6 +3246,7 @@ "advanced-settings": "Advanced settings", "data-settings": "Data settings", "no-data-display-message": "\"No data to display\" alternative message", + "data-page-size": "Maximum entities per datasource", "settings-component-not-found": "Settings form component not found for selector '{{selector}}'" }, "widget-type": { From f30e769ecc49e1b2c2e7e2872f0dec57e913d704 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 16 May 2022 14:39:43 +0300 Subject: [PATCH 2/4] Fix EntityViewController test --- .../BaseEntityViewControllerTest.java | 53 ++++++++++++------- 1 file changed, 33 insertions(+), 20 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseEntityViewControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseEntityViewControllerTest.java index 155c768cb0..615914c6d7 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseEntityViewControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseEntityViewControllerTest.java @@ -62,6 +62,7 @@ import java.util.Set; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; +import static java.util.concurrent.TimeUnit.HOURS; import static java.util.concurrent.TimeUnit.MILLISECONDS; import static java.util.concurrent.TimeUnit.SECONDS; import static org.assertj.core.api.Assertions.assertThat; @@ -368,7 +369,7 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes public void testTheCopyOfAttrsIntoTSForTheView() throws Exception { Set expectedActualAttributesSet = Set.of("caKey1", "caKey2", "caKey3", "caKey4"); Set actualAttributesSet = - getAttributesByKeys("{\"caKey1\":\"value1\", \"caKey2\":true, \"caKey3\":42.0, \"caKey4\":73}", expectedActualAttributesSet); + putAttributesAndWait("{\"caKey1\":\"value1\", \"caKey2\":true, \"caKey3\":42.0, \"caKey4\":73}", expectedActualAttributesSet); log.debug("got correct actualAttributesSet, saving new entity view..."); EntityView savedView = getNewSavedEntityView("Test entity view"); @@ -389,13 +390,15 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes @Test public void testTheCopyOfAttrsOutOfTSForTheView() throws Exception { + long now = System.currentTimeMillis(); Set expectedActualAttributesSet = Set.of("caKey1", "caKey2", "caKey3", "caKey4"); Set actualAttributesSet = - getAttributesByKeys("{\"caKey1\":\"value1\", \"caKey2\":true, \"caKey3\":42.0, \"caKey4\":73}", expectedActualAttributesSet); + putAttributesAndWait("{\"caKey1\":\"value1\", \"caKey2\":true, \"caKey3\":42.0, \"caKey4\":73}", expectedActualAttributesSet); - List> valueTelemetryOfDevices = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + testDevice.getId().getId().toString() + - "/values/attributes?keys=" + String.join(",", actualAttributesSet), new TypeReference<>() { + List> values = doGetAsyncTyped("/api/plugins/telemetry/DEVICE/" + testDevice.getId() + + "/values/attributes?keys=" + String.join(",", expectedActualAttributesSet), new TypeReference<>() { }); + assertEquals(expectedActualAttributesSet.size(), values.size()); EntityView view = new EntityView(); view.setEntityId(testDevice.getId()); @@ -403,12 +406,12 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes view.setName("Test entity view"); view.setType("default"); view.setKeys(telemetry); - view.setStartTimeMs((long) getValue(valueTelemetryOfDevices, "lastActivityTime") * 10); - view.setEndTimeMs((long) getValue(valueTelemetryOfDevices, "lastActivityTime") / 10); + view.setStartTimeMs(now - HOURS.toMillis(1)); + view.setEndTimeMs(now - 1); EntityView savedView = doPost("/api/entityView", view, EntityView.class); - List> values = doGetAsyncTyped("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() + - "/values/attributes?keys=" + String.join(",", actualAttributesSet), new TypeReference<>() { + values = doGetAsyncTyped("/api/plugins/telemetry/ENTITY_VIEW/" + savedView.getId().getId().toString() + + "/values/attributes?keys=" + String.join(",", expectedActualAttributesSet), new TypeReference<>() { }); assertEquals(0, values.size()); } @@ -431,19 +434,19 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes uploadTelemetry("{\"tsKey1\":\"value1\", \"tsKey2\":true, \"tsKey3\":40.0}", accessToken); getWsClient().waitForUpdate(); - long startTimeMs = System.currentTimeMillis(); + long startTimeMs = getCurTsButNotPrevTs(now); getWsClient().registerWaitForUpdate(); uploadTelemetry("{\"tsKey1\":\"value2\", \"tsKey2\":false, \"tsKey3\":80.0}", accessToken); getWsClient().waitForUpdate(); - Thread.sleep(3); + long middleOfTestMs = getCurTsButNotPrevTs(startTimeMs); getWsClient().registerWaitForUpdate(); uploadTelemetry("{\"tsKey1\":\"value3\", \"tsKey2\":false, \"tsKey3\":120.0}", accessToken); getWsClient().waitForUpdate(); - long endTimeMs = System.currentTimeMillis(); + long endTimeMs = getCurTsButNotPrevTs(middleOfTestMs); getWsClient().registerWaitForUpdate(); uploadTelemetry("{\"tsKey1\":\"value4\", \"tsKey2\":true, \"tsKey3\":160.0}", accessToken); getWsClient().waitForUpdate(); @@ -455,15 +458,25 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes EntityView savedView = doPost("/api/entityView", view, EntityView.class); String entityViewId = savedView.getId().getId().toString(); - Map>> expectedValues = getTelemetryValues("DEVICE", deviceId, keys, 0L, (startTimeMs + endTimeMs) / 2); - Assert.assertEquals(2, expectedValues.get("tsKey1").size()); - Assert.assertEquals(2, expectedValues.get("tsKey2").size()); - Assert.assertEquals(2, expectedValues.get("tsKey3").size()); + Map>> actualDeviceValues = getTelemetryValues("DEVICE", deviceId, keys, 0L, middleOfTestMs); + Assert.assertEquals(2, actualDeviceValues.get("tsKey1").size()); + Assert.assertEquals(2, actualDeviceValues.get("tsKey2").size()); + Assert.assertEquals(2, actualDeviceValues.get("tsKey3").size()); + + Map>> actualEntityViewValues = getTelemetryValues("ENTITY_VIEW", entityViewId, keys, 0L, middleOfTestMs); + Assert.assertEquals(1, actualEntityViewValues.get("tsKey1").size()); + Assert.assertEquals(1, actualEntityViewValues.get("tsKey2").size()); + Assert.assertEquals(1, actualEntityViewValues.get("tsKey3").size()); + } - Map>> actualValues = getTelemetryValues("ENTITY_VIEW", entityViewId, keys, 0L, (startTimeMs + endTimeMs) / 2); - Assert.assertEquals(1, actualValues.get("tsKey1").size()); - Assert.assertEquals(1, actualValues.get("tsKey2").size()); - Assert.assertEquals(1, actualValues.get("tsKey3").size()); + private static long getCurTsButNotPrevTs(long prevTs) throws InterruptedException { + long result = System.currentTimeMillis(); + if (prevTs == result) { + Thread.sleep(1); + return getCurTsButNotPrevTs(prevTs); + } else { + return result; + } } private void uploadTelemetry(String strKvs, String accessToken) throws Exception { @@ -504,7 +517,7 @@ public abstract class BaseEntityViewControllerTest extends AbstractControllerTes }); } - private Set getAttributesByKeys(String stringKV, Set expectedKeySet) throws Exception { + private Set putAttributesAndWait(String stringKV, Set expectedKeySet) throws Exception { DeviceTypeFilter dtf = new DeviceTypeFilter(testDevice.getType(), testDevice.getName()); List keysToSubscribe = expectedKeySet.stream() .map(key -> new EntityKey(EntityKeyType.CLIENT_ATTRIBUTE, key)) From ec7a36cd896c2779ed23e5192afaa43a01106a8f Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 16 May 2022 16:58:29 +0300 Subject: [PATCH 3/4] Fix race conditions on WS subscribe commands processing --- ...efaultTbEntityDataSubscriptionService.java | 75 +++++++++++-------- .../subscription/TbAbstractSubCtx.java | 18 ++++- .../subscription/TbAlarmDataSubCtx.java | 8 +- .../subscription/TbEntityCountSubCtx.java | 7 +- .../subscription/TbEntityDataSubCtx.java | 10 +-- .../DefaultTelemetryWebSocketService.java | 6 +- .../controller/BaseWebsocketApiTest.java | 14 +--- .../controller/TbTestWebSocketClient.java | 4 +- 8 files changed, 78 insertions(+), 64 deletions(-) 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 e0eb4820fd..82d9eb779a 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 @@ -5,7 +5,7 @@ * 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 + * 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, @@ -228,7 +228,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } } else if (!theCtx.isInitialDataSent()) { EntityDataUpdate update = new EntityDataUpdate(theCtx.getCmdId(), theCtx.getData(), null, theCtx.getMaxEntitiesPerDataSubscription()); - wsService.sendWsMsg(theCtx.getSessionId(), update); + theCtx.sendWsMsg(update); theCtx.setInitialDataSent(true); } } catch (RuntimeException e) { @@ -287,7 +287,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc ctx.clearEntitySubscriptions(); if (entities.isEmpty()) { AlarmDataUpdate update = new AlarmDataUpdate(cmd.getCmdId(), new PageData<>(), null, 0, 0); - wsService.sendWsMsg(ctx.getSessionId(), update); + ctx.sendWsMsg(update); } else { ctx.fetchAlarms(); ctx.createLatestValuesSubscriptions(cmd.getQuery().getLatestValues()); @@ -420,22 +420,26 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } } catch (InterruptedException | ExecutionException e) { log.warn("[{}][{}][{}] Failed to fetch historical data", ctx.getSessionId(), ctx.getCmdId(), entityData.getEntityId(), e); - wsService.sendWsMsg(ctx.getSessionId(), - new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to fetch historical data!")); + ctx.sendWsMsg(new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to fetch historical data!")); } }); - EntityDataUpdate update; - if (!ctx.isInitialDataSent()) { - update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription()); - ctx.setInitialDataSent(true); - } else { - update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription()); - } - wsService.sendWsMsg(ctx.getSessionId(), update); - if (subscribe) { - ctx.createTimeseriesSubscriptions(keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()), cmd.getStartTs(), cmd.getEndTs()); + ctx.getWsLock().lock(); + try { + EntityDataUpdate update; + if (!ctx.isInitialDataSent()) { + update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription()); + ctx.setInitialDataSent(true); + } else { + update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription()); + } + if (subscribe) { + ctx.createTimeseriesSubscriptions(keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()), cmd.getStartTs(), cmd.getEndTs()); + } + ctx.sendWsMsg(update); + ctx.getData().getData().forEach(ed -> ed.getTimeseries().clear()); + } finally { + ctx.getWsLock().unlock(); } - ctx.getData().getData().forEach(ed -> ed.getTimeseries().clear()); return ctx; }, wsCallBackExecutor); } @@ -464,7 +468,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc ListenableFuture> missingTsData = tsService.findLatest(ctx.getTenantId(), entityData.getEntityId(), missingTsKeys); missingTelemetryFutures.put(entityData, Futures.transform(missingTsData, this::toTsValue, MoreExecutors.directExecutor())); } - Futures.addCallback(Futures.allAsList(missingTelemetryFutures.values()), new FutureCallback>>() { + Futures.addCallback(Futures.allAsList(missingTelemetryFutures.values()), new FutureCallback<>() { @Override public void onSuccess(@Nullable List> result) { missingTelemetryFutures.forEach((key, value) -> { @@ -475,30 +479,39 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } }); EntityDataUpdate update; - if (!ctx.isInitialDataSent()) { - update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription()); - ctx.setInitialDataSent(true); - } else { - update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription()); + ctx.getWsLock().lock(); + try { + ctx.createLatestValuesSubscriptions(latestCmd.getKeys()); + if (!ctx.isInitialDataSent()) { + update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription()); + ctx.setInitialDataSent(true); + } else { + update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription()); + } + ctx.sendWsMsg(update); + } finally { + ctx.getWsLock().unlock(); } - wsService.sendWsMsg(ctx.getSessionId(), update); - ctx.createLatestValuesSubscriptions(latestCmd.getKeys()); } @Override public void onFailure(Throwable t) { log.warn("[{}][{}] Failed to process websocket command: {}:{}", ctx.getSessionId(), ctx.getCmdId(), ctx.getQuery(), latestCmd, t); - wsService.sendWsMsg(ctx.getSessionId(), - new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to process websocket command!")); + ctx.sendWsMsg(new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to process websocket command!")); } }, wsCallBackExecutor); } else { - if (!ctx.isInitialDataSent()) { - EntityDataUpdate update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription()); - wsService.sendWsMsg(ctx.getSessionId(), update); - ctx.setInitialDataSent(true); + ctx.getWsLock().lock(); + try { + ctx.createLatestValuesSubscriptions(latestCmd.getKeys()); + if (!ctx.isInitialDataSent()) { + EntityDataUpdate update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription()); + ctx.sendWsMsg(update); + ctx.setInitialDataSent(true); + } + } finally { + ctx.getWsLock().unlock(); } - ctx.createLatestValuesSubscriptions(latestCmd.getKeys()); } } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java index 749825f644..a8d2472eb0 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java @@ -5,7 +5,7 @@ * 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 + * 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, @@ -41,6 +41,7 @@ import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; +import org.thingsboard.server.service.telemetry.cmd.v2.CmdUpdate; import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; import java.util.ArrayList; @@ -52,14 +53,18 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; @Slf4j @Data public abstract class TbAbstractSubCtx { + @Getter + protected final Lock wsLock = new ReentrantLock(true); protected final String serviceId; protected final SubscriptionServiceStatistics stats; - protected final TelemetryWebSocketService wsService; + private final TelemetryWebSocketService wsService; protected final EntityService entityService; protected final TbLocalSubscriptionService localSubscriptionService; protected final AttributesService attributesService; @@ -314,4 +319,13 @@ public abstract class TbAbstractSubCtx { private final String sourceAttribute; } + public void sendWsMsg(CmdUpdate update) { + wsLock.lock(); + try { + wsService.sendWsMsg(sessionRef.getSessionId(), update); + } finally { + wsLock.unlock(); + } + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java index d87f1557ff..8655a91e96 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java @@ -5,7 +5,7 @@ * 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 + * 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, @@ -117,7 +117,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { } else { update = new AlarmDataUpdate(cmdId, new PageData<>(), null, maxEntitiesPerAlarmSubscription, data.getTotalElements()); } - wsService.sendWsMsg(getSessionId(), update); + sendWsMsg(update); } public void fetchData() { @@ -198,7 +198,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { return alarm; }).collect(Collectors.toList()); if (!update.isEmpty()) { - wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, update, maxEntitiesPerAlarmSubscription, data.getTotalElements())); + sendWsMsg(new AlarmDataUpdate(cmdId, null, update, maxEntitiesPerAlarmSubscription, data.getTotalElements())); } } else { log.trace("[{}][{}][{}][{}] Received stale subscription update: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), keyType, subscriptionUpdate); @@ -222,7 +222,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { AlarmData updated = new AlarmData(alarm, current.getOriginatorName(), current.getEntityId()); updated.getLatest().putAll(current.getLatest()); alarmsMap.put(alarmId, updated); - wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, Collections.singletonList(updated), maxEntitiesPerAlarmSubscription, data.getTotalElements())); + sendWsMsg(new AlarmDataUpdate(cmdId, null, Collections.singletonList(updated), maxEntitiesPerAlarmSubscription, data.getTotalElements())); } else { fetchAlarms(); } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityCountSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityCountSubCtx.java index 53df7c4aaf..7fc87eb17c 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityCountSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityCountSubCtx.java @@ -17,14 +17,11 @@ package org.thingsboard.server.service.subscription; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.query.EntityCountQuery; -import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.entity.EntityService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketService; import org.thingsboard.server.service.telemetry.TelemetryWebSocketSessionRef; import org.thingsboard.server.service.telemetry.cmd.v2.EntityCountUpdate; -import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; -import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; @Slf4j public class TbEntityCountSubCtx extends TbAbstractSubCtx { @@ -40,7 +37,7 @@ public class TbEntityCountSubCtx extends TbAbstractSubCtx { @Override public void fetchData() { result = (int) entityService.countEntitiesByQuery(getTenantId(), getCustomerId(), query); - wsService.sendWsMsg(sessionRef.getSessionId(), new EntityCountUpdate(cmdId, result)); + sendWsMsg(new EntityCountUpdate(cmdId, result)); } @Override @@ -48,7 +45,7 @@ public class TbEntityCountSubCtx extends TbAbstractSubCtx { int newCount = (int) entityService.countEntitiesByQuery(getTenantId(), getCustomerId(), query); if (newCount != result) { result = newCount; - wsService.sendWsMsg(sessionRef.getSessionId(), new EntityCountUpdate(cmdId, result)); + sendWsMsg(new EntityCountUpdate(cmdId, result)); } } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java index b8211e50cc..6423604c4e 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java @@ -5,7 +5,7 @@ * 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 + * 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, @@ -51,7 +51,7 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { @Getter @Setter - private boolean initialDataSent; + private volatile boolean initialDataSent; private TimeSeriesCmd curTsCmd; private LatestValueCmd latestValueCmd; @Getter @@ -121,7 +121,7 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { if (!latestUpdate.isEmpty()) { Map> latestMap = Collections.singletonMap(keyType, latestUpdate); entityData = new EntityData(entityId, latestMap, null); - wsService.sendWsMsg(sessionId, new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData), maxEntitiesPerDataSubscription)); + sendWsMsg(new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData), maxEntitiesPerDataSubscription)); } } @@ -162,7 +162,7 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { Map tsMap = new HashMap<>(); tsUpdate.forEach((key, tsValue) -> tsMap.put(key, tsValue.toArray(new TsValue[tsValue.size()]))); EntityData entityData = new EntityData(entityId, null, tsMap); - wsService.sendWsMsg(sessionId, new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData), maxEntitiesPerDataSubscription)); + sendWsMsg(new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData), maxEntitiesPerDataSubscription)); } } @@ -219,9 +219,9 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { } } } - wsService.sendWsMsg(sessionRef.getSessionId(), new EntityDataUpdate(cmdId, data, null, maxEntitiesPerDataSubscription)); subIdsToCancel.forEach(subId -> localSubscriptionService.cancelSubscription(getSessionId(), subId)); subsToAdd.forEach(localSubscriptionService::addSubscription); + sendWsMsg(new EntityDataUpdate(cmdId, data, null, maxEntitiesPerDataSubscription)); } public void setCurrentCmd(EntityDataCmd cmd) { 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 67fe9346a7..5a50a2ea22 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 @@ -442,7 +442,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi private void handleWsAttributesSubscriptionByKeys(TelemetryWebSocketSessionRef sessionRef, AttributesSubscriptionCmd cmd, String sessionId, EntityId entityId, List keys) { - FutureCallback> callback = new FutureCallback>() { + FutureCallback> callback = new FutureCallback<>() { @Override public void onSuccess(List data) { List attributesData = data.stream().map(d -> new BasicTsKvEntry(d.getLastUpdateTs(), d)).collect(Collectors.toList()); @@ -542,7 +542,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi private void handleWsAttributesSubscription(TelemetryWebSocketSessionRef sessionRef, AttributesSubscriptionCmd cmd, String sessionId, EntityId entityId) { - FutureCallback> callback = new FutureCallback>() { + FutureCallback> callback = new FutureCallback<>() { @Override public void onSuccess(List data) { List attributesData = data.stream().map(d -> new BasicTsKvEntry(d.getLastUpdateTs(), d)).collect(Collectors.toList()); @@ -666,7 +666,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi } private FutureCallback> getSubscriptionCallback(final TelemetryWebSocketSessionRef sessionRef, final TimeseriesSubscriptionCmd cmd, final String sessionId, final EntityId entityId, final long startTs, final List keys) { - return new FutureCallback>() { + return new FutureCallback<>() { @Override public void onSuccess(List data) { sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); 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 90d6ac7b9d..f7838d9689 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java @@ -102,7 +102,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { List tsData = Arrays.asList(dataPoint1, dataPoint2, dataPoint3); sendTelemetry(device, tsData); - Thread.sleep(100); update = getWsClient().sendHistoryCmd(keys, now, TimeUnit.HOURS.toMillis(1), dtf); @@ -136,7 +135,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { List tsData = Arrays.asList(dataPoint1, dataPoint2, dataPoint3); sendTelemetry(device, tsData); - Thread.sleep(100); update = getWsClient().subscribeTsUpdate(List.of("temperature"), now, TimeUnit.HOURS.toMillis(1)); Assert.assertEquals(1, update.getCmdId()); @@ -153,7 +151,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { now = System.currentTimeMillis(); TsKvEntry dataPoint4 = new BasicTsKvEntry(now, new LongDataEntry("temperature", 45L)); getWsClient().registerWaitForUpdate(); - Thread.sleep(100); sendTelemetry(device, Arrays.asList(dataPoint4)); String msg = getWsClient().waitForUpdate(); @@ -309,13 +306,12 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { Assert.assertEquals(0, pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature").getTs()); Assert.assertEquals("", pageData.getData().get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("temperature").getValue()); + getWsClient().registerWaitForUpdate(); TsKvEntry dataPoint1 = new BasicTsKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("temperature", 42L)); List tsData = Arrays.asList(dataPoint1); sendTelemetry(device, tsData); - Thread.sleep(100); - - update = getWsClient().subscribeLatestUpdate(keys, dtf); + update = getWsClient().parseDataReply(getWsClient().waitForUpdate()); Assert.assertEquals(1, update.getCmdId()); @@ -329,7 +325,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { now = System.currentTimeMillis(); TsKvEntry dataPoint2 = new BasicTsKvEntry(now, new LongDataEntry("temperature", 52L)); - getWsClient().registerWaitForUpdate(); sendTelemetry(device, Arrays.asList(dataPoint2)); update = getWsClient().parseDataReply(getWsClient().waitForUpdate()); @@ -371,7 +366,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { Assert.assertEquals("", pageData.getData().get(0).getLatest().get(EntityKeyType.SERVER_ATTRIBUTE).get("serverAttributeKey").getValue()); getWsClient().registerWaitForUpdate(); - Thread.sleep(500); AttributeKvEntry dataPoint1 = new BaseAttributeKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("serverAttributeKey", 42L)); List tsData = Arrays.asList(dataPoint1); @@ -394,7 +388,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { AttributeKvEntry dataPoint2 = new BaseAttributeKvEntry(now, new LongDataEntry("serverAttributeKey", 52L)); getWsClient().registerWaitForUpdate(); - Thread.sleep(500); sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint2)); msg = getWsClient().waitForUpdate(); Assert.assertNotNull(msg); @@ -411,14 +404,12 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { //Sending update from the past, while latest value has new timestamp; getWsClient().registerWaitForUpdate(); - Thread.sleep(500); sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint1)); msg = getWsClient().waitForUpdate(TimeUnit.SECONDS.toMillis(1)); Assert.assertNull(msg); //Sending duplicate update again getWsClient().registerWaitForUpdate(); - Thread.sleep(500); sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, Arrays.asList(dataPoint2)); msg = getWsClient().waitForUpdate(TimeUnit.SECONDS.toMillis(1)); Assert.assertNull(msg); @@ -456,7 +447,6 @@ public abstract class BaseWebsocketApiTest extends AbstractControllerTest { getWsClient().registerWaitForUpdate(); AttributeKvEntry dataPoint1 = new BaseAttributeKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("serverAttributeKey", 42L)); List tsData = Arrays.asList(dataPoint1); - Thread.sleep(100); sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, tsData); 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 6af6176b06..e2031e701f 100644 --- a/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java +++ b/application/src/test/java/org/thingsboard/server/controller/TbTestWebSocketClient.java @@ -44,8 +44,8 @@ import java.util.concurrent.TimeUnit; public class TbTestWebSocketClient extends WebSocketClient { private volatile String lastMsg; - private CountDownLatch reply; - private CountDownLatch update; + private volatile CountDownLatch reply; + private volatile CountDownLatch update; public TbTestWebSocketClient(URI serverUri) { super(serverUri); From 097862b59f8aa0c1d92e1e97f719142d4bea1b9d Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Mon, 16 May 2022 17:00:34 +0300 Subject: [PATCH 4/4] Licnese headers --- .../subscription/DefaultTbEntityDataSubscriptionService.java | 2 +- .../server/service/subscription/TbAbstractSubCtx.java | 2 +- .../server/service/subscription/TbAlarmDataSubCtx.java | 2 +- .../server/service/subscription/TbEntityDataSubCtx.java | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) 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 82d9eb779a..507ddb3834 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 @@ -5,7 +5,7 @@ * 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 + * 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, diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java index a8d2472eb0..15e7f754c5 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractSubCtx.java @@ -5,7 +5,7 @@ * 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 + * 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, diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java index 8655a91e96..5804002c77 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java @@ -5,7 +5,7 @@ * 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 + * 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, diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java index 6423604c4e..385661f177 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java @@ -5,7 +5,7 @@ * 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 + * 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,