|
|
|
@ -19,6 +19,8 @@ import com.google.common.util.concurrent.Futures; |
|
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
|
import org.junit.Before; |
|
|
|
import org.junit.Test; |
|
|
|
import org.mockito.Mockito; |
|
|
|
import org.springframework.test.util.ReflectionTestUtils; |
|
|
|
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; |
|
|
|
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; |
|
|
|
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult; |
|
|
|
@ -44,13 +46,13 @@ public class AbstractChunkedAggregationTimeseriesDaoTest { |
|
|
|
final int LIMIT = 1; |
|
|
|
final String TEMP = "temp"; |
|
|
|
final String DESC = "DESC"; |
|
|
|
AbstractChunkedAggregationTimeseriesDao tsDao; |
|
|
|
private AbstractChunkedAggregationTimeseriesDao tsDao; |
|
|
|
|
|
|
|
@Before |
|
|
|
public void setUp() throws Exception { |
|
|
|
tsDao = spy(AbstractChunkedAggregationTimeseriesDao.class); |
|
|
|
Optional<TsKvEntry> optionalListenableFuture = Optional.of(mock(TsKvEntry.class)); |
|
|
|
willReturn(optionalListenableFuture).given(tsDao).findAndAggregateAsync(any(), anyString(), anyLong(), anyLong(), anyLong(), any()); |
|
|
|
willReturn(Futures.immediateFuture(optionalListenableFuture)).given(tsDao).findAndAggregateAsync(any(), anyString(), anyLong(), anyLong(), anyLong(), any()); |
|
|
|
willReturn(Futures.immediateFuture(mock(ReadTsKvQueryResult.class))).given(tsDao).getReadTsKvQueryResultFuture(any(), any()); |
|
|
|
} |
|
|
|
|
|
|
|
@ -58,11 +60,11 @@ public class AbstractChunkedAggregationTimeseriesDaoTest { |
|
|
|
public void givenIntervalNotMultiplePeriod_whenAggregateCount_thenLastIntervalShorterThanOthersAndEqualsEndTs() { |
|
|
|
ReadTsKvQuery query = new BaseReadTsKvQuery(TEMP, 1, 3000, 2000, LIMIT, COUNT, DESC); |
|
|
|
ReadTsKvQuery subQueryFirst = new BaseReadTsKvQuery(TEMP, 1, 2001, 1001, LIMIT, COUNT, DESC); |
|
|
|
ReadTsKvQuery subQuerySecond = new BaseReadTsKvQuery(TEMP, 2001, 3001, 2501, LIMIT, COUNT, DESC); |
|
|
|
ReadTsKvQuery subQuerySecond = new BaseReadTsKvQuery(TEMP, 2001, 3000, 2501, LIMIT, COUNT, DESC); |
|
|
|
tsDao.findAllAsync(SYS_TENANT_ID, SYS_TENANT_ID, query); |
|
|
|
verify(tsDao, times(2)).findAndAggregateAsync(any(), any(), anyLong(), anyLong(), anyLong(), any()); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, subQueryFirst.getKey(), 1, 2001, getTsForReadTsKvQuery(1, 2001), COUNT); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, subQuerySecond.getKey(), 2001, 3000 + 1, getTsForReadTsKvQuery(2001, 3001), COUNT); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, subQuerySecond.getKey(), 2001, 3000, getTsForReadTsKvQuery(2001, 3000), COUNT); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
@ -72,19 +74,17 @@ public class AbstractChunkedAggregationTimeseriesDaoTest { |
|
|
|
willCallRealMethod().given(tsDao).findAllAsync(SYS_TENANT_ID, SYS_TENANT_ID, query); |
|
|
|
assertThat(tsDao.findAllAsync(SYS_TENANT_ID, SYS_TENANT_ID, query)).isNotNull(); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(any(), any(), anyLong(), anyLong(), anyLong(), any()); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, subQueryFirst.getKey(), 1, 3000 + 1, getTsForReadTsKvQuery(1, 3001), COUNT); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, subQueryFirst.getKey(), 1, 3000, getTsForReadTsKvQuery(1, 3000), COUNT); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
public void givenIntervalNotMultiplePeriod_whenAggregateCount_thenIntervalEqualsPeriodMinusOne() { |
|
|
|
ReadTsKvQuery query = new BaseReadTsKvQuery(TEMP, 1, 3000, 2999, LIMIT, COUNT, DESC); |
|
|
|
ReadTsKvQuery subQueryFirst = new BaseReadTsKvQuery(TEMP, 1, 3000, 1500, LIMIT, COUNT, DESC); |
|
|
|
ReadTsKvQuery subQuerySecond = new BaseReadTsKvQuery(TEMP, 3000, 3001, 3000, LIMIT, COUNT, DESC); |
|
|
|
willCallRealMethod().given(tsDao).findAllAsync(SYS_TENANT_ID, SYS_TENANT_ID, query); |
|
|
|
tsDao.findAllAsync(SYS_TENANT_ID, SYS_TENANT_ID, query); |
|
|
|
verify(tsDao, times(2)).findAndAggregateAsync(any(), any(), anyLong(), anyLong(), anyLong(), any()); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(any(), any(), anyLong(), anyLong(), anyLong(), any()); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, subQueryFirst.getKey(), 1, 3000, getTsForReadTsKvQuery(1, 3000), COUNT); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, subQuerySecond.getKey(), 3000, 3001, getTsForReadTsKvQuery(3000, 3001), COUNT); |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
@ -95,7 +95,7 @@ public class AbstractChunkedAggregationTimeseriesDaoTest { |
|
|
|
willCallRealMethod().given(tsDao).findAllAsync(SYS_TENANT_ID, SYS_TENANT_ID, query); |
|
|
|
tsDao.findAllAsync(SYS_TENANT_ID, SYS_TENANT_ID, query); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(any(), any(), anyLong(), anyLong(), anyLong(), any()); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, subQueryFirst.getKey(), 1, 3001, getTsForReadTsKvQuery(1, 3001), COUNT); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, subQueryFirst.getKey(), 1, 3000, getTsForReadTsKvQuery(1, 3000), COUNT); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
@ -135,7 +135,7 @@ public class AbstractChunkedAggregationTimeseriesDaoTest { |
|
|
|
willCallRealMethod().given(tsDao).findAllAsync(SYS_TENANT_ID, SYS_TENANT_ID, query); |
|
|
|
tsDao.findAllAsync(SYS_TENANT_ID, SYS_TENANT_ID, query); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(any(), any(), anyLong(), anyLong(), anyLong(), any()); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, subQueryFirst.getKey(), 1, 3001, getTsForReadTsKvQuery(1, 3001), COUNT); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, subQueryFirst.getKey(), 1, 3000, getTsForReadTsKvQuery(1, 3000), COUNT); |
|
|
|
} |
|
|
|
|
|
|
|
@Test |
|
|
|
@ -145,8 +145,7 @@ public class AbstractChunkedAggregationTimeseriesDaoTest { |
|
|
|
tsDao.findAllAsync(SYS_TENANT_ID, SYS_TENANT_ID, query); |
|
|
|
verify(tsDao, times(1000)).findAndAggregateAsync(any(), any(), anyLong(), anyLong(), anyLong(), any()); |
|
|
|
for (long i = 1; i <= 3000; i += 3) { |
|
|
|
ReadTsKvQuery querySub = new BaseReadTsKvQuery(TEMP, i, i + 3, i + (i + 3 - i) / 2, LIMIT, COUNT, DESC); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, querySub.getKey(), i, i + 3, getTsForReadTsKvQuery(i, i + 3), COUNT); |
|
|
|
verify(tsDao, times(1)).findAndAggregateAsync(SYS_TENANT_ID, TEMP, i, Math.min(i + 3, 3000), getTsForReadTsKvQuery(i, i + 3), COUNT); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|