@ -20,6 +20,7 @@ import com.fasterxml.jackson.databind.node.TextNode;
import org.junit.After ;
import org.junit.Before ;
import org.junit.Test ;
import org.mockito.ArgumentMatcher ;
import org.springframework.beans.factory.annotation.Autowired ;
import org.springframework.boot.test.mock.mockito.SpyBean ;
import org.springframework.test.context.TestPropertySource ;
@ -33,6 +34,7 @@ import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.alarm.AlarmSeverity ;
import org.thingsboard.server.common.data.event.EventType ;
import org.thingsboard.server.common.data.event.LifecycleEvent ;
import org.thingsboard.server.common.data.housekeeper.HousekeeperTask ;
import org.thingsboard.server.common.data.housekeeper.HousekeeperTaskType ;
import org.thingsboard.server.common.data.id.EntityId ;
import org.thingsboard.server.common.data.id.RuleChainId ;
@ -56,33 +58,46 @@ import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.dao.alarm.AlarmService ;
import org.thingsboard.server.dao.attributes.AttributesService ;
import org.thingsboard.server.dao.event.EventService ;
import org.thingsboard.server.dao.housekeeper.HousekeeperService ;
import org.thingsboard.server.dao.rule.RuleChainService ;
import org.thingsboard.server.dao.service.DaoSqlTest ;
import org.thingsboard.server.dao.timeseries.TimeseriesService ;
import org.thingsboard.server.gen.transport.TransportProtos.HousekeeperTaskProto ;
import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg ;
import org.thingsboard.server.service.housekeeper.processor.TelemetryDeletionTaskProcessor ;
import java.util.Arrays ;
import java.util.Collections ;
import java.util.List ;
import java.util.Optional ;
import java.util.concurrent.TimeUnit ;
import java.util.concurrent.TimeoutException ;
import java.util.function.Function ;
import java.util.function.Predicate ;
import java.util.stream.Collectors ;
import static org.assertj.core.api.Assertions.assertThat ;
import static org.awaitility.Awaitility.await ;
import static org.mockito.ArgumentMatchers.any ;
import static org.mockito.ArgumentMatchers.argThat ;
import static org.mockito.Mockito.doCallRealMethod ;
import static org.mockito.Mockito.doThrow ;
import static org.mockito.Mockito.verify ;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status ;
@DaoSqlTest
@TestPropertySource ( properties = {
"transport.http.enabled=true"
"queue.core.housekeeper.enabled=true" ,
"transport.http.enabled=true" ,
"queue.core.housekeeper.reprocessing-start-delay-sec=1" ,
"queue.core.housekeeper.task-reprocessing-delay-sec=2" ,
"queue.core.housekeeper.poll-interval-ms=1000"
} )
public class HousekeeperServiceTest extends AbstractControllerTest {
@SpyBean
private HousekeeperService housekeeperService ;
private DefaultHousekeeperService housekeeperService ;
@SpyBean
private HousekeeperReprocessingService housekeeperReprocessingService ;
@Autowired
private EventService eventService ;
@Autowired
@ -93,6 +108,8 @@ public class HousekeeperServiceTest extends AbstractControllerTest {
private RuleChainService ruleChainService ;
@Autowired
private AlarmService alarmService ;
@SpyBean
private TelemetryDeletionTaskProcessor telemetryDeletionTaskProcessor ;
private TenantId tenantId ;
@ -112,7 +129,7 @@ public class HousekeeperServiceTest extends AbstractControllerTest {
@Test
public void whenDeviceIsDeleted_thenCleanUpRelatedData ( ) throws Exception {
Device device = createDevice ( "test" , "test " ) ;
Device device = createDevice ( "wekfwepf" , "wekfwepf " ) ;
createRelatedData ( device . getId ( ) ) ;
doDelete ( "/api/device/" + device . getId ( ) ) . andExpect ( status ( ) . isOk ( ) ) ;
@ -143,7 +160,7 @@ public class HousekeeperServiceTest extends AbstractControllerTest {
@Test
public void whenUserIsDeleted_thenCleanUpRelatedData ( ) throws Exception {
Device device = createDevice ( "test" , "test " ) ;
Device device = createDevice ( "vneoruvhwe" , "vneoruvhwe " ) ;
UserId userId = customerUserId ;
createRelatedData ( userId ) ;
Alarm alarm = Alarm . builder ( )
@ -172,7 +189,7 @@ public class HousekeeperServiceTest extends AbstractControllerTest {
createRelatedData ( tenantId ) ;
Device device = createDevice ( "test" , "test " ) ;
Device device = createDevice ( "oi324rujoi" , "oi324rujoi " ) ;
createRelatedData ( device . getId ( ) ) ;
RuleChainMetaData ruleChainMetaData = createRuleChain ( ) ;
@ -199,6 +216,45 @@ public class HousekeeperServiceTest extends AbstractControllerTest {
} ) ;
}
@Test
public void whenTaskProcessingFails_thenReprocessUntilSuccessful ( ) throws Exception {
TimeoutException error = new TimeoutException ( "Test timeout" ) ;
doThrow ( error ) . when ( telemetryDeletionTaskProcessor ) . process ( any ( ) ) ;
Device device = createDevice ( "vep9ruv32" , "vep9ruv32" ) ;
createRelatedData ( device . getId ( ) ) ;
doDelete ( "/api/device/" + device . getId ( ) ) . andExpect ( status ( ) . isOk ( ) ) ;
int attempts = 3 ;
await ( ) . atMost ( 30 , TimeUnit . SECONDS ) . untilAsserted ( ( ) - > {
verify ( housekeeperService ) . processTask ( argThat ( verifyTaskSubmission ( device . getId ( ) , HousekeeperTaskType . DELETE_TELEMETRY ,
task - > task . getErrorsCount ( ) = = 0 ) ) ) ;
for ( int i = 1 ; i < = attempts ; i + + ) {
int attempt = i ;
verify ( housekeeperReprocessingService ) . submitForReprocessing ( argThat ( verifyTaskSubmission ( device . getId ( ) , HousekeeperTaskType . DELETE_TELEMETRY ,
task - > task . getErrorsCount ( ) > 0 & & task . getAttempt ( ) = = attempt ) ) , argThat ( e - > e . getMessage ( ) . equals ( error . getMessage ( ) ) ) ) ;
verify ( housekeeperService ) . processTask ( argThat ( verifyTaskSubmission ( device . getId ( ) , HousekeeperTaskType . DELETE_TELEMETRY ,
task - > task . getErrorsCount ( ) > 0 & & task . getAttempt ( ) = = attempt ) ) ) ;
}
} ) ;
assertThat ( getTimeseriesHistory ( device . getId ( ) ) ) . isNotEmpty ( ) ;
doCallRealMethod ( ) . when ( telemetryDeletionTaskProcessor ) . process ( any ( ) ) ; // fixing the code
await ( ) . atMost ( 30 , TimeUnit . SECONDS ) . untilAsserted ( ( ) - > {
assertThat ( getTimeseriesHistory ( device . getId ( ) ) ) . isEmpty ( ) ;
} ) ;
}
private ArgumentMatcher < ToHousekeeperServiceMsg > verifyTaskSubmission ( EntityId entityId , HousekeeperTaskType taskType ,
Predicate < HousekeeperTaskProto > additionalCheck ) {
return msg - > {
HousekeeperTask task = JacksonUtil . fromString ( msg . getTask ( ) . getValue ( ) , HousekeeperTask . class ) ;
return task . getEntityId ( ) . equals ( entityId ) & & task . getTaskType ( ) = = taskType & & additionalCheck . test ( msg . getTask ( ) ) ;
} ;
}
private void createRelatedData ( EntityId entityId ) throws Exception {
createTelemetry ( entityId ) ;
for ( String scope : List . of ( DataConstants . SERVER_SCOPE , DataConstants . SHARED_SCOPE , DataConstants . CLIENT_SCOPE ) ) {