Browse Source

InMemoryStorage: added getLagTotal to be able to await while queue have any messages (no guarantee that messages was processed).

pull/5383/head
Sergey Matvienko 5 years ago
parent
commit
d98419106b
  1. 4
      common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java
  2. 39
      common/queue/src/test/java/org/thingsboard/server/queue/memory/InMemoryStorageTest.java

4
common/queue/src/main/java/org/thingsboard/server/queue/memory/InMemoryStorage.java

@ -42,6 +42,10 @@ public final class InMemoryStorage {
});
}
public int getLagTotal() {
return storage.values().stream().map(BlockingQueue::size).reduce(0, Integer::sum);
}
public static InMemoryStorage getInstance() {
if (instance == null) {
synchronized (InMemoryStorage.class) {

39
common/queue/src/test/java/org/thingsboard/server/queue/memory/InMemoryStorageTest.java

@ -0,0 +1,39 @@
package org.thingsboard.server.queue.memory;
import org.junit.After;
import org.junit.Before;
import org.junit.Test;
import org.thingsboard.server.queue.TbQueueMsg;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.mock;
public class InMemoryStorageTest {
InMemoryStorage storage = InMemoryStorage.getInstance();
@Before
public void setUp() {
storage.cleanup();
}
@After
public void tearDown() {
storage.cleanup();
}
@Test
public void givenStorage_whenGetLagTotal_thenReturnInteger() throws InterruptedException {
assertThat(storage.getLagTotal()).isEqualTo(0);
storage.put("main", mock(TbQueueMsg.class));
assertThat(storage.getLagTotal()).isEqualTo(1);
storage.put("main", mock(TbQueueMsg.class));
assertThat(storage.getLagTotal()).isEqualTo(2);
storage.put("hp", mock(TbQueueMsg.class));
assertThat(storage.getLagTotal()).isEqualTo(3);
storage.get("main");
assertThat(storage.getLagTotal()).isEqualTo(1);
storage.cleanup();
assertThat(storage.getLagTotal()).isEqualTo(0);
}
}
Loading…
Cancel
Save