Browse Source

Do not copy latest to entity views when data was not saved on the main entity

pull/12611/head
Dmytro Skarzhynets 2 years ago
parent
commit
54afdaa07a
No known key found for this signature in database GPG Key ID: 2B51652F224037DF
  1. 10
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java
  2. 41
      application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java
  3. 12
      common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java
  4. 6
      common/util/src/main/java/org/thingsboard/common/util/DonAsynchron.java

10
application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java

@ -151,7 +151,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, request.getEntries()));
}
if (strategy.saveLatest()) {
copyLatestToEntityViews(tenantId, entityId, request.getEntries());
addMainCallback(saveFuture, __ -> copyLatestToEntityViews(tenantId, entityId, request.getEntries()));
}
return saveFuture;
}
@ -207,8 +207,8 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
}
private void copyLatestToEntityViews(TenantId tenantId, EntityId entityId, List<TsKvEntry> ts) {
if (EntityType.DEVICE.equals(entityId.getEntityType()) || EntityType.ASSET.equals(entityId.getEntityType())) {
Futures.addCallback(this.tbEntityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId),
if (entityId.getEntityType().isOneOf(EntityType.DEVICE, EntityType.ASSET)) {
Futures.addCallback(tbEntityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId),
new FutureCallback<>() {
@Override
public void onSuccess(@Nullable List<EntityView> result) {
@ -312,6 +312,10 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer
addMainCallback(saveFuture, result -> callback.onSuccess(null), callback::onFailure);
}
private <S> void addMainCallback(ListenableFuture<S> saveFuture, Consumer<S> onSuccess) {
DonAsynchron.withCallback(saveFuture, onSuccess, null, tsCallBackExecutor);
}
private <S> void addMainCallback(ListenableFuture<S> saveFuture, Consumer<S> onSuccess, Consumer<Throwable> onFailure) {
DonAsynchron.withCallback(saveFuture, onSuccess, onFailure, tsCallBackExecutor);
}

41
application/src/test/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionServiceTest.java

@ -71,6 +71,7 @@ import java.util.concurrent.ExecutorService;
import java.util.stream.LongStream;
import java.util.stream.Stream;
import static com.google.common.util.concurrent.Futures.immediateFailedFuture;
import static com.google.common.util.concurrent.Futures.immediateFuture;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.BDDMockito.given;
@ -318,6 +319,46 @@ class DefaultTelemetrySubscriptionServiceTest {
then(subscriptionManagerService).shouldHaveNoInteractions();
}
@Test
void shouldNotCopyLatestToEntityViewWhenTimeseriesSaveFailedOnMainEntity() {
// GIVEN
var entityView = new EntityView(new EntityViewId(UUID.randomUUID()));
entityView.setTenantId(tenantId);
entityView.setCustomerId(customerId);
entityView.setEntityId(entityId);
entityView.setKeys(new TelemetryEntityView(sampleTelemetry.stream().map(KvEntry::getKey).toList(), new AttributesEntityView()));
// mock that there is one entity view
lenient().when(tbEntityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId)).thenReturn(immediateFuture(List.of(entityView)));
// mock that save latest call for entity view is successful
lenient().when(tsService.saveLatest(tenantId, entityView.getId(), sampleTelemetry)).thenReturn(immediateFuture(listOfNNumbers(sampleTelemetry.size())));
// mock TPI for entity view
lenient().when(partitionService.resolve(ServiceType.TB_CORE, tenantId, entityView.getId())).thenReturn(tpi);
var request = TimeseriesSaveRequest.builder()
.tenantId(tenantId)
.customerId(customerId)
.entityId(entityId)
.entries(sampleTelemetry)
.ttl(sampleTtl)
.strategy(new TimeseriesSaveRequest.Strategy(true, true, false))
.callback(emptyCallback)
.build();
given(tsService.save(tenantId, entityId, sampleTelemetry, sampleTtl)).willReturn(immediateFailedFuture(new RuntimeException("failed to save latest on main entity")));
// WHEN
telemetryService.saveTimeseries(request);
// THEN
// should save only time series for the main entity
then(tsService).should().save(tenantId, entityId, sampleTelemetry, sampleTtl);
then(tsService).shouldHaveNoMoreInteractions();
// should not send any WS updates
then(subscriptionManagerService).shouldHaveNoInteractions();
}
@ParameterizedTest
@MethodSource("booleanCombinations")
void shouldCallCorrectApiBasedOnBooleanFlagsInTheSaveRequest(boolean saveTimeseries, boolean saveLatest, boolean sendWsUpdate) {

12
common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java

@ -84,4 +84,16 @@ public enum EntityType {
this.tableName = tableName;
}
public boolean isOneOf(EntityType... types) {
if (types == null) {
return false;
}
for (EntityType type : types) {
if (this == type) {
return true;
}
}
return false;
}
}

6
common/util/src/main/java/org/thingsboard/common/util/DonAsynchron.java

@ -36,6 +36,9 @@ public class DonAsynchron {
FutureCallback<T> callback = new FutureCallback<T>() {
@Override
public void onSuccess(T result) {
if (onSuccess == null) {
return;
}
try {
onSuccess.accept(result);
} catch (Throwable th) {
@ -45,6 +48,9 @@ public class DonAsynchron {
@Override
public void onFailure(Throwable t) {
if (onFailure == null) {
return;
}
onFailure.accept(t);
}
};

Loading…
Cancel
Save