diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java index 732ddb76aa..a3a048fe9e 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java @@ -120,6 +120,7 @@ import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.device.ClaimDevicesService; import org.thingsboard.server.dao.tenant.TenantProfileService; import org.thingsboard.server.dao.timeseries.TimeseriesService; +import org.thingsboard.server.queue.memory.InMemoryStorage; import org.thingsboard.server.service.entitiy.tenant.profile.TbTenantProfileService; import org.thingsboard.server.service.security.auth.jwt.RefreshTokenRequest; import org.thingsboard.server.service.security.auth.rest.LoginRequest; @@ -246,6 +247,9 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest { @SpyBean protected MailService mailService; + @Autowired + protected InMemoryStorage storage; + @Rule public TestRule watcher = new TestWatcher() { protected void starting(Description description) { @@ -372,8 +376,8 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest { } catch (Exception e) { throw new RuntimeException(e); } - Awaitility.await("tenant cleanup finish").atMost(30, TimeUnit.SECONDS) - .until(() -> attributesService.find(TenantId.SYS_TENANT_ID, tenantId, AttributeScope.SERVER_SCOPE, "test").get().isEmpty()); + Awaitility.await("all tasks processed").atMost(TIMEOUT, TimeUnit.SECONDS).during(300, TimeUnit.MILLISECONDS) + .until(() -> storage.getLag("tb_housekeeper") == 0); } private List getAllTenants() throws Exception { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/memory/DefaultInMemoryStorage.java b/common/queue/src/main/java/org/thingsboard/server/queue/memory/DefaultInMemoryStorage.java index 9c7dc02b4f..f46a05ca26 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/memory/DefaultInMemoryStorage.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/memory/DefaultInMemoryStorage.java @@ -20,8 +20,10 @@ import org.springframework.stereotype.Component; import org.thingsboard.server.queue.TbQueueMsg; import java.util.ArrayList; +import java.util.Collection; import java.util.Collections; import java.util.List; +import java.util.Optional; import java.util.concurrent.BlockingQueue; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.LinkedBlockingQueue; @@ -47,6 +49,11 @@ public final class DefaultInMemoryStorage implements InMemoryStorage { return storage.values().stream().map(BlockingQueue::size).reduce(0, Integer::sum); } + @Override + public int getLag(String topic) { + return Optional.ofNullable(storage.get(topic)).map(Collection::size).orElse(0); + } + @Override public boolean put(String topic, TbQueueMsg msg) { return storage.computeIfAbsent(topic, (t) -> new LinkedBlockingQueue<>()).add(msg); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java index 0e753d1a81..81b174a0d3 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java @@ -25,6 +25,8 @@ public interface InMemoryStorage { int getLagTotal(); + int getLag(String topic); + boolean put(String topic, TbQueueMsg msg); List get(String topic) throws InterruptedException;