From a04eac6015bc98e626f446cfa9b6c4be802df626 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 27 Apr 2021 18:26:52 +0300 Subject: [PATCH 01/28] added logs and timing metrics for DefaultTbQueueRequestTemplate, TbKafkaConsumerTemplate (todo revert or refactor) --- .../queue/common/DefaultTbQueueRequestTemplate.java | 8 ++++---- .../server/queue/kafka/TbKafkaConsumerTemplate.java | 10 ++++++++++ 2 files changed, 14 insertions(+), 4 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index b171f2d8d8..e07be52491 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -92,10 +92,10 @@ public class DefaultTbQueueRequestTemplate responses = responseTemplate.poll(pollInterval); - if (responses.size() > 0) { - log.trace("Polling responses completed, consumer records count [{}]", responses.size()); - } + final int pendingRequestsCount = pendingRequests.size(); + log.trace("Starting template pool topic {}, for pendingRequests {}", responseTemplate.getTopic(), pendingRequestsCount); + List responses = responseTemplate.poll(pollInterval); //poll js responses + log.trace("Completed template poll topic {}, for pendingRequests [{}], received [{}]", responseTemplate.getTopic(), pendingRequestsCount, responses.size()); responses.forEach(response -> { byte[] requestIdHeader = response.getHeaders().get(REQUEST_ID_HEADER); UUID requestId; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java index fe4300a421..816a732d41 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerTemplate.java @@ -21,6 +21,7 @@ import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.springframework.util.StopWatch; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueMsg; import org.thingsboard.server.queue.common.AbstractTbQueueConsumerTemplate; @@ -82,7 +83,16 @@ public class TbKafkaConsumerTemplate extends AbstractTbQue @Override protected List> doPoll(long durationInMillis) { + StopWatch stopWatch = new StopWatch(); + stopWatch.start(); + + log.trace("poll topic {} maxDuration {}", getTopic(), durationInMillis); + ConsumerRecords records = consumer.poll(Duration.ofMillis(durationInMillis)); + + stopWatch.stop(); + log.trace("poll topic {} took {}ms", getTopic(), stopWatch.getTotalTimeMillis()); + if (records.isEmpty()) { return Collections.emptyList(); } else { From 67c9025a0625270f02626421d5a47d6c6c79ee3c Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 5 May 2021 18:01:52 +0300 Subject: [PATCH 02/28] improved logs for timeouts for DefaultTbQueueRequestTemplate --- .../server/queue/common/DefaultTbQueueRequestTemplate.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index e07be52491..084b43f4b6 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -121,7 +121,7 @@ public class DefaultTbQueueRequestTemplate staleRequest = pendingRequests.remove(key); if (staleRequest != null) { - log.trace("[{}] Request timeout detected, expTime [{}], tickTs [{}]", key, staleRequest.expTime, tickTs); + log.info("[{}] Request timeout detected, expTime [{}], tickTs [{}]", key, staleRequest.expTime, tickTs); staleRequest.future.setException(new TimeoutException()); } } @@ -129,7 +129,7 @@ public class DefaultTbQueueRequestTemplate Date: Wed, 5 May 2021 20:10:51 +0300 Subject: [PATCH 03/28] refactored for tests DefaultTbQueueRequestTemplate class --- .../common/DefaultTbQueueRequestTemplate.java | 112 ++++++++++-------- 1 file changed, 64 insertions(+), 48 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index 084b43f4b6..2c779d5d5a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -85,59 +85,75 @@ public class DefaultTbQueueRequestTemplate { - long nextCleanupMs = 0L; - while (!stopped) { - try { - final int pendingRequestsCount = pendingRequests.size(); - log.trace("Starting template pool topic {}, for pendingRequests {}", responseTemplate.getTopic(), pendingRequestsCount); - List responses = responseTemplate.poll(pollInterval); //poll js responses - log.trace("Completed template poll topic {}, for pendingRequests [{}], received [{}]", responseTemplate.getTopic(), pendingRequestsCount, responses.size()); - responses.forEach(response -> { - byte[] requestIdHeader = response.getHeaders().get(REQUEST_ID_HEADER); - UUID requestId; - if (requestIdHeader == null) { - log.error("[{}] Missing requestId in header and body", response); - } else { - requestId = bytesToUuid(requestIdHeader); - log.trace("[{}] Response received: {}", requestId, response); - ResponseMetaData expectedResponse = pendingRequests.remove(requestId); - if (expectedResponse == null) { - log.trace("[{}] Invalid or stale request", requestId); - } else { - expectedResponse.future.set(response); + executor.submit(this::fetchAndProcessResponses); + } + + void fetchAndProcessResponses() { + long nextCleanupMs = 0L; + while (!stopped) { + try { + final int pendingRequestsCount = pendingRequests.size(); + log.trace("Starting template pool topic {}, for pendingRequests {}", responseTemplate.getTopic(), pendingRequestsCount); + List responses = doPoll(); //poll js responses + log.trace("Completed template poll topic {}, for pendingRequests [{}], received [{}]", responseTemplate.getTopic(), pendingRequestsCount, responses.size()); + responses.forEach(this::processResponse); + responseTemplate.commit(); + tickTs = System.currentTimeMillis(); + tickSize = pendingRequests.size(); + if (nextCleanupMs < tickTs) { + //cleanup; + pendingRequests.forEach((key, value) -> { + if (value.expTime < tickTs) { + ResponseMetaData staleRequest = pendingRequests.remove(key); + if (staleRequest != null) { + setTimeoutException(key, staleRequest); } } }); - responseTemplate.commit(); - tickTs = System.currentTimeMillis(); - tickSize = pendingRequests.size(); - if (nextCleanupMs < tickTs) { - //cleanup; - pendingRequests.forEach((key, value) -> { - if (value.expTime < tickTs) { - ResponseMetaData staleRequest = pendingRequests.remove(key); - if (staleRequest != null) { - log.info("[{}] Request timeout detected, expTime [{}], tickTs [{}]", key, staleRequest.expTime, tickTs); - staleRequest.future.setException(new TimeoutException()); - } - } - }); - nextCleanupMs = tickTs + maxRequestTimeout; - } - } catch (Throwable e) { - log.warn("Failed to obtain responses from queue. Going to sleep " + pollInterval + "ms", e); - try { - Thread.sleep(pollInterval); - } catch (InterruptedException e2) { - log.trace("Failed to wait until the server has capacity to handle new responses", e2); - } + nextCleanupMs = tickTs + maxRequestTimeout; } + } catch (Throwable e) { + log.warn("Failed to obtain responses from queue. Going to sleep " + pollInterval + "ms", e); + sleep(); } - }); + } + } + + List doPoll() { + return responseTemplate.poll(pollInterval); + } + + void sleep() { + try { + Thread.sleep(pollInterval); + } catch (InterruptedException e2) { + log.trace("Failed to wait until the server has capacity to handle new responses", e2); + } + } + + void setTimeoutException(UUID key, ResponseMetaData staleRequest) { + log.info("[{}] Request timeout detected, expTime [{}], tickTs [{}]", key, staleRequest.expTime, tickTs); + staleRequest.future.setException(new TimeoutException()); + } + + void processResponse(Response response) { + byte[] requestIdHeader = response.getHeaders().get(REQUEST_ID_HEADER); + UUID requestId; + if (requestIdHeader == null) { + log.error("[{}] Missing requestId in header and body", response); + } else { + requestId = bytesToUuid(requestIdHeader); + log.trace("[{}] Response received: {}", requestId, response); + ResponseMetaData expectedResponse = pendingRequests.remove(requestId); + if (expectedResponse == null) { + log.warn("[{}] Invalid or stale request, response: {}", requestId, response); + } else { + expectedResponse.future.set(response); + } + } } @Override @@ -174,7 +190,7 @@ public class DefaultTbQueueRequestTemplate future = SettableFuture.create(); ResponseMetaData responseMetaData = new ResponseMetaData<>(tickTs + maxRequestTimeout, future); pendingRequests.putIfAbsent(requestId, responseMetaData); - log.trace("[{}] Sending request, key [{}], expTime [{}]", requestId, request.getKey(), responseMetaData.expTime); + log.trace("[{}] Sending request, key [{}], expTime [{}], request {}", requestId, request.getKey(), responseMetaData.expTime, request); if (messagesStats != null) { messagesStats.incrementTotal(); } @@ -184,7 +200,7 @@ public class DefaultTbQueueRequestTemplate Date: Thu, 6 May 2021 11:48:42 +0300 Subject: [PATCH 04/28] test: added init stop test for DefaultTbQueueRequestTemplate --- .../common/DefaultTbQueueRequestTemplate.java | 22 ++-- .../DefaultTbQueueRequestTemplateTest.java | 113 ++++++++++++++++++ 2 files changed, 123 insertions(+), 12 deletions(-) create mode 100644 common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index 2c779d5d5a..a51e143e26 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -31,6 +31,7 @@ import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.TbQueueRequestTemplate; import org.thingsboard.server.common.stats.MessagesStats; +import javax.annotation.Nullable; import java.util.List; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; @@ -47,14 +48,14 @@ public class DefaultTbQueueRequestTemplate requestTemplate; private final TbQueueConsumer responseTemplate; private final ConcurrentMap> pendingRequests; - private final boolean internalExecutor; + final boolean internalExecutor; private final ExecutorService executor; private final long maxRequestTimeout; private final long maxPendingRequests; private final long pollInterval; - private volatile long tickTs = 0L; - private volatile long tickSize = 0L; - private volatile boolean stopped = false; + volatile long tickTs = 0L; + volatile long tickSize = 0L; + volatile boolean stopped = false; private MessagesStats messagesStats; @@ -65,7 +66,7 @@ public class DefaultTbQueueRequestTemplate requestTemplate; + @Mock + TbQueueConsumer responseTemplate; + + long maxRequestTimeout = 20; + long maxPendingRequests = 10000; + long pollInterval = 25; + String topic = "js-responses-tb-node-0"; + + DefaultTbQueueRequestTemplate inst; + + @Before + public void setUp() throws Exception { + willReturn(topic).given(responseTemplate).getTopic(); + inst = new DefaultTbQueueRequestTemplate( + queueAdmin, requestTemplate, responseTemplate, + maxRequestTimeout, maxPendingRequests, pollInterval, executor); + + } + + @Test + public void givenExternalExecutor_whenStop_thenDoNotShutdownExecutor() { + inst = Mockito.spy(inst); + willDoNothing().given(inst).fetchAndProcessResponses(); + Assert.assertFalse(inst.stopped); + + inst.init(); + Assert.assertNotEquals(inst.tickTs, 0); + Assert.assertFalse(inst.internalExecutor); + verify(queueAdmin, times(1)).createTopicIfNotExists(topic); + verify(requestTemplate, times(1)).init(); + verify(responseTemplate, times(1)).subscribe(); + verify(executor, times(1)).submit(any(Runnable.class)); + + inst.stop(); + Assert.assertTrue(inst.stopped); + verify(responseTemplate, times(1)).unsubscribe(); + verify(requestTemplate, times(1)).stop(); + verify(executor, never()).shutdownNow(); + } + + @Test + public void fetchAndProcessResponses() { + inst.init(); + + inst.stop(); + } +} \ No newline at end of file From 872717828cc0e9338f6d1e986abcb67a56e0d06f Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 6 May 2021 14:54:27 +0300 Subject: [PATCH 05/28] test: refactored and added mainLoop test DefaultTbQueueRequestTemplate --- .../DefaultTbQueueRequestTemplateTest.java | 76 ++++++++++++++----- 1 file changed, 56 insertions(+), 20 deletions(-) diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java index 769b66fcd4..413b90c5b9 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java +++ b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java @@ -31,27 +31,27 @@ package org.thingsboard.server.queue.common; import org.junit.After; -import org.junit.Assert; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; -import org.mockito.BDDMockito; import org.mockito.Mock; -import org.mockito.Mockito; -import org.mockito.internal.util.reflection.Whitebox; import org.mockito.runners.MockitoJUnitRunner; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueMsg; import org.thingsboard.server.queue.TbQueueProducer; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; +import java.util.concurrent.TimeUnit; import static org.junit.Assert.*; +import static org.mockito.BDDMockito.willAnswer; import static org.mockito.BDDMockito.willDoNothing; import static org.mockito.BDDMockito.willReturn; import static org.mockito.Matchers.any; import static org.mockito.Mockito.never; +import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -61,16 +61,17 @@ public class DefaultTbQueueRequestTemplateTest { @Mock TbQueueAdmin queueAdmin; @Mock - ExecutorService executor; - @Mock TbQueueProducer requestTemplate; @Mock TbQueueConsumer responseTemplate; + @Mock + ExecutorService executorMock; + ExecutorService executor; + String topic = "js-responses-tb-node-0"; long maxRequestTimeout = 20; long maxPendingRequests = 10000; long pollInterval = 25; - String topic = "js-responses-tb-node-0"; DefaultTbQueueRequestTemplate inst; @@ -79,35 +80,70 @@ public class DefaultTbQueueRequestTemplateTest { willReturn(topic).given(responseTemplate).getTopic(); inst = new DefaultTbQueueRequestTemplate( queueAdmin, requestTemplate, responseTemplate, - maxRequestTimeout, maxPendingRequests, pollInterval, executor); + maxRequestTimeout, maxPendingRequests, pollInterval, executorMock); + + } + @After + public void tearDown() throws Exception { + if (executor != null) { + executor.shutdownNow(); + } } @Test - public void givenExternalExecutor_whenStop_thenDoNotShutdownExecutor() { - inst = Mockito.spy(inst); - willDoNothing().given(inst).fetchAndProcessResponses(); - Assert.assertFalse(inst.stopped); + public void givenInstance_whenVerifyInitialParameters_thenOK() { + assertEquals(maxPendingRequests, inst.maxPendingRequests); + assertEquals(maxRequestTimeout, inst.maxRequestTimeout); + assertEquals(pollInterval, inst.pollInterval); + assertEquals(executorMock, inst.executor); + assertFalse(inst.stopped); + assertFalse(inst.internalExecutor); + } + + @Test + public void givenExternalExecutor_whenInitStop_thenOK() { + inst = spy(inst); + willDoNothing().given(inst).mainLoop(); inst.init(); - Assert.assertNotEquals(inst.tickTs, 0); - Assert.assertFalse(inst.internalExecutor); + assertNotEquals(0, inst.tickTs); + assertEquals(0, inst.nextCleanupMs); verify(queueAdmin, times(1)).createTopicIfNotExists(topic); verify(requestTemplate, times(1)).init(); verify(responseTemplate, times(1)).subscribe(); - verify(executor, times(1)).submit(any(Runnable.class)); + verify(executorMock, times(1)).submit(any(Runnable.class)); inst.stop(); - Assert.assertTrue(inst.stopped); + assertTrue(inst.stopped); verify(responseTemplate, times(1)).unsubscribe(); verify(requestTemplate, times(1)).stop(); - verify(executor, never()).shutdownNow(); + verify(executorMock, never()).shutdownNow(); } @Test - public void fetchAndProcessResponses() { - inst.init(); + public void givenMainLoop_whenLoopFewTimes_thenVerifyInvocationCount() throws InterruptedException { + inst = spy(inst); + executor = inst.createExecutor(); + CountDownLatch latch = new CountDownLatch(5); + willDoNothing().given(inst).sleep(); + willAnswer(invocation -> { + if (latch.getCount() == 1) { + inst.stop(); //stop the loop in natural way + } + if (latch.getCount() == 3 || latch.getCount() == 4) { + latch.countDown(); + throw new RuntimeException("test catch block"); + } + latch.countDown(); + return null; + }).given(inst).fetchAndProcessResponses(); - inst.stop(); + executor.submit(inst::mainLoop); + latch.await(10, TimeUnit.SECONDS); + + verify(inst, times(5)).fetchAndProcessResponses(); + verify(inst, times(2)).sleep(); } + } \ No newline at end of file From d5fffa5002528db63302749b354aff8cd488ce5d Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 6 May 2021 15:50:07 +0300 Subject: [PATCH 06/28] test: added failed test case with overflow maxPendingRequests for DefaultTbQueueRequestTemplate.send(). another refactoring for easy mocking. --- .../DefaultTbQueueRequestTemplateTest.java | 53 ++++++++++++++++--- 1 file changed, 47 insertions(+), 6 deletions(-) diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java index 413b90c5b9..ef9d0835d2 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java +++ b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java @@ -39,17 +39,22 @@ import org.mockito.runners.MockitoJUnitRunner; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueMsg; +import org.thingsboard.server.queue.TbQueueMsgHeaders; import org.thingsboard.server.queue.TbQueueProducer; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.TimeUnit; -import static org.junit.Assert.*; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertNotEquals; +import static org.junit.Assert.assertTrue; import static org.mockito.BDDMockito.willAnswer; import static org.mockito.BDDMockito.willDoNothing; import static org.mockito.BDDMockito.willReturn; import static org.mockito.Matchers.any; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; @@ -70,7 +75,7 @@ public class DefaultTbQueueRequestTemplateTest { ExecutorService executor; String topic = "js-responses-tb-node-0"; long maxRequestTimeout = 20; - long maxPendingRequests = 10000; + long maxPendingRequests = 32; long pollInterval = 25; DefaultTbQueueRequestTemplate inst; @@ -78,9 +83,9 @@ public class DefaultTbQueueRequestTemplateTest { @Before public void setUp() throws Exception { willReturn(topic).given(responseTemplate).getTopic(); - inst = new DefaultTbQueueRequestTemplate( + inst = spy(new DefaultTbQueueRequestTemplate( queueAdmin, requestTemplate, responseTemplate, - maxRequestTimeout, maxPendingRequests, pollInterval, executorMock); + maxRequestTimeout, maxPendingRequests, pollInterval, executorMock)); } @@ -103,7 +108,6 @@ public class DefaultTbQueueRequestTemplateTest { @Test public void givenExternalExecutor_whenInitStop_thenOK() { - inst = spy(inst); willDoNothing().given(inst).mainLoop(); inst.init(); @@ -123,7 +127,6 @@ public class DefaultTbQueueRequestTemplateTest { @Test public void givenMainLoop_whenLoopFewTimes_thenVerifyInvocationCount() throws InterruptedException { - inst = spy(inst); executor = inst.createExecutor(); CountDownLatch latch = new CountDownLatch(5); willDoNothing().given(inst).sleep(); @@ -146,4 +149,42 @@ public class DefaultTbQueueRequestTemplateTest { verify(inst, times(2)).sleep(); } + @Test + public void givenMessages_whenSend_thenOK() { + willDoNothing().given(inst).sendToRequestTemplate(any(), any(), any(), any()); + inst.init(); + int msgCount = 10; + for (int i = 0; i < msgCount; i++) { + inst.send(getRequestMsgMock()); + } + assertEquals(msgCount, inst.pendingRequests.size()); + verify(inst, times(msgCount)).sendToRequestTemplate(any(), any(), any(), any()); + } + + @Test + public void givenMessagesOverMaxPendingRequests_whenSend_thenImmediateFailedFutureForTheOfRequests() { + willDoNothing().given(inst).sendToRequestTemplate(any(), any(), any(), any()); + inst.init(); + assertEquals(0, inst.tickSize); + int msgOverflowCount = 10; + for (int i = 0; i < inst.maxPendingRequests; i++) { + assertFalse(inst.send(getRequestMsgMock()).isDone()); //SettableFuture future - pending only + } + for (int i = 0; i < msgOverflowCount; i++) { + assertTrue("max pending requests overflow", inst.send(getRequestMsgMock()).isDone()); //overflow, immediate failed future + } + assertEquals(inst.maxPendingRequests, inst.pendingRequests.size()); + verify(inst, times((int) inst.maxPendingRequests)).sendToRequestTemplate(any(), any(), any(), any()); + } + + @Test + public void givenNothing_whenFetchAndProcessResponsesWithTimeout_thenFail() { + + } + + TbQueueMsg getRequestMsgMock() { + TbQueueMsg requestMsg = mock(TbQueueMsg.class); + willReturn(mock(TbQueueMsgHeaders.class)).given(requestMsg).getHeaders(); + return requestMsg; + } } \ No newline at end of file From 9666986156a0005f2d934646d625b29c02aa091f Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 7 May 2021 10:04:41 +0300 Subject: [PATCH 07/28] test: added failed test for FetchAndProcessResponses when request removed as staled too early DefaultTbQueueRequestTemplate --- .../common/DefaultTbQueueRequestTemplate.java | 131 ++++++++++++------ .../DefaultTbQueueRequestTemplateTest.java | 60 ++++++-- 2 files changed, 138 insertions(+), 53 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index a51e143e26..bfda553538 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -19,6 +19,7 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.SettableFuture; import lombok.Builder; +import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -47,15 +48,16 @@ public class DefaultTbQueueRequestTemplate requestTemplate; private final TbQueueConsumer responseTemplate; - private final ConcurrentMap> pendingRequests; + final ConcurrentMap> pendingRequests; final boolean internalExecutor; - private final ExecutorService executor; - private final long maxRequestTimeout; - private final long maxPendingRequests; - private final long pollInterval; + final ExecutorService executor; + final long maxRequestTimeout; + final long maxPendingRequests; + final long pollInterval; volatile long tickTs = 0L; volatile long tickSize = 0L; volatile boolean stopped = false; + long nextCleanupMs = 0L; private MessagesStats messagesStats; @@ -75,44 +77,26 @@ public class DefaultTbQueueRequestTemplate responses = doPoll(); //poll js responses - log.trace("Completed template poll topic {}, for pendingRequests [{}], received [{}]", responseTemplate.getTopic(), pendingRequestsCount, responses.size()); - responses.forEach(this::processResponse); - responseTemplate.commit(); - tickTs = System.currentTimeMillis(); - tickSize = pendingRequests.size(); - if (nextCleanupMs < tickTs) { - //cleanup; - pendingRequests.forEach((key, value) -> { - if (value.expTime < tickTs) { - ResponseMetaData staleRequest = pendingRequests.remove(key); - if (staleRequest != null) { - setTimeoutException(key, staleRequest); - } - } - }); - nextCleanupMs = tickTs + maxRequestTimeout; - } + fetchAndProcessResponses(); } catch (Throwable e) { log.warn("Failed to obtain responses from queue. Going to sleep " + pollInterval + "ms", e); sleep(); @@ -120,6 +104,36 @@ public class DefaultTbQueueRequestTemplate responses = doPoll(); //poll js responses + //if (responses.size() > 0) { + log.trace("Completed template poll topic {}, for pendingRequests [{}], received [{}]", responseTemplate.getTopic(), pendingRequestsCount, responses.size()); + //} + responses.forEach(this::processResponse); //this can take a long time + responseTemplate.commit(); + tickTs = getCurrentTime(); + tickSize = pendingRequests.size(); + if (nextCleanupMs < tickTs) { + //cleanup; + pendingRequests.forEach((key, value) -> { + if (value.expTime < tickTs) { + ResponseMetaData staleRequest = pendingRequests.remove(key); + if (staleRequest != null) { + setTimeoutException(key, staleRequest, tickTs); + } + } + }); + setupNextCleanup(); + } + } + + void setupNextCleanup() { + nextCleanupMs = tickTs + maxRequestTimeout; + log.info("setupNextCleanup {}", nextCleanupMs); + } + List doPoll() { return responseTemplate.poll(pollInterval); } @@ -132,8 +146,13 @@ public class DefaultTbQueueRequestTemplate staleRequest) { - log.info("[{}] Request timeout detected, expTime [{}], tickTs [{}]", key, staleRequest.expTime, tickTs); + void setTimeoutException(UUID key, ResponseMetaData staleRequest, long tickTs) { + if (tickTs >= staleRequest.getSubmitTime() + staleRequest.getTimeout()) { + log.info("Request timeout detected, tickTs [{}], {}, key [{}]", tickTs, staleRequest, key); + } else { + log.error("Request timeout detected, tickTs [{}], {}, key [{}]", tickTs, staleRequest, key); + } + staleRequest.future.setException(new TimeoutException()); } @@ -144,10 +163,10 @@ public class DefaultTbQueueRequestTemplate expectedResponse = pendingRequests.remove(requestId); if (expectedResponse == null) { - log.warn("[{}] Invalid or stale request, response: {}", requestId, response); + log.warn("[{}] Invalid or stale request, response: {}", requestId, String.valueOf(response).replace("\n", " ")); } else { expectedResponse.future.set(response); } @@ -184,11 +203,22 @@ public class DefaultTbQueueRequestTemplate future = SettableFuture.create(); - ResponseMetaData responseMetaData = new ResponseMetaData<>(tickTs + maxRequestTimeout, future); + ResponseMetaData responseMetaData = new ResponseMetaData<>(tickTs + maxRequestTimeout, future, currentTime, maxRequestTimeout); + log.info("pending {}", responseMetaData); pendingRequests.putIfAbsent(requestId, responseMetaData); - log.trace("[{}] Sending request, key [{}], expTime [{}], request {}", requestId, request.getKey(), responseMetaData.expTime, request); + sendToRequestTemplate(request, requestId, future, responseMetaData); + return future; + } + + long getCurrentTime() { + return System.currentTimeMillis(); + } + + void sendToRequestTemplate(Request request, UUID requestId, SettableFuture future, ResponseMetaData responseMetaData) { + log.trace("[{}] Sending request, key [{}], expTime [{}], request {}", requestId, request.getKey(), responseMetaData.expTime, String.valueOf(request).replace("\n", " ")); if (messagesStats != null) { messagesStats.incrementTotal(); } @@ -198,7 +228,7 @@ public class DefaultTbQueueRequestTemplate { + @Getter + static class ResponseMetaData { + private final long submitTime; + private final long timeout; private final long expTime; private final SettableFuture future; - ResponseMetaData(long ts, SettableFuture future) { + ResponseMetaData(long ts, SettableFuture future, long submitTime, long timeout) { + this.submitTime = submitTime; + this.timeout = timeout; this.expTime = ts; this.future = future; } + + @Override + public String toString() { + return "ResponseMetaData{" + + "submitTime=" + submitTime + + ", calculatedExpTime=" + (submitTime + timeout) + + ", expTime=" + expTime + + ", deltaMs=" + (expTime - submitTime) + + ", future=" + future + + '}'; + } } } diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java index ef9d0835d2..715b4eab21 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java +++ b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java @@ -29,23 +29,29 @@ * OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART. */ package org.thingsboard.server.queue.common; - +import lombok.extern.slf4j.Slf4j; import org.junit.After; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.runners.MockitoJUnitRunner; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueMsg; -import org.thingsboard.server.queue.TbQueueMsgHeaders; import org.thingsboard.server.queue.TbQueueProducer; +import java.util.Collections; +import java.util.List; +import java.util.UUID; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.greaterThanOrEqualTo; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotEquals; @@ -54,12 +60,17 @@ import static org.mockito.BDDMockito.willAnswer; import static org.mockito.BDDMockito.willDoNothing; import static org.mockito.BDDMockito.willReturn; import static org.mockito.Matchers.any; +import static org.mockito.Mockito.RETURNS_DEEP_STUBS; +import static org.mockito.Mockito.atLeastOnce; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; +import static org.hamcrest.MatcherAssert.assertThat; + +@Slf4j @RunWith(MockitoJUnitRunner.class) public class DefaultTbQueueRequestTemplateTest { @@ -74,9 +85,9 @@ public class DefaultTbQueueRequestTemplateTest { ExecutorService executor; String topic = "js-responses-tb-node-0"; - long maxRequestTimeout = 20; + long maxRequestTimeout = 10; long maxPendingRequests = 32; - long pollInterval = 25; + long pollInterval = 5; DefaultTbQueueRequestTemplate inst; @@ -171,20 +182,49 @@ public class DefaultTbQueueRequestTemplateTest { assertFalse(inst.send(getRequestMsgMock()).isDone()); //SettableFuture future - pending only } for (int i = 0; i < msgOverflowCount; i++) { - assertTrue("max pending requests overflow", inst.send(getRequestMsgMock()).isDone()); //overflow, immediate failed future + assertFalse("max pending requests overflow", inst.send(getRequestMsgMock()).isDone()); //overflow, immediate failed future } - assertEquals(inst.maxPendingRequests, inst.pendingRequests.size()); + assertThat(inst.pendingRequests.size(), equalTo(inst.maxPendingRequests)); verify(inst, times((int) inst.maxPendingRequests)).sendToRequestTemplate(any(), any(), any(), any()); } @Test - public void givenNothing_whenFetchAndProcessResponsesWithTimeout_thenFail() { + public void givenNothing_whenSendAndFetchAndProcessResponsesWithTimeout_thenFail() { + //given + AtomicLong currentTime = new AtomicLong(); + willAnswer(x -> { + log.info("currentTime={}", currentTime.get()); + return currentTime.get(); + }).given(inst).getCurrentTime(); + inst.init(); + inst.setupNextCleanup(); + willReturn(Collections.emptyList()).given(inst).doPoll(); + willDoNothing().given(inst).processResponse(any()); + + //when + for (int i = 0; i <= inst.maxRequestTimeout*2; i++) { + currentTime.incrementAndGet(); + assertFalse(inst.send(getRequestMsgMock()).isDone()); //SettableFuture future - pending only + if (i % (inst.maxRequestTimeout * 3 / 2) == 0) { + inst.fetchAndProcessResponses(); + } + } + + //then + ArgumentCaptor argumentCaptorResp = ArgumentCaptor.forClass(DefaultTbQueueRequestTemplate.ResponseMetaData.class); + ArgumentCaptor argumentCaptorUUID = ArgumentCaptor.forClass(UUID.class); + ArgumentCaptor argumentCaptorLong = ArgumentCaptor.forClass(Long.class); + verify(inst, atLeastOnce()).setTimeoutException(argumentCaptorUUID.capture(), argumentCaptorResp.capture(), argumentCaptorLong.capture()); + + List responseMetaDataList = argumentCaptorResp.getAllValues(); + List tickTsList = argumentCaptorLong.getAllValues(); + for (int i = 0; i < responseMetaDataList.size(); i++) { + assertThat("tickTs >= calculatedExpTime", tickTsList.get(i), greaterThanOrEqualTo(responseMetaDataList.get(i).getSubmitTime() + responseMetaDataList.get(i).getTimeout())); + } } TbQueueMsg getRequestMsgMock() { - TbQueueMsg requestMsg = mock(TbQueueMsg.class); - willReturn(mock(TbQueueMsgHeaders.class)).given(requestMsg).getHeaders(); - return requestMsg; + return mock(TbQueueMsg.class, RETURNS_DEEP_STUBS); } } \ No newline at end of file From 28235732c6140dfc9a5cb81d4c2e39b8ac217532 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 7 May 2021 12:07:17 +0300 Subject: [PATCH 08/28] test: fixed class DefaultTbQueueRequestTemplate --- .../common/DefaultTbQueueRequestTemplate.java | 90 +++++++++++-------- .../DefaultTbQueueRequestTemplateTest.java | 37 ++++---- 2 files changed, 73 insertions(+), 54 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index bfda553538..d5179385e6 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -36,10 +36,12 @@ import javax.annotation.Nullable; import java.util.List; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; @Slf4j public class DefaultTbQueueRequestTemplate extends AbstractTbQueueTemplate @@ -48,16 +50,15 @@ public class DefaultTbQueueRequestTemplate requestTemplate; private final TbQueueConsumer responseTemplate; - final ConcurrentMap> pendingRequests; + final ConcurrentHashMap> pendingRequests = new ConcurrentHashMap<>(); final boolean internalExecutor; final ExecutorService executor; - final long maxRequestTimeout; + final long maxRequestTimeoutNs; final long maxPendingRequests; final long pollInterval; - volatile long tickTs = 0L; - volatile long tickSize = 0L; volatile boolean stopped = false; - long nextCleanupMs = 0L; + long nextCleanupNs = 0L; + private final Lock cleanerLock = new ReentrantLock(); private MessagesStats messagesStats; @@ -72,8 +73,7 @@ public class DefaultTbQueueRequestTemplate(); - this.maxRequestTimeout = maxRequestTimeout; + this.maxRequestTimeoutNs = TimeUnit.MILLISECONDS.toNanos(maxRequestTimeout); this.maxPendingRequests = maxPendingRequests; this.pollInterval = pollInterval; this.internalExecutor = (executor == null); @@ -88,7 +88,6 @@ public class DefaultTbQueueRequestTemplate responses = doPoll(); //poll js responses //if (responses.size() > 0) { @@ -113,25 +112,36 @@ public class DefaultTbQueueRequestTemplate { - if (value.expTime < tickTs) { - ResponseMetaData staleRequest = pendingRequests.remove(key); - if (staleRequest != null) { - setTimeoutException(key, staleRequest, tickTs); + tryCleanStaleRequests(); + } + + private boolean tryCleanStaleRequests() { + if (!cleanerLock.tryLock()) { + return false; + } + try { + log.trace("tryCleanStaleRequest..."); + final long currentNs = getCurrentClockNs(); + if (nextCleanupNs < currentNs) { + pendingRequests.forEach((key, value) -> { + if (value.expTime < currentNs) { + ResponseMetaData staleRequest = pendingRequests.remove(key); + if (staleRequest != null) { + setTimeoutException(key, staleRequest, currentNs); + } } - } - }); - setupNextCleanup(); + }); + setupNextCleanup(); + } + } finally { + cleanerLock.unlock(); } + return true; } void setupNextCleanup() { - nextCleanupMs = tickTs + maxRequestTimeout; - log.info("setupNextCleanup {}", nextCleanupMs); + nextCleanupNs = getCurrentClockNs() + maxRequestTimeoutNs; + log.info("setupNextCleanup {}", nextCleanupNs); } List doPoll() { @@ -146,11 +156,11 @@ public class DefaultTbQueueRequestTemplate staleRequest, long tickTs) { - if (tickTs >= staleRequest.getSubmitTime() + staleRequest.getTimeout()) { - log.info("Request timeout detected, tickTs [{}], {}, key [{}]", tickTs, staleRequest, key); + void setTimeoutException(UUID key, ResponseMetaData staleRequest, long currentNs) { + if (currentNs >= staleRequest.getSubmitTime() + staleRequest.getTimeout()) { + log.info("Request timeout detected, currentNs [{}], {}, key [{}]", currentNs, staleRequest, key); } else { - log.error("Request timeout detected, tickTs [{}], {}, key [{}]", tickTs, staleRequest, key); + log.error("Request timeout detected, currentNs [{}], {}, key [{}]", currentNs, staleRequest, key); } staleRequest.future.setException(new TimeoutException()); @@ -197,23 +207,31 @@ public class DefaultTbQueueRequestTemplate send(Request request) { - if (tickSize > maxPendingRequests) { + if (pendingRequests.mappingCount() >= maxPendingRequests) { + log.warn("Pending request map is full [{}]! Consider to increase maxPendingRequests or increase processing performance", maxPendingRequests); return Futures.immediateFailedFuture(new RuntimeException("Pending request map is full!")); } UUID requestId = UUID.randomUUID(); request.getHeaders().put(REQUEST_ID_HEADER, uuidToBytes(requestId)); request.getHeaders().put(RESPONSE_TOPIC_HEADER, stringToBytes(responseTemplate.getTopic())); - long currentTime = getCurrentTime(); - request.getHeaders().put(REQUEST_TIME, longToBytes(currentTime)); + request.getHeaders().put(REQUEST_TIME, longToBytes(getCurrentTimeMs())); + long currentClockNs = getCurrentClockNs(); SettableFuture future = SettableFuture.create(); - ResponseMetaData responseMetaData = new ResponseMetaData<>(tickTs + maxRequestTimeout, future, currentTime, maxRequestTimeout); - log.info("pending {}", responseMetaData); - pendingRequests.putIfAbsent(requestId, responseMetaData); + ResponseMetaData responseMetaData = new ResponseMetaData<>(currentClockNs + maxRequestTimeoutNs, future, currentClockNs, maxRequestTimeoutNs); + log.info("pending {}", responseMetaData); //TODO trace + if (pendingRequests.putIfAbsent(requestId, responseMetaData) != null) { + log.warn("Pending request already exists [{}]!", maxPendingRequests); + return Futures.immediateFailedFuture(new RuntimeException("Pending request already exists !" + requestId)); + } sendToRequestTemplate(request, requestId, future, responseMetaData); return future; } - long getCurrentTime() { + long getCurrentClockNs() { + return System.nanoTime(); //MONOTONIC clock instead wall clock + } + + long getCurrentTimeMs() { //Wall clock to send Ts to the an external service return System.currentTimeMillis(); } @@ -261,8 +279,8 @@ public class DefaultTbQueueRequestTemplate { log.info("currentTime={}", currentTime.get()); return currentTime.get(); - }).given(inst).getCurrentTime(); + }).given(inst).getCurrentClockNs(); inst.init(); inst.setupNextCleanup(); willReturn(Collections.emptyList()).given(inst).doPoll(); willDoNothing().given(inst).processResponse(any()); //when - for (int i = 0; i <= inst.maxRequestTimeout*2; i++) { - currentTime.incrementAndGet(); + long stepNs = TimeUnit.MILLISECONDS.toNanos(1); + for (long i = 0; i <= inst.maxRequestTimeoutNs * 2; i = i + stepNs) { + currentTime.addAndGet(stepNs); assertFalse(inst.send(getRequestMsgMock()).isDone()); //SettableFuture future - pending only - if (i % (inst.maxRequestTimeout * 3 / 2) == 0) { + if (i % (inst.maxRequestTimeoutNs * 3 / 2) == 0) { inst.fetchAndProcessResponses(); } } From 928b8f0fd9ca39fe7ebe874c7ea195153779b96d Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 7 May 2021 12:26:44 +0300 Subject: [PATCH 09/28] test: refactored for assertThat for DefaultTbQueueRequestTemplateTest --- .../common/DefaultTbQueueRequestTemplate.java | 5 ++--- .../DefaultTbQueueRequestTemplateTest.java | 17 ++++++++--------- 2 files changed, 10 insertions(+), 12 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index d5179385e6..c31dc3a899 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -158,11 +158,10 @@ public class DefaultTbQueueRequestTemplate staleRequest, long currentNs) { if (currentNs >= staleRequest.getSubmitTime() + staleRequest.getTimeout()) { - log.info("Request timeout detected, currentNs [{}], {}, key [{}]", currentNs, staleRequest, key); + log.warn("Request timeout detected, currentNs [{}], {}, key [{}]", currentNs, staleRequest, key); } else { log.error("Request timeout detected, currentNs [{}], {}, key [{}]", currentNs, staleRequest, key); } - staleRequest.future.setException(new TimeoutException()); } @@ -218,7 +217,7 @@ public class DefaultTbQueueRequestTemplate future = SettableFuture.create(); ResponseMetaData responseMetaData = new ResponseMetaData<>(currentClockNs + maxRequestTimeoutNs, future, currentClockNs, maxRequestTimeoutNs); - log.info("pending {}", responseMetaData); //TODO trace + log.trace("pending {}", responseMetaData); if (pendingRequests.putIfAbsent(requestId, responseMetaData) != null) { log.warn("Pending request already exists [{}]!", maxPendingRequests); return Futures.immediateFailedFuture(new RuntimeException("Pending request already exists !" + requestId)); diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java index 33dfd82989..261ff9c7be 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java +++ b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java @@ -87,7 +87,7 @@ public class DefaultTbQueueRequestTemplateTest { ExecutorService executor; String topic = "js-responses-tb-node-0"; long maxRequestTimeout = 10; - long maxPendingRequests = 1000; + long maxPendingRequests = 32; long pollInterval = 5; DefaultTbQueueRequestTemplate inst; @@ -124,14 +124,14 @@ public class DefaultTbQueueRequestTemplateTest { inst.init(); //assertNotEquals(0, inst.tickTs); - assertEquals(0, inst.nextCleanupNs); + assertThat(inst.nextCleanupNs, equalTo(0L)); verify(queueAdmin, times(1)).createTopicIfNotExists(topic); verify(requestTemplate, times(1)).init(); verify(responseTemplate, times(1)).subscribe(); verify(executorMock, times(1)).submit(any(Runnable.class)); inst.stop(); - assertTrue(inst.stopped); + assertThat(inst.stopped, is(true)); verify(responseTemplate, times(1)).unsubscribe(); verify(requestTemplate, times(1)).stop(); verify(executorMock, never()).shutdownNow(); @@ -165,11 +165,11 @@ public class DefaultTbQueueRequestTemplateTest { public void givenMessages_whenSend_thenOK() { willDoNothing().given(inst).sendToRequestTemplate(any(), any(), any(), any()); inst.init(); - int msgCount = 10; + final int msgCount = 10; for (int i = 0; i < msgCount; i++) { inst.send(getRequestMsgMock()); } - assertEquals(msgCount, inst.pendingRequests.mappingCount()); + assertThat(inst.pendingRequests.mappingCount(), equalTo((long) msgCount)); verify(inst, times(msgCount)).sendToRequestTemplate(any(), any(), any(), any()); } @@ -179,10 +179,10 @@ public class DefaultTbQueueRequestTemplateTest { inst.init(); int msgOverflowCount = 10; for (int i = 0; i < inst.maxPendingRequests; i++) { - assertFalse(inst.send(getRequestMsgMock()).isDone()); //SettableFuture future - pending only + assertThat(inst.send(getRequestMsgMock()).isDone(), is(false)); //SettableFuture future - pending only } for (int i = 0; i < msgOverflowCount; i++) { - assertTrue("max pending requests overflow", inst.send(getRequestMsgMock()).isDone()); //overflow, immediate failed future + assertThat("max pending requests overflow", inst.send(getRequestMsgMock()).isDone(), is(true)); //overflow, immediate failed future } assertThat(inst.pendingRequests.mappingCount(), equalTo(inst.maxPendingRequests)); verify(inst, times((int) inst.maxPendingRequests)).sendToRequestTemplate(any(), any(), any(), any()); @@ -205,7 +205,7 @@ public class DefaultTbQueueRequestTemplateTest { long stepNs = TimeUnit.MILLISECONDS.toNanos(1); for (long i = 0; i <= inst.maxRequestTimeoutNs * 2; i = i + stepNs) { currentTime.addAndGet(stepNs); - assertFalse(inst.send(getRequestMsgMock()).isDone()); //SettableFuture future - pending only + assertThat(inst.send(getRequestMsgMock()).isDone(), is(false)); //SettableFuture future - pending only if (i % (inst.maxRequestTimeoutNs * 3 / 2) == 0) { inst.fetchAndProcessResponses(); } @@ -222,7 +222,6 @@ public class DefaultTbQueueRequestTemplateTest { for (int i = 0; i < responseMetaDataList.size(); i++) { assertThat("tickTs >= calculatedExpTime", tickTsList.get(i), greaterThanOrEqualTo(responseMetaDataList.get(i).getSubmitTime() + responseMetaDataList.get(i).getTimeout())); } - } TbQueueMsg getRequestMsgMock() { From 1e066f21569aa7824177a6080577ac83d29ef2a7 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 7 May 2021 12:41:50 +0300 Subject: [PATCH 10/28] added timeout parameter to the TbQueueRequestTemplate.send( ) --- .../server/queue/TbQueueRequestTemplate.java | 2 ++ .../common/DefaultTbQueueRequestTemplate.java | 18 +++++++++++++++--- 2 files changed, 17 insertions(+), 3 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java index 5dc89a9c26..192f8e1675 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/TbQueueRequestTemplate.java @@ -24,6 +24,8 @@ public interface TbQueueRequestTemplate send(Request request); + ListenableFuture send(Request request, long timeoutNs); + void stop(); void setMessagesStats(MessagesStats messagesStats); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index c31dc3a899..5b03e1bc2d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -206,6 +206,11 @@ public class DefaultTbQueueRequestTemplate send(Request request) { + return send(request, this.maxRequestTimeoutNs); + } + + @Override + public ListenableFuture send(Request request, long requestTimeoutNs) { if (pendingRequests.mappingCount() >= maxPendingRequests) { log.warn("Pending request map is full [{}]! Consider to increase maxPendingRequests or increase processing performance", maxPendingRequests); return Futures.immediateFailedFuture(new RuntimeException("Pending request map is full!")); @@ -216,7 +221,7 @@ public class DefaultTbQueueRequestTemplate future = SettableFuture.create(); - ResponseMetaData responseMetaData = new ResponseMetaData<>(currentClockNs + maxRequestTimeoutNs, future, currentClockNs, maxRequestTimeoutNs); + ResponseMetaData responseMetaData = new ResponseMetaData<>(currentClockNs + requestTimeoutNs, future, currentClockNs, requestTimeoutNs); log.trace("pending {}", responseMetaData); if (pendingRequests.putIfAbsent(requestId, responseMetaData) != null) { log.warn("Pending request already exists [{}]!", maxPendingRequests); @@ -226,11 +231,18 @@ public class DefaultTbQueueRequestTemplate Date: Fri, 7 May 2021 14:07:20 +0300 Subject: [PATCH 11/28] sleep(poolInterval) replaced with parkNanos(1) for DefaultTbQueueRequestTemplate --- .../server/actors/service/DefaultActorService.java | 10 +--------- .../queue/common/DefaultTbQueueRequestTemplate.java | 10 ++++------ 2 files changed, 5 insertions(+), 15 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java index 7a28c205bc..6f359c00b0 100644 --- a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java +++ b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java @@ -102,15 +102,7 @@ public class DefaultActorService extends TbApplicationEventListener staleRequest, long currentNs) { From 16a55b16e4f7c0f678cd10088b4f8ae94c5d5921 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 7 May 2021 15:56:45 +0300 Subject: [PATCH 12/28] js invoke service: logs added when script going to disable due to exception (special log for TimeoutException) --- .../service/script/AbstractJsInvokeService.java | 15 ++++++++++++--- .../script/AbstractNashornJsInvokeService.java | 2 +- .../service/script/RemoteJsInvokeService.java | 7 ++++--- .../common/DefaultTbQueueRequestTemplate.java | 2 +- 4 files changed, 18 insertions(+), 8 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java index 2f0378ae9d..b089f52e76 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java @@ -30,6 +30,7 @@ import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicInteger; /** @@ -84,8 +85,10 @@ public abstract class AbstractJsInvokeService implements JsInvokeService { apiUsageClient.report(tenantId, customerId, ApiUsageRecordKey.JS_EXEC_COUNT, 1); return doInvokeFunction(scriptId, functionName, args); } else { - return Futures.immediateFailedFuture( - new RuntimeException("Script invocation is blocked due to maximum error count " + getMaxErrors() + "!")); + String message = "Script invocation is blocked due to maximum error count " + + getMaxErrors() + ", scriptId " + scriptId + "!"; + log.warn(message); + return Futures.immediateFailedFuture(new RuntimeException(message)); } } else { return Futures.immediateFailedFuture(new RuntimeException("JS Execution is disabled due to API limits!")); @@ -117,7 +120,13 @@ public abstract class AbstractJsInvokeService implements JsInvokeService { protected abstract long getMaxBlacklistDuration(); - protected void onScriptExecutionError(UUID scriptId) { + protected void onScriptExecutionError(UUID scriptId, Throwable t) { + if (t instanceof TimeoutException || (t.getCause() != null && t.getCause() instanceof TimeoutException)) { + log.warn("Script has TimeoutException and will increment counter {} on disabledFunctions for id {}", //TODO remove after test + disabledFunctions.computeIfAbsent(scriptId, key -> new DisableListInfo()).get(), + scriptId); + //return; //timeout is not a good reason to disable function + } disabledFunctions.computeIfAbsent(scriptId, key -> new DisableListInfo()).incrementAndGet(); } diff --git a/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java index 9985ac60a2..f9fd3db977 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java @@ -160,7 +160,7 @@ public abstract class AbstractNashornJsInvokeService extends AbstractJsInvokeSer return ((Invocable) engine).invokeFunction(functionName, args); } } catch (Exception e) { - onScriptExecutionError(scriptId); + onScriptExecutionError(scriptId, e); throw new ExecutionException(e); } }); diff --git a/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java index 334a471973..b8bfeea75b 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java @@ -193,7 +193,7 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { @Override public void onFailure(Throwable t) { - onScriptExecutionError(scriptId); + onScriptExecutionError(scriptId, t); if (t instanceof TimeoutException || (t.getCause() != null && t.getCause() instanceof TimeoutException)) { queueTimeoutMsgs.incrementAndGet(); } @@ -205,9 +205,10 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { if (invokeResult.getSuccess()) { return invokeResult.getResult(); } else { - onScriptExecutionError(scriptId); + final RuntimeException e = new RuntimeException(invokeResult.getErrorDetails()); + onScriptExecutionError(scriptId, e); log.debug("[{}] Failed to compile script due to [{}]: {}", scriptId, invokeResult.getErrorCode().name(), invokeResult.getErrorDetails()); - throw new RuntimeException(invokeResult.getErrorDetails()); + throw e; } }, callbackExecutor); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index f6a5b33289..7b0e7a14cd 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -109,7 +109,7 @@ public class DefaultTbQueueRequestTemplate responses = doPoll(); //poll js responses //if (responses.size() > 0) { - log.trace("Completed template poll topic {}, for pendingRequests [{}], received [{}]", responseTemplate.getTopic(), pendingRequestsCount, responses.size()); + log.info("Completed template poll topic {}, for pendingRequests [{}], received [{}]", responseTemplate.getTopic(), pendingRequestsCount, responses.size()); //TODO reduce verbose after test //} responses.forEach(this::processResponse); //this can take a long time responseTemplate.commit(); From 24eca345768009b674ad827937247bcdd8cba523 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 5 May 2021 08:22:38 +0300 Subject: [PATCH 13/28] new utility class tbStopWatch added to improve performance measurement experience --- common/util/pom.xml | 5 ++ .../thingsboard/common/util/TbStopWatch.java | 62 +++++++++++++++++++ 2 files changed, 67 insertions(+) create mode 100644 common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java diff --git a/common/util/pom.xml b/common/util/pom.xml index 6f66da5790..128bc478d5 100644 --- a/common/util/pom.xml +++ b/common/util/pom.xml @@ -36,6 +36,11 @@ + + org.springframework + spring-core + ${spring.version} + com.google.guava guava diff --git a/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java b/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java new file mode 100644 index 0000000000..b7dbd0052c --- /dev/null +++ b/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java @@ -0,0 +1,62 @@ +/** + * ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL + * + * Copyright © 2016-2021 ThingsBoard, Inc. All Rights Reserved. + * + * NOTICE: All information contained herein is, and remains + * the property of ThingsBoard, Inc. and its suppliers, + * if any. The intellectual and technical concepts contained + * herein are proprietary to ThingsBoard, Inc. + * and its suppliers and may be covered by U.S. and Foreign Patents, + * patents in process, and are protected by trade secret or copyright law. + * + * Dissemination of this information or reproduction of this material is strictly forbidden + * unless prior written permission is obtained from COMPANY. + * + * Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees, + * managers or contractors who have executed Confidentiality and Non-disclosure agreements + * explicitly covering such access. + * + * The copyright notice above does not evidence any actual or intended publication + * or disclosure of this source code, which includes + * information that is confidential and/or proprietary, and is a trade secret, of COMPANY. + * ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE, + * OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT + * THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED, + * AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES. + * THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION + * DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS, + * OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART. + */ +package org.thingsboard.common.util; + +import org.springframework.util.StopWatch; + +/** + * Utility method that extends Spring Framework StopWatch + * It is a MONOTONIC time stopwatch. + * It is a replacement for any measurements with a wall-clock like System.currentTimeMillis() + * It is not affected by leap second, day-light saving and wall-clock adjustments by manual or network time synchronization + * The main features is a single call for common use cases: + * - create and start: TbStopWatch sw = TbStopWatch.startNew() + * - stop and get: sw.stopAndGetTotalTimeMillis() or sw.stopAndGetLastTaskTimeMillis() + * */ +public class TbStopWatch extends StopWatch { + + public static TbStopWatch startNew(){ + TbStopWatch stopWatch = new TbStopWatch(); + stopWatch.start(); + return stopWatch; + } + + public long stopAndGetTotalTimeMillis(){ + stop(); + return getTotalTimeMillis(); + } + + public long stopAndGetLastTaskTimeMillis(){ + stop(); + return getLastTaskTimeMillis(); + } + +} From e2aa4be7415790fdf371023fe22a4c25e5fcc605 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 7 May 2021 16:26:20 +0300 Subject: [PATCH 14/28] TbStopWatch: methods for nanos added --- .../java/org/thingsboard/common/util/TbStopWatch.java | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java b/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java index b7dbd0052c..4f5d1346a9 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java +++ b/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java @@ -54,9 +54,19 @@ public class TbStopWatch extends StopWatch { return getTotalTimeMillis(); } + public long stopAndGetTotalTimeNanos(){ + stop(); + return getLastTaskTimeNanos(); + } + public long stopAndGetLastTaskTimeMillis(){ stop(); return getLastTaskTimeMillis(); } + public long stopAndGetLastTaskTimeNanos(){ + stop(); + return getLastTaskTimeNanos(); + } + } From 9daa43a115bd584a409ea1279ca4ef6b4b0ae5cb Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 7 May 2021 16:34:35 +0300 Subject: [PATCH 15/28] queue request template: sleep on exception shortened according to the stopwatch. test adjusted --- .../queue/common/DefaultTbQueueRequestTemplate.java | 12 +++++++----- .../common/DefaultTbQueueRequestTemplateTest.java | 10 +++++----- 2 files changed, 12 insertions(+), 10 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index 7b0e7a14cd..a385ef17a8 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -21,6 +21,7 @@ import com.google.common.util.concurrent.SettableFuture; import lombok.Builder; import lombok.Getter; import lombok.extern.slf4j.Slf4j; +import org.thingsboard.common.util.TbStopWatch; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.queue.TbQueueAdmin; @@ -95,11 +96,13 @@ public class DefaultTbQueueRequestTemplate staleRequest, long currentNs) { diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java index 261ff9c7be..e60e4071f3 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java +++ b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java @@ -54,13 +54,13 @@ import java.util.concurrent.atomic.AtomicLong; import static org.hamcrest.Matchers.equalTo; import static org.hamcrest.Matchers.greaterThanOrEqualTo; import static org.hamcrest.Matchers.is; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertFalse; -import static org.junit.Assert.assertTrue; +import static org.hamcrest.Matchers.lessThan; import static org.mockito.BDDMockito.willAnswer; import static org.mockito.BDDMockito.willDoNothing; import static org.mockito.BDDMockito.willReturn; import static org.mockito.Matchers.any; +import static org.mockito.Matchers.anyLong; +import static org.mockito.Matchers.longThat; import static org.mockito.Mockito.RETURNS_DEEP_STUBS; import static org.mockito.Mockito.atLeastOnce; import static org.mockito.Mockito.mock; @@ -141,7 +141,7 @@ public class DefaultTbQueueRequestTemplateTest { public void givenMainLoop_whenLoopFewTimes_thenVerifyInvocationCount() throws InterruptedException { executor = inst.createExecutor(); CountDownLatch latch = new CountDownLatch(5); - willDoNothing().given(inst).sleep(); + willDoNothing().given(inst).sleep(anyLong()); willAnswer(invocation -> { if (latch.getCount() == 1) { inst.stop(); //stop the loop in natural way @@ -158,7 +158,7 @@ public class DefaultTbQueueRequestTemplateTest { latch.await(10, TimeUnit.SECONDS); verify(inst, times(5)).fetchAndProcessResponses(); - verify(inst, times(2)).sleep(); + verify(inst, times(2)).sleep(longThat(lessThan(TimeUnit.MILLISECONDS.toNanos(inst.pollInterval)))); } @Test From c3aec4d077220aef22f369c6d6ee692205b75a47 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Fri, 7 May 2021 15:59:42 +0300 Subject: [PATCH 16/28] revert initDispatcherExecutor --- .../server/actors/service/DefaultActorService.java | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java index 6f359c00b0..98f739a24b 100644 --- a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java +++ b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java @@ -102,7 +102,16 @@ public class DefaultActorService extends TbApplicationEventListener Date: Tue, 11 May 2021 08:01:21 +0300 Subject: [PATCH 17/28] log WARN for onScriptExecutionError to prevent silent ban (actually, system going to idle if the script is a part of a rule chain). restored TRACE level on fetchAndProcessResponses --- .../service/script/AbstractJsInvokeService.java | 17 ++++++++++------- .../common/DefaultTbQueueRequestTemplate.java | 6 ++---- 2 files changed, 12 insertions(+), 11 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java index b089f52e76..dc01e0d97e 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java @@ -121,13 +121,16 @@ public abstract class AbstractJsInvokeService implements JsInvokeService { protected abstract long getMaxBlacklistDuration(); protected void onScriptExecutionError(UUID scriptId, Throwable t) { - if (t instanceof TimeoutException || (t.getCause() != null && t.getCause() instanceof TimeoutException)) { - log.warn("Script has TimeoutException and will increment counter {} on disabledFunctions for id {}", //TODO remove after test - disabledFunctions.computeIfAbsent(scriptId, key -> new DisableListInfo()).get(), - scriptId); - //return; //timeout is not a good reason to disable function - } - disabledFunctions.computeIfAbsent(scriptId, key -> new DisableListInfo()).incrementAndGet(); + DisableListInfo disableListInfo = disabledFunctions.computeIfAbsent(scriptId, key -> new DisableListInfo()); + log.warn("Script has exception and will increment counter {} on disabledFunctions for id {}, exception {}, cause {}", + disableListInfo.get(), scriptId, t.getClass(), t.getCause()); +// if (t instanceof TimeoutException || (t.getCause() != null && t.getCause() instanceof TimeoutException)) { +// log.warn("Script has TimeoutException and will increment counter {} on disabledFunctions for id {}", //TODO remove after test +// disableListInfo.get(), +// scriptId); +// //return; //timeout is not a good reason to disable function +// } + disableListInfo.incrementAndGet(); } private String generateJsScript(JsScriptType scriptType, String functionName, String scriptBody, String... argNames) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index a385ef17a8..abeb739160 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -109,11 +109,9 @@ public class DefaultTbQueueRequestTemplate responses = doPoll(); //poll js responses - //if (responses.size() > 0) { - log.info("Completed template poll topic {}, for pendingRequests [{}], received [{}]", responseTemplate.getTopic(), pendingRequestsCount, responses.size()); //TODO reduce verbose after test - //} + log.trace("Completed template poll topic {}, for pendingRequests [{}], received [{}] responses", responseTemplate.getTopic(), pendingRequestsCount, responses.size()); responses.forEach(this::processResponse); //this can take a long time responseTemplate.commit(); tryCleanStaleRequests(); From d3e4102282b581b055ec34c0fa555711a19ae472 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 11 May 2021 08:11:15 +0300 Subject: [PATCH 18/28] removed version from pom.xml from common/util for spring-core (version already defined through the root dependency management) --- common/util/pom.xml | 1 - 1 file changed, 1 deletion(-) diff --git a/common/util/pom.xml b/common/util/pom.xml index 128bc478d5..8a0fa1705b 100644 --- a/common/util/pom.xml +++ b/common/util/pom.xml @@ -39,7 +39,6 @@ org.springframework spring-core - ${spring.version} com.google.guava From da65eaab473d1f89131b71b0710d0c3c53e130cf Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 11 May 2021 08:20:16 +0300 Subject: [PATCH 19/28] added spring-core for the root dependency management --- pom.xml | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/pom.xml b/pom.xml index b7d0eb71f6..5f88903559 100755 --- a/pom.xml +++ b/pom.xml @@ -1009,6 +1009,11 @@ + + org.springframework + spring-core + ${spring.version} + org.springframework.boot spring-boot-starter-web From ddfb57db588dc0e1de67dd662af49a50c2ed0cd2 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 11 May 2021 09:30:58 +0300 Subject: [PATCH 20/28] log detailed error onScriptExecutionError for debug --- .../server/service/script/AbstractJsInvokeService.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java index dc01e0d97e..8a6aeb399b 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java @@ -123,7 +123,8 @@ public abstract class AbstractJsInvokeService implements JsInvokeService { protected void onScriptExecutionError(UUID scriptId, Throwable t) { DisableListInfo disableListInfo = disabledFunctions.computeIfAbsent(scriptId, key -> new DisableListInfo()); log.warn("Script has exception and will increment counter {} on disabledFunctions for id {}, exception {}, cause {}", - disableListInfo.get(), scriptId, t.getClass(), t.getCause()); + disableListInfo.get(), scriptId, t, t.getCause()); + log.error("onScriptExecutionError", t); // if (t instanceof TimeoutException || (t.getCause() != null && t.getCause() instanceof TimeoutException)) { // log.warn("Script has TimeoutException and will increment counter {} on disabledFunctions for id {}", //TODO remove after test // disableListInfo.get(), From 1d33c09127165907def60b09bca41a11b60d8d85 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Tue, 11 May 2021 16:22:10 +0300 Subject: [PATCH 21/28] JsInvokeService - log js body on script failure to investigate easy --- .../server/service/script/AbstractJsInvokeService.java | 7 +++---- .../service/script/AbstractNashornJsInvokeService.java | 2 +- .../server/service/script/RemoteJsInvokeService.java | 10 +++++----- 3 files changed, 9 insertions(+), 10 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java index 8a6aeb399b..cae4a8e05c 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/AbstractJsInvokeService.java @@ -120,11 +120,10 @@ public abstract class AbstractJsInvokeService implements JsInvokeService { protected abstract long getMaxBlacklistDuration(); - protected void onScriptExecutionError(UUID scriptId, Throwable t) { + protected void onScriptExecutionError(UUID scriptId, Throwable t, String scriptBody) { DisableListInfo disableListInfo = disabledFunctions.computeIfAbsent(scriptId, key -> new DisableListInfo()); - log.warn("Script has exception and will increment counter {} on disabledFunctions for id {}, exception {}, cause {}", - disableListInfo.get(), scriptId, t, t.getCause()); - log.error("onScriptExecutionError", t); + log.warn("Script has exception and will increment counter {} on disabledFunctions for id {}, exception {}, cause {}, scriptBody {}", + disableListInfo.get(), scriptId, t, t.getCause(), scriptBody); // if (t instanceof TimeoutException || (t.getCause() != null && t.getCause() instanceof TimeoutException)) { // log.warn("Script has TimeoutException and will increment counter {} on disabledFunctions for id {}", //TODO remove after test // disableListInfo.get(), diff --git a/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java index f9fd3db977..15a3cf1c15 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/AbstractNashornJsInvokeService.java @@ -160,7 +160,7 @@ public abstract class AbstractNashornJsInvokeService extends AbstractJsInvokeSer return ((Invocable) engine).invokeFunction(functionName, args); } } catch (Exception e) { - onScriptExecutionError(scriptId, e); + onScriptExecutionError(scriptId, e, functionName); throw new ExecutionException(e); } }); diff --git a/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java index b8bfeea75b..5b3b1f5e23 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java @@ -18,7 +18,6 @@ package org.thingsboard.server.service.script; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.MoreExecutors; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; @@ -161,7 +160,8 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { @Override protected ListenableFuture doInvokeFunction(UUID scriptId, String functionName, Object[] args) { - String scriptBody = scriptIdToBodysMap.get(scriptId); + log.trace("doInvokeFunction js-request for uuid {} with timeout {}ms", scriptId, maxRequestsTimeout); + final String scriptBody = scriptIdToBodysMap.get(scriptId); if (scriptBody == null) { return Futures.immediateFailedFuture(new RuntimeException("No script body found for scriptId: [" + scriptId + "]!")); } @@ -170,7 +170,7 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { .setScriptIdLSB(scriptId.getLeastSignificantBits()) .setFunctionName(functionName) .setTimeout((int) maxRequestsTimeout) - .setScriptBody(scriptIdToBodysMap.get(scriptId)); + .setScriptBody(scriptBody); for (Object arg : args) { jsRequestBuilder.addArgs(arg.toString()); @@ -193,7 +193,7 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { @Override public void onFailure(Throwable t) { - onScriptExecutionError(scriptId, t); + onScriptExecutionError(scriptId, t, scriptBody); if (t instanceof TimeoutException || (t.getCause() != null && t.getCause() instanceof TimeoutException)) { queueTimeoutMsgs.incrementAndGet(); } @@ -206,7 +206,7 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { return invokeResult.getResult(); } else { final RuntimeException e = new RuntimeException(invokeResult.getErrorDetails()); - onScriptExecutionError(scriptId, e); + onScriptExecutionError(scriptId, e, scriptBody); log.debug("[{}] Failed to compile script due to [{}]: {}", scriptId, invokeResult.getErrorCode().name(), invokeResult.getErrorDetails()); throw e; } From 7e551cc92a94aef3727f5732837846c588a4805b Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 13 May 2021 20:23:24 +0300 Subject: [PATCH 22/28] log: reduced severity to trace for DefaultTbQueueRequestTemplate.setupNextCleanup() --- .../server/queue/common/DefaultTbQueueRequestTemplate.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index abeb739160..7fa3ac5989 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -143,7 +143,7 @@ public class DefaultTbQueueRequestTemplate doPoll() { From ed44ac8f19e14a652978123147f9869e478c2958 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 17 Jun 2021 13:27:19 +0300 Subject: [PATCH 23/28] test fixed for DefaultTbQueueRequestTemplateTest --- .../queue/common/DefaultTbQueueRequestTemplateTest.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java index e60e4071f3..c73be715fe 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java +++ b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java @@ -37,7 +37,7 @@ import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.ArgumentCaptor; import org.mockito.Mock; -import org.mockito.runners.MockitoJUnitRunner; +import org.mockito.junit.MockitoJUnitRunner; import org.thingsboard.server.queue.TbQueueAdmin; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueMsg; @@ -55,12 +55,11 @@ import static org.hamcrest.Matchers.equalTo; import static org.hamcrest.Matchers.greaterThanOrEqualTo; import static org.hamcrest.Matchers.is; import static org.hamcrest.Matchers.lessThan; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.BDDMockito.willAnswer; import static org.mockito.BDDMockito.willDoNothing; import static org.mockito.BDDMockito.willReturn; -import static org.mockito.Matchers.any; -import static org.mockito.Matchers.anyLong; -import static org.mockito.Matchers.longThat; import static org.mockito.Mockito.RETURNS_DEEP_STUBS; import static org.mockito.Mockito.atLeastOnce; import static org.mockito.Mockito.mock; @@ -70,6 +69,7 @@ import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.hamcrest.MatcherAssert.assertThat; +import static org.mockito.hamcrest.MockitoHamcrest.longThat; @Slf4j @RunWith(MockitoJUnitRunner.class) From 63c906250c39799a8b0142c529d6b21d71d52477 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 17 Jun 2021 13:27:52 +0300 Subject: [PATCH 24/28] log cleanup DefaultTbQueueRequestTemplate --- .../server/queue/common/DefaultTbQueueRequestTemplate.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java index 7fa3ac5989..7463c11df5 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplate.java @@ -245,7 +245,7 @@ public class DefaultTbQueueRequestTemplate future, ResponseMetaData responseMetaData) { - log.trace("[{}] Sending request, key [{}], expTime [{}], request {}", requestId, request.getKey(), responseMetaData.expTime, String.valueOf(request).replace("\n", " ")); + log.trace("[{}] Sending request, key [{}], expTime [{}], request {}", requestId, request.getKey(), responseMetaData.expTime, request); if (messagesStats != null) { messagesStats.incrementTotal(); } @@ -255,7 +255,7 @@ public class DefaultTbQueueRequestTemplate Date: Thu, 17 Jun 2021 13:34:06 +0300 Subject: [PATCH 25/28] code cleanup --- .../server/actors/service/DefaultActorService.java | 1 - .../server/service/script/AbstractJsInvokeService.java | 6 ------ 2 files changed, 7 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java index 98f739a24b..7a28c205bc 100644 --- a/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java +++ b/application/src/main/java/org/thingsboard/server/actors/service/DefaultActorService.java @@ -111,7 +111,6 @@ public class DefaultActorService extends TbApplicationEventListener new DisableListInfo()); log.warn("Script has exception and will increment counter {} on disabledFunctions for id {}, exception {}, cause {}, scriptBody {}", disableListInfo.get(), scriptId, t, t.getCause(), scriptBody); -// if (t instanceof TimeoutException || (t.getCause() != null && t.getCause() instanceof TimeoutException)) { -// log.warn("Script has TimeoutException and will increment counter {} on disabledFunctions for id {}", //TODO remove after test -// disableListInfo.get(), -// scriptId); -// //return; //timeout is not a good reason to disable function -// } disableListInfo.incrementAndGet(); } From 75603e58bec47a8b30292cdaa57b03831547514f Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 17 Jun 2021 13:34:54 +0300 Subject: [PATCH 26/28] TbStopWatch updated licence header for CE --- .../thingsboard/common/util/TbStopWatch.java | 35 ++++++------------- 1 file changed, 10 insertions(+), 25 deletions(-) diff --git a/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java b/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java index 4f5d1346a9..90f58ce7f2 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java +++ b/common/util/src/main/java/org/thingsboard/common/util/TbStopWatch.java @@ -1,32 +1,17 @@ /** - * ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL + * Copyright © 2016-2021 The Thingsboard Authors * - * Copyright © 2016-2021 ThingsBoard, Inc. All Rights Reserved. + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at * - * NOTICE: All information contained herein is, and remains - * the property of ThingsBoard, Inc. and its suppliers, - * if any. The intellectual and technical concepts contained - * herein are proprietary to ThingsBoard, Inc. - * and its suppliers and may be covered by U.S. and Foreign Patents, - * patents in process, and are protected by trade secret or copyright law. + * http://www.apache.org/licenses/LICENSE-2.0 * - * Dissemination of this information or reproduction of this material is strictly forbidden - * unless prior written permission is obtained from COMPANY. - * - * Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees, - * managers or contractors who have executed Confidentiality and Non-disclosure agreements - * explicitly covering such access. - * - * The copyright notice above does not evidence any actual or intended publication - * or disclosure of this source code, which includes - * information that is confidential and/or proprietary, and is a trade secret, of COMPANY. - * ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE, - * OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT - * THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED, - * AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES. - * THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION - * DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS, - * OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART. + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. */ package org.thingsboard.common.util; From aaedb9e879ccd1ad09b14d5cc717ec981f6d1927 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 17 Jun 2021 13:35:20 +0300 Subject: [PATCH 27/28] DefaultTbQueueRequestTemplateTest updated licence header for CE --- .../DefaultTbQueueRequestTemplateTest.java | 35 ++++++------------- 1 file changed, 10 insertions(+), 25 deletions(-) diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java index c73be715fe..e50e204adb 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java +++ b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java @@ -1,32 +1,17 @@ /** - * ThingsBoard, Inc. ("COMPANY") CONFIDENTIAL + * Copyright © 2016-2021 The Thingsboard Authors * - * Copyright © 2016-2021 ThingsBoard, Inc. All Rights Reserved. + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at * - * NOTICE: All information contained herein is, and remains - * the property of ThingsBoard, Inc. and its suppliers, - * if any. The intellectual and technical concepts contained - * herein are proprietary to ThingsBoard, Inc. - * and its suppliers and may be covered by U.S. and Foreign Patents, - * patents in process, and are protected by trade secret or copyright law. + * http://www.apache.org/licenses/LICENSE-2.0 * - * Dissemination of this information or reproduction of this material is strictly forbidden - * unless prior written permission is obtained from COMPANY. - * - * Access to the source code contained herein is hereby forbidden to anyone except current COMPANY employees, - * managers or contractors who have executed Confidentiality and Non-disclosure agreements - * explicitly covering such access. - * - * The copyright notice above does not evidence any actual or intended publication - * or disclosure of this source code, which includes - * information that is confidential and/or proprietary, and is a trade secret, of COMPANY. - * ANY REPRODUCTION, MODIFICATION, DISTRIBUTION, PUBLIC PERFORMANCE, - * OR PUBLIC DISPLAY OF OR THROUGH USE OF THIS SOURCE CODE WITHOUT - * THE EXPRESS WRITTEN CONSENT OF COMPANY IS STRICTLY PROHIBITED, - * AND IN VIOLATION OF APPLICABLE LAWS AND INTERNATIONAL TREATIES. - * THE RECEIPT OR POSSESSION OF THIS SOURCE CODE AND/OR RELATED INFORMATION - * DOES NOT CONVEY OR IMPLY ANY RIGHTS TO REPRODUCE, DISCLOSE OR DISTRIBUTE ITS CONTENTS, - * OR TO MANUFACTURE, USE, OR SELL ANYTHING THAT IT MAY DESCRIBE, IN WHOLE OR IN PART. + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. */ package org.thingsboard.server.queue.common; From bec228bc83d1541ed69c92ebb77fd13bb606eb47 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Thu, 17 Jun 2021 13:42:37 +0300 Subject: [PATCH 28/28] test: removed unnecessary stubbing on DefaultTbQueueRequestTemplateTest --- .../queue/common/DefaultTbQueueRequestTemplateTest.java | 4 ---- 1 file changed, 4 deletions(-) diff --git a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java index e50e204adb..9979e9ac43 100644 --- a/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java +++ b/common/queue/src/test/java/org/thingsboard/server/queue/common/DefaultTbQueueRequestTemplateTest.java @@ -105,10 +105,7 @@ public class DefaultTbQueueRequestTemplateTest { @Test public void givenExternalExecutor_whenInitStop_thenOK() { - willDoNothing().given(inst).mainLoop(); - inst.init(); - //assertNotEquals(0, inst.tickTs); assertThat(inst.nextCleanupNs, equalTo(0L)); verify(queueAdmin, times(1)).createTopicIfNotExists(topic); verify(requestTemplate, times(1)).init(); @@ -184,7 +181,6 @@ public class DefaultTbQueueRequestTemplateTest { inst.init(); inst.setupNextCleanup(); willReturn(Collections.emptyList()).given(inst).doPoll(); - willDoNothing().given(inst).processResponse(any()); //when long stepNs = TimeUnit.MILLISECONDS.toNanos(1);