From 8b637c9e9456963e1e24ec0e6306903d2e765c1f Mon Sep 17 00:00:00 2001 From: vparomskiy Date: Fri, 23 Mar 2018 16:54:48 +0200 Subject: [PATCH 1/3] add cassandra properties in test suite --- .../server/actors/plugin/PluginActorMessageProcessor.java | 1 + application/src/main/resources/thingsboard.yml | 4 ++-- .../quota/inmemory/HostRequestIntervalRegistry.java | 5 +++-- .../server/dao/timeseries/CassandraBaseTimeseriesDao.java | 3 --- .../thingsboard/server/dao/util/BufferedRateLimiter.java | 8 ++++++-- dao/src/test/resources/cassandra-test.properties | 5 +++++ 6 files changed, 17 insertions(+), 9 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/plugin/PluginActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/plugin/PluginActorMessageProcessor.java index 7dfe9a5d81..2af30d4b49 100644 --- a/application/src/main/java/org/thingsboard/server/actors/plugin/PluginActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/plugin/PluginActorMessageProcessor.java @@ -106,6 +106,7 @@ public class PluginActorMessageProcessor extends ComponentMsgProcessor try { pluginImpl.process(trustedCtx, msg.getRuleTenantId(), msg.getRuleId(), msg.getMsg()); } catch (Exception ex) { + logger.debug("[{}] Failed to process RuleToPlugin msg: [{}] [{}]", tenantId, msg.getMsg(), ex); RuleToPluginMsg ruleMsg = msg.getMsg(); MsgType responceMsgType = MsgType.RULE_ENGINE_ERROR; Integer requestId = 0; diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 07face1059..c47fc28fbd 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -133,7 +133,7 @@ quota: intervalMin: 2 database: - type: "${DATABASE_TYPE:cassandra}" # cassandra OR sql + type: "${DATABASE_TYPE:sql}" # cassandra OR sql # Cassandra driver configuration parameters cassandra: @@ -226,7 +226,7 @@ caffeine: specs: relations: timeToLiveInMinutes: 1440 - maxSize: 0 + maxSize: 100000 deviceCredentials: timeToLiveInMinutes: 1440 maxSize: 100000 diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/HostRequestIntervalRegistry.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/HostRequestIntervalRegistry.java index 8d254a08c6..3782ed22ed 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/HostRequestIntervalRegistry.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/HostRequestIntervalRegistry.java @@ -61,13 +61,14 @@ public class HostRequestIntervalRegistry { } public long tick(String clientHostId) { + IntervalCount intervalCount = hostCounts.computeIfAbsent(clientHostId, s -> new IntervalCount(intervalDurationMs)); + long currentCount = intervalCount.resetIfExpiredAndTick(); if (whiteList.contains(clientHostId)) { return 0; } else if (blackList.contains(clientHostId)) { return Long.MAX_VALUE; } - IntervalCount intervalCount = hostCounts.computeIfAbsent(clientHostId, s -> new IntervalCount(intervalDurationMs)); - return intervalCount.resetIfExpiredAndTick(); + return currentCount; } public void clean() { diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index cf141711a5..cda4b1669b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -439,8 +439,6 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem private PreparedStatement getLatestStmt() { if (latestInsertStmt == null) { -// latestInsertStmt = new PreparedStatement[DataType.values().length]; -// for (DataType type : DataType.values()) { latestInsertStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_LATEST_CF + "(" + ModelConstants.ENTITY_TYPE_COLUMN + "," + ModelConstants.ENTITY_ID_COLUMN + @@ -451,7 +449,6 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem "," + ModelConstants.LONG_VALUE_COLUMN + "," + ModelConstants.DOUBLE_VALUE_COLUMN + ")" + " VALUES(?, ?, ?, ?, ?, ?, ?, ?)"); -// } } return latestInsertStmt; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateLimiter.java b/dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateLimiter.java index de07dbfa47..2acd623a37 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateLimiter.java +++ b/dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateLimiter.java @@ -29,6 +29,7 @@ import java.util.concurrent.atomic.AtomicInteger; @Component @Slf4j +@NoSqlDao public class BufferedRateLimiter implements AsyncRateLimiter { private final ListeningExecutorService pool = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(10)); @@ -113,6 +114,9 @@ public class BufferedRateLimiter implements AsyncRateLimiter { lockedFuture.cancelFuture(); return Futures.immediateFailedFuture(new IllegalStateException("Rate Limit Buffer is full. Reject")); } + if(permits.get() < permitsLimit) { + reprocessQueue(); + } return lockedFuture.future; } catch (InterruptedException e) { return Futures.immediateFailedFuture(new IllegalStateException("Rate Limit Task interrupted. Reject")); @@ -130,8 +134,8 @@ public class BufferedRateLimiter implements AsyncRateLimiter { expiredCount++; } } - log.info("Permits maxBuffer is [{}] max concurrent [{}] expired [{}]", maxQueueSize.getAndSet(0), - maxGrantedPermissions.getAndSet(0), expiredCount); + log.info("Permits maxBuffer is [{}] max concurrent [{}] expired [{}] current granted [{}]", maxQueueSize.getAndSet(0), + maxGrantedPermissions.getAndSet(0), expiredCount, permits.get()); } private class LockedFuture { diff --git a/dao/src/test/resources/cassandra-test.properties b/dao/src/test/resources/cassandra-test.properties index 82fcbe1949..737687f053 100644 --- a/dao/src/test/resources/cassandra-test.properties +++ b/dao/src/test/resources/cassandra-test.properties @@ -47,3 +47,8 @@ cassandra.query.default_fetch_size=2000 cassandra.query.ts_key_value_partitioning=HOURS cassandra.query.max_limit_per_request=1000 +cassandra.query.buffer_size=100000 +cassandra.query.concurrent_limit=1000 +cassandra.query.permit_max_wait_time=20000 +cassandra.query.rate_limit_print_interval_ms=30000 + From 54b272c04f9370a9f4e106cade3c511959fc9a18 Mon Sep 17 00:00:00 2001 From: vparomskiy Date: Tue, 27 Mar 2018 13:16:48 +0300 Subject: [PATCH 2/3] return permit if request expired/canceled --- .../dao/exception/BufferLimitException.java | 25 +++++++++++++++ .../dao/nosql/RateLimitedResultSetFuture.java | 14 +++++--- .../server/dao/util/BufferedRateLimiter.java | 21 +++++++++--- .../nosql/RateLimitedResultSetFutureTest.java | 32 +++++++++++++++++-- .../dao/util/BufferedRateLimiterTest.java | 5 +-- 5 files changed, 82 insertions(+), 15 deletions(-) create mode 100644 dao/src/main/java/org/thingsboard/server/dao/exception/BufferLimitException.java diff --git a/dao/src/main/java/org/thingsboard/server/dao/exception/BufferLimitException.java b/dao/src/main/java/org/thingsboard/server/dao/exception/BufferLimitException.java new file mode 100644 index 0000000000..3334dc62a9 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/exception/BufferLimitException.java @@ -0,0 +1,25 @@ +/** + * Copyright © 2016-2018 The Thingsboard Authors + * + * 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 + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * 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.dao.exception; + +public class BufferLimitException extends RuntimeException { + + private static final long serialVersionUID = 4513762009041887588L; + + public BufferLimitException() { + super("Rate Limit Buffer is full"); + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/nosql/RateLimitedResultSetFuture.java b/dao/src/main/java/org/thingsboard/server/dao/nosql/RateLimitedResultSetFuture.java index 2674c6ddea..d2505632d7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/nosql/RateLimitedResultSetFuture.java +++ b/dao/src/main/java/org/thingsboard/server/dao/nosql/RateLimitedResultSetFuture.java @@ -24,6 +24,7 @@ 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.Uninterruptibles; +import org.thingsboard.server.dao.exception.BufferLimitException; import org.thingsboard.server.dao.util.AsyncRateLimiter; import javax.annotation.Nullable; @@ -35,9 +36,15 @@ public class RateLimitedResultSetFuture implements ResultSetFuture { private final ListenableFuture rateLimitFuture; public RateLimitedResultSetFuture(Session session, AsyncRateLimiter rateLimiter, Statement statement) { - this.rateLimitFuture = rateLimiter.acquireAsync(); + this.rateLimitFuture = Futures.withFallback(rateLimiter.acquireAsync(), t -> { + if (!(t instanceof BufferLimitException)) { + rateLimiter.release(); + } + return Futures.immediateFailedFuture(t); + }); this.originalFuture = Futures.transform(rateLimitFuture, (Function) i -> executeAsyncWithRelease(rateLimiter, session, statement)); + } @Override @@ -108,10 +115,7 @@ public class RateLimitedResultSetFuture implements ResultSetFuture { try { ResultSetFuture resultSetFuture = Uninterruptibles.getUninterruptibly(originalFuture); resultSetFuture.addListener(listener, executor); - } catch (CancellationException e) { - cancel(false); - return; - } catch (ExecutionException e) { + } catch (CancellationException | ExecutionException e) { Futures.immediateFailedFuture(e).addListener(listener, executor); } }, executor); diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateLimiter.java b/dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateLimiter.java index 2acd623a37..03eb46f1ab 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateLimiter.java +++ b/dao/src/main/java/org/thingsboard/server/dao/util/BufferedRateLimiter.java @@ -23,6 +23,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; +import org.thingsboard.server.dao.exception.BufferLimitException; import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; @@ -41,6 +42,9 @@ public class BufferedRateLimiter implements AsyncRateLimiter { private final AtomicInteger maxQueueSize = new AtomicInteger(); private final AtomicInteger maxGrantedPermissions = new AtomicInteger(); + private final AtomicInteger totalGranted = new AtomicInteger(); + private final AtomicInteger totalReleased = new AtomicInteger(); + private final AtomicInteger totalRequested = new AtomicInteger(); public BufferedRateLimiter(@Value("${cassandra.query.buffer_size}") int queueLimit, @Value("${cassandra.query.concurrent_limit}") int permitsLimit, @@ -53,11 +57,13 @@ public class BufferedRateLimiter implements AsyncRateLimiter { @Override public ListenableFuture acquireAsync() { + totalRequested.incrementAndGet(); if (queue.isEmpty()) { if (permits.incrementAndGet() <= permitsLimit) { if (permits.get() > maxGrantedPermissions.get()) { maxGrantedPermissions.set(permits.get()); } + totalGranted.incrementAndGet(); return Futures.immediateFuture(null); } permits.decrementAndGet(); @@ -69,6 +75,7 @@ public class BufferedRateLimiter implements AsyncRateLimiter { @Override public void release() { permits.decrementAndGet(); + totalReleased.incrementAndGet(); reprocessQueue(); } @@ -80,6 +87,7 @@ public class BufferedRateLimiter implements AsyncRateLimiter { } LockedFuture lockedFuture = queue.poll(); if (lockedFuture != null) { + totalGranted.incrementAndGet(); lockedFuture.latch.countDown(); } else { permits.decrementAndGet(); @@ -112,17 +120,17 @@ public class BufferedRateLimiter implements AsyncRateLimiter { LockedFuture lockedFuture = createLockedFuture(); if (!queue.offer(lockedFuture, 1, TimeUnit.SECONDS)) { lockedFuture.cancelFuture(); - return Futures.immediateFailedFuture(new IllegalStateException("Rate Limit Buffer is full. Reject")); + return Futures.immediateFailedFuture(new BufferLimitException()); } if(permits.get() < permitsLimit) { reprocessQueue(); } return lockedFuture.future; } catch (InterruptedException e) { - return Futures.immediateFailedFuture(new IllegalStateException("Rate Limit Task interrupted. Reject")); + return Futures.immediateFailedFuture(new BufferLimitException()); } } - return Futures.immediateFailedFuture(new IllegalStateException("Rate Limit Buffer is full. Reject")); + return Futures.immediateFailedFuture(new BufferLimitException()); } @Scheduled(fixedDelayString = "${cassandra.query.rate_limit_print_interval_ms}") @@ -134,8 +142,11 @@ public class BufferedRateLimiter implements AsyncRateLimiter { expiredCount++; } } - log.info("Permits maxBuffer is [{}] max concurrent [{}] expired [{}] current granted [{}]", maxQueueSize.getAndSet(0), - maxGrantedPermissions.getAndSet(0), expiredCount, permits.get()); + log.info("Permits maxBuffer [{}] maxPermits [{}] expired [{}] currPermits [{}] currBuffer [{}] " + + "totalPermits [{}] totalRequests [{}] totalReleased [{}]", + maxQueueSize.getAndSet(0), maxGrantedPermissions.getAndSet(0), expiredCount, + permits.get(), queue.size(), + totalGranted.getAndSet(0), totalRequested.getAndSet(0), totalReleased.getAndSet(0)); } private class LockedFuture { diff --git a/dao/src/test/java/org/thingsboard/server/dao/nosql/RateLimitedResultSetFutureTest.java b/dao/src/test/java/org/thingsboard/server/dao/nosql/RateLimitedResultSetFutureTest.java index fa62c2b9b0..f49668d3fd 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/nosql/RateLimitedResultSetFutureTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/nosql/RateLimitedResultSetFutureTest.java @@ -19,16 +19,17 @@ import com.datastax.driver.core.*; import com.datastax.driver.core.exceptions.UnsupportedFeatureException; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.MoreExecutors; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Mock; import org.mockito.Mockito; import org.mockito.runners.MockitoJUnitRunner; import org.mockito.stubbing.Answer; +import org.thingsboard.server.dao.exception.BufferLimitException; import org.thingsboard.server.dao.util.AsyncRateLimiter; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.TimeoutException; +import java.util.concurrent.*; import static org.junit.Assert.*; import static org.mockito.Mockito.*; @@ -53,7 +54,7 @@ public class RateLimitedResultSetFutureTest { @Test public void doNotReleasePermissionIfRateLimitFutureFailed() throws InterruptedException { - when(rateLimiter.acquireAsync()).thenReturn(Futures.immediateFailedFuture(new IllegalArgumentException())); + when(rateLimiter.acquireAsync()).thenReturn(Futures.immediateFailedFuture(new BufferLimitException())); resultSetFuture = new RateLimitedResultSetFuture(session, rateLimiter, statement); Thread.sleep(1000L); verify(rateLimiter).acquireAsync(); @@ -153,4 +154,29 @@ public class RateLimitedResultSetFutureTest { verify(rateLimiter, times(1)).release(); } + @Test + public void expiredQueryReturnPermit() throws InterruptedException, ExecutionException { + CountDownLatch latch = new CountDownLatch(1); + ListenableFuture future = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(1)).submit(() -> { + latch.await(); + return null; + }); + when(rateLimiter.acquireAsync()).thenReturn(future); + resultSetFuture = new RateLimitedResultSetFuture(session, rateLimiter, statement); + + ListenableFuture transform = Futures.transform(resultSetFuture, ResultSet::one); +// TimeUnit.MILLISECONDS.sleep(200); + future.cancel(false); + latch.countDown(); + + try { + transform.get(); + fail(); + } catch (Exception e) { + assertTrue(e instanceof ExecutionException); + } + verify(rateLimiter, times(1)).acquireAsync(); + verify(rateLimiter, times(1)).release(); + } + } \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/util/BufferedRateLimiterTest.java b/dao/src/test/java/org/thingsboard/server/dao/util/BufferedRateLimiterTest.java index 5bfc3b6e95..67c3ce8d73 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/util/BufferedRateLimiterTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/util/BufferedRateLimiterTest.java @@ -17,6 +17,7 @@ package org.thingsboard.server.dao.util; import com.google.common.util.concurrent.*; import org.junit.Test; +import org.thingsboard.server.dao.exception.BufferLimitException; import javax.annotation.Nullable; import java.util.concurrent.ExecutionException; @@ -61,8 +62,8 @@ public class BufferedRateLimiterTest { } catch (Exception e) { assertTrue(e instanceof ExecutionException); Throwable actualCause = e.getCause(); - assertTrue(actualCause instanceof IllegalStateException); - assertEquals("Rate Limit Buffer is full. Reject", actualCause.getMessage()); + assertTrue(actualCause instanceof BufferLimitException); + assertEquals("Rate Limit Buffer is full", actualCause.getMessage()); } } From 2f6995fcedb0f795429a3d4b0ad1c7fb687314c7 Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Tue, 27 Mar 2018 16:59:41 +0300 Subject: [PATCH 3/3] Fix for UI build. --- ui/package.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/ui/package.json b/ui/package.json index ad95ef4a7a..ad9a7a67a3 100644 --- a/ui/package.json +++ b/ui/package.json @@ -15,7 +15,7 @@ }, "dependencies": { "@flowjs/ng-flow": "^2.7.1", - "ace-builds": "^1.2.5", + "ace-builds": "1.3.1", "angular": "1.5.8", "angular-animate": "1.5.8", "angular-aria": "1.5.8",