From 0baa4b3ac944c22cdd5949794d6015c93589124d Mon Sep 17 00:00:00 2001 From: vparomskiy Date: Thu, 17 May 2018 09:57:48 +0300 Subject: [PATCH 1/5] implement API limit for tenant based on Rule Messages count --- .../service/queue/DefaultMsgQueueService.java | 9 +++++ .../src/main/resources/thingsboard.yml | 25 ++++++++++-- ...Service.java => AbstractQuotaService.java} | 36 +++-------------- .../transport/quota/RequestLimitPolicy.java | 15 +++++++ .../host/HostIntervalRegistryCleaner.java | 14 +++++++ .../host/HostIntervalRegistryLogger.java | 37 +++++++++++++++++ .../host/HostRequestIntervalRegistry.java | 37 +++++++++++++++++ .../{ => host}/HostRequestLimitPolicy.java | 21 ++++------ .../quota/host/HostRequestsQuotaService.java | 37 +++++++++++++++++ .../inmemory/IntervalRegistryCleaner.java | 10 ++--- .../inmemory/IntervalRegistryLogger.java | 32 ++++----------- ...try.java => KeyBasedIntervalRegistry.java} | 40 ++++--------------- .../tenant/TenantIntervalRegistryCleaner.java | 14 +++++++ .../tenant/TenantIntervalRegistryLogger.java | 37 +++++++++++++++++ .../tenant/TenantMsgsIntervalRegistry.java | 16 ++++++++ .../quota/tenant/TenantQuotaService.java | 15 +++++++ .../tenant/TenantRequestLimitPolicy.java | 13 ++++++ .../quota/HostRequestLimitPolicyTest.java | 1 + .../quota/HostRequestsQuotaServiceTest.java | 8 ++-- .../HostRequestIntervalRegistryTest.java | 4 +- .../inmemory/IntervalRegistryLoggerTest.java | 4 +- .../transport/coap/CoapTransportService.java | 3 +- .../transport/http/DeviceApiController.java | 3 +- .../transport/mqtt/MqttTransportService.java | 4 +- 24 files changed, 312 insertions(+), 123 deletions(-) rename common/transport/src/main/java/org/thingsboard/server/common/transport/quota/{HostRequestsQuotaService.java => AbstractQuotaService.java} (50%) create mode 100644 common/transport/src/main/java/org/thingsboard/server/common/transport/quota/RequestLimitPolicy.java create mode 100644 common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryCleaner.java create mode 100644 common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryLogger.java create mode 100644 common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestIntervalRegistry.java rename common/transport/src/main/java/org/thingsboard/server/common/transport/quota/{ => host}/HostRequestLimitPolicy.java (73%) create mode 100644 common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestsQuotaService.java rename common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/{HostRequestIntervalRegistry.java => KeyBasedIntervalRegistry.java} (52%) create mode 100644 common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryCleaner.java create mode 100644 common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryLogger.java create mode 100644 common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantMsgsIntervalRegistry.java create mode 100644 common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantQuotaService.java create mode 100644 common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantRequestLimitPolicy.java diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java index a67278cc6d..a4558eb0a4 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultMsgQueueService.java @@ -23,6 +23,7 @@ import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.common.transport.quota.tenant.TenantQuotaService; import org.thingsboard.server.dao.queue.MsgQueue; import javax.annotation.PostConstruct; @@ -48,6 +49,9 @@ public class DefaultMsgQueueService implements MsgQueueService { @Autowired private MsgQueue msgQueue; + @Autowired + private TenantQuotaService quotaService; + private ScheduledExecutorService cleanupExecutor; private Map pendingCountPerTenant = new ConcurrentHashMap<>(); @@ -70,6 +74,11 @@ public class DefaultMsgQueueService implements MsgQueueService { @Override public ListenableFuture put(TenantId tenantId, TbMsg msg, UUID nodeId, long clusterPartition) { + if(quotaService.isQuotaExceeded(tenantId.getId().toString())) { + log.warn("Tenant TbMsg Quota exceeded for [{}:{}] . Reject", tenantId.getId()); + return Futures.immediateFailedFuture(new RuntimeException("Tenant TbMsg Quota exceeded")); + } + AtomicLong pendingMsgCount = pendingCountPerTenant.computeIfAbsent(tenantId, key -> new AtomicLong()); if (pendingMsgCount.incrementAndGet() < queueMaxSize) { return msgQueue.put(tenantId, msg, nodeId, clusterPartition); diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index dda6e1595f..63c5c4448d 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -129,9 +129,28 @@ quota: whitelist: "${QUOTA_HOST_WHITELIST:localhost,127.0.0.1}" # Array of blacklist hosts blacklist: "${QUOTA_HOST_BLACKLIST:}" - log: - topSize: 10 - intervalMin: 2 + log: + topSize: 10 + intervalMin: 2 + rule: + tenant: + # Max allowed number of API requests in interval for single tenant + limit: "${QUOTA_TENANT_LIMIT:100000}" + # Interval duration + intervalMs: "${QUOTA_TENANT_INTERVAL_MS:60000}" + # Maximum silence duration for tenant after which Tenant removed from QuotaService. Must be bigger than intervalMs + ttlMs: "${QUOTA_TENANT_TTL_MS:60000}" + # Interval for scheduled task that cleans expired records. TTL is used for expiring + cleanPeriodMs: "${QUOTA_TENANT_CLEAN_PERIOD_MS:300000}" + # Enable Host API Limits + enabled: "${QUOTA_TENANT_ENABLED:false}" + # Array of whitelist tenants + whitelist: "${QUOTA_TENANT_WHITELIST:}" + # Array of blacklist tenants + blacklist: "${QUOTA_HOST_BLACKLIST:}" + log: + topSize: 10 + intervalMin: 2 database: type: "${DATABASE_TYPE:sql}" # cassandra OR sql diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/HostRequestsQuotaService.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/AbstractQuotaService.java similarity index 50% rename from common/transport/src/main/java/org/thingsboard/server/common/transport/quota/HostRequestsQuotaService.java rename to common/transport/src/main/java/org/thingsboard/server/common/transport/quota/AbstractQuotaService.java index c1f045df8f..b2e32bacae 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/HostRequestsQuotaService.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/AbstractQuotaService.java @@ -1,47 +1,23 @@ -/** - * 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.common.transport.quota; -import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.stereotype.Service; -import org.thingsboard.server.common.transport.quota.inmemory.HostRequestIntervalRegistry; import org.thingsboard.server.common.transport.quota.inmemory.IntervalRegistryCleaner; import org.thingsboard.server.common.transport.quota.inmemory.IntervalRegistryLogger; +import org.thingsboard.server.common.transport.quota.inmemory.KeyBasedIntervalRegistry; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; -/** - * @author Vitaliy Paromskiy - * @version 1.0 - */ -@Service -@Slf4j -public class HostRequestsQuotaService implements QuotaService { +public class AbstractQuotaService implements QuotaService { - private final HostRequestIntervalRegistry requestRegistry; - private final HostRequestLimitPolicy requestsPolicy; + private final KeyBasedIntervalRegistry requestRegistry; + private final RequestLimitPolicy requestsPolicy; private final IntervalRegistryCleaner registryCleaner; private final IntervalRegistryLogger registryLogger; private final boolean enabled; - public HostRequestsQuotaService(HostRequestIntervalRegistry requestRegistry, HostRequestLimitPolicy requestsPolicy, + public AbstractQuotaService(KeyBasedIntervalRegistry requestRegistry, RequestLimitPolicy requestsPolicy, IntervalRegistryCleaner registryCleaner, IntervalRegistryLogger registryLogger, - @Value("${quota.host.enabled}") boolean enabled) { + boolean enabled) { this.requestRegistry = requestRegistry; this.requestsPolicy = requestsPolicy; this.registryCleaner = registryCleaner; diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/RequestLimitPolicy.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/RequestLimitPolicy.java new file mode 100644 index 0000000000..72b6b493a8 --- /dev/null +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/RequestLimitPolicy.java @@ -0,0 +1,15 @@ +package org.thingsboard.server.common.transport.quota; + + +public abstract class RequestLimitPolicy { + + private final long limit; + + public RequestLimitPolicy(long limit) { + this.limit = limit; + } + + public boolean isValid(long currentValue) { + return currentValue <= limit; + } +} diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryCleaner.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryCleaner.java new file mode 100644 index 0000000000..1afd88e529 --- /dev/null +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryCleaner.java @@ -0,0 +1,14 @@ +package org.thingsboard.server.common.transport.quota.host; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.transport.quota.inmemory.IntervalRegistryCleaner; + +@Component +public class HostIntervalRegistryCleaner extends IntervalRegistryCleaner { + + public HostIntervalRegistryCleaner(HostRequestIntervalRegistry intervalRegistry, + @Value("${quota.host.cleanPeriodMs}") long cleanPeriodMs) { + super(intervalRegistry, cleanPeriodMs); + } +} diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryLogger.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryLogger.java new file mode 100644 index 0000000000..09dee68d18 --- /dev/null +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryLogger.java @@ -0,0 +1,37 @@ +package org.thingsboard.server.common.transport.quota.host; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.transport.quota.inmemory.IntervalRegistryLogger; + +import java.util.Map; +import java.util.concurrent.TimeUnit; + +@Component +@Slf4j +public class HostIntervalRegistryLogger extends IntervalRegistryLogger { + + private final long logIntervalMin; + + public HostIntervalRegistryLogger(@Value("${quota.host.log.topSize}") int topSize, + @Value("${quota.host.log.intervalMin}") long logIntervalMin, + HostRequestIntervalRegistry intervalRegistry) { + super(topSize, logIntervalMin, intervalRegistry); + this.logIntervalMin = logIntervalMin; + } + + protected void log(Map top, int uniqHosts, long requestsCount) { + long rps = requestsCount / TimeUnit.MINUTES.toSeconds(logIntervalMin); + StringBuilder builder = new StringBuilder("Quota Statistic : "); + builder.append("uniqHosts : ").append(uniqHosts).append("; "); + builder.append("requestsCount : ").append(requestsCount).append("; "); + builder.append("RPS : ").append(rps).append(" "); + builder.append("top -> "); + for (Map.Entry host : top.entrySet()) { + builder.append(host.getKey()).append(" : ").append(host.getValue()).append("; "); + } + + log.info(builder.toString()); + } +} diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestIntervalRegistry.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestIntervalRegistry.java new file mode 100644 index 0000000000..3d588d5948 --- /dev/null +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestIntervalRegistry.java @@ -0,0 +1,37 @@ +/** + * 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.common.transport.quota.host; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.transport.quota.inmemory.KeyBasedIntervalRegistry; + +/** + * @author Vitaliy Paromskiy + * @version 1.0 + */ +@Component +@Slf4j +public class HostRequestIntervalRegistry extends KeyBasedIntervalRegistry { + + public HostRequestIntervalRegistry(@Value("${quota.host.intervalMs}") long intervalDurationMs, + @Value("${quota.host.ttlMs}") long ttlMs, + @Value("${quota.host.whitelist}") String whiteList, + @Value("${quota.host.blacklist}") String blackList) { + super(intervalDurationMs, ttlMs, whiteList, blackList, "host"); + } +} diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/HostRequestLimitPolicy.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestLimitPolicy.java similarity index 73% rename from common/transport/src/main/java/org/thingsboard/server/common/transport/quota/HostRequestLimitPolicy.java rename to common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestLimitPolicy.java index cf1c4e85ae..849e78fbbb 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/HostRequestLimitPolicy.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestLimitPolicy.java @@ -1,38 +1,33 @@ /** * 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 - * + *

+ * 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.common.transport.quota; +package org.thingsboard.server.common.transport.quota.host; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; +import org.thingsboard.server.common.transport.quota.RequestLimitPolicy; /** * @author Vitaliy Paromskiy * @version 1.0 */ @Component -public class HostRequestLimitPolicy { - - private final long limit; +public class HostRequestLimitPolicy extends RequestLimitPolicy { public HostRequestLimitPolicy(@Value("${quota.host.limit}") long limit) { - this.limit = limit; - } - - public boolean isValid(long currentValue) { - return currentValue <= limit; + super(limit); } } diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestsQuotaService.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestsQuotaService.java new file mode 100644 index 0000000000..0acd7b2b7d --- /dev/null +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestsQuotaService.java @@ -0,0 +1,37 @@ +/** + * 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.common.transport.quota.host; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Service; +import org.thingsboard.server.common.transport.quota.AbstractQuotaService; + +/** + * @author Vitaliy Paromskiy + * @version 1.0 + */ +@Service +@Slf4j +public class HostRequestsQuotaService extends AbstractQuotaService { + + public HostRequestsQuotaService(HostRequestIntervalRegistry requestRegistry, HostRequestLimitPolicy requestsPolicy, + HostIntervalRegistryCleaner registryCleaner, HostIntervalRegistryLogger registryLogger, + @Value("${quota.host.enabled}") boolean enabled) { + super(requestRegistry, requestsPolicy, registryCleaner, registryLogger, enabled); + } + +} diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryCleaner.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryCleaner.java index a227d2abba..0c510ff815 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryCleaner.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryCleaner.java @@ -16,10 +16,7 @@ package org.thingsboard.server.common.transport.quota.inmemory; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.stereotype.Component; -import javax.annotation.PreDestroy; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; @@ -28,15 +25,14 @@ import java.util.concurrent.TimeUnit; * @author Vitaliy Paromskiy * @version 1.0 */ -@Component @Slf4j -public class IntervalRegistryCleaner { +public abstract class IntervalRegistryCleaner { - private final HostRequestIntervalRegistry intervalRegistry; + private final KeyBasedIntervalRegistry intervalRegistry; private final long cleanPeriodMs; private ScheduledExecutorService executor; - public IntervalRegistryCleaner(HostRequestIntervalRegistry intervalRegistry, @Value("${quota.host.cleanPeriodMs}") long cleanPeriodMs) { + public IntervalRegistryCleaner(KeyBasedIntervalRegistry intervalRegistry, long cleanPeriodMs) { this.intervalRegistry = intervalRegistry; this.cleanPeriodMs = cleanPeriodMs; } diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLogger.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLogger.java index 8b34a6becc..1b02f02100 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLogger.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLogger.java @@ -1,12 +1,12 @@ /** * 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 - * + *

+ * 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. @@ -17,8 +17,6 @@ package org.thingsboard.server.common.transport.quota.inmemory; import com.google.common.collect.MinMaxPriorityQueue; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.stereotype.Component; import java.util.Comparator; import java.util.Map; @@ -32,17 +30,15 @@ import java.util.stream.Collectors; * @author Vitaliy Paromskiy * @version 1.0 */ -@Component @Slf4j -public class IntervalRegistryLogger { +public abstract class IntervalRegistryLogger { private final int topSize; - private final HostRequestIntervalRegistry intervalRegistry; + private final KeyBasedIntervalRegistry intervalRegistry; private final long logIntervalMin; private ScheduledExecutorService executor; - public IntervalRegistryLogger(@Value("${quota.log.topSize}") int topSize, @Value("${quota.log.intervalMin}") long logIntervalMin, - HostRequestIntervalRegistry intervalRegistry) { + public IntervalRegistryLogger(int topSize, long logIntervalMin, KeyBasedIntervalRegistry intervalRegistry) { this.topSize = topSize; this.logIntervalMin = logIntervalMin; this.intervalRegistry = intervalRegistry; @@ -79,17 +75,5 @@ public class IntervalRegistryLogger { return topQueue.stream().collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); } - private void log(Map top, int uniqHosts, long requestsCount) { - long rps = requestsCount / TimeUnit.MINUTES.toSeconds(logIntervalMin); - StringBuilder builder = new StringBuilder("Quota Statistic : "); - builder.append("uniqHosts : ").append(uniqHosts).append("; "); - builder.append("requestsCount : ").append(requestsCount).append("; "); - builder.append("RPS : ").append(rps).append(" "); - builder.append("top -> "); - for (Map.Entry host : top.entrySet()) { - builder.append(host.getKey()).append(" : ").append(host.getValue()).append("; "); - } - - log.info(builder.toString()); - } + protected abstract void log(Map top, int uniqHosts, long requestsCount); } 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/KeyBasedIntervalRegistry.java similarity index 52% rename from common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/HostRequestIntervalRegistry.java rename to common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/KeyBasedIntervalRegistry.java index 3782ed22ed..9884316ec5 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/KeyBasedIntervalRegistry.java @@ -1,39 +1,16 @@ -/** - * 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.common.transport.quota.inmemory; import com.google.common.collect.Sets; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.stereotype.Component; -import javax.annotation.PostConstruct; import java.util.Map; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.stream.Collectors; -/** - * @author Vitaliy Paromskiy - * @version 1.0 - */ -@Component @Slf4j -public class HostRequestIntervalRegistry { +public abstract class KeyBasedIntervalRegistry { private final Map hostCounts = new ConcurrentHashMap<>(); private final long intervalDurationMs; @@ -41,23 +18,20 @@ public class HostRequestIntervalRegistry { private final Set whiteList; private final Set blackList; - public HostRequestIntervalRegistry(@Value("${quota.host.intervalMs}") long intervalDurationMs, - @Value("${quota.host.ttlMs}") long ttlMs, - @Value("${quota.host.whitelist}") String whiteList, - @Value("${quota.host.blacklist}") String blackList) { + public KeyBasedIntervalRegistry(long intervalDurationMs, long ttlMs, String whiteList, String blackList, String name) { this.intervalDurationMs = intervalDurationMs; this.ttlMs = ttlMs; this.whiteList = Sets.newHashSet(StringUtils.split(whiteList, ',')); this.blackList = Sets.newHashSet(StringUtils.split(blackList, ',')); + validate(name); } - @PostConstruct - public void init() { + private void validate(String name) { if (ttlMs < intervalDurationMs) { - log.warn("TTL for IntervalRegistry [{}] smaller than interval duration [{}]", ttlMs, intervalDurationMs); + log.warn("TTL for {} IntervalRegistry [{}] smaller than interval duration [{}]", name, ttlMs, intervalDurationMs); } - log.info("Start Host Quota Service with whitelist {}", whiteList); - log.info("Start Host Quota Service with blacklist {}", blackList); + log.info("Start {} KeyBasedIntervalRegistry with whitelist {}", name, whiteList); + log.info("Start {} KeyBasedIntervalRegistry with blacklist {}", name, blackList); } public long tick(String clientHostId) { diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryCleaner.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryCleaner.java new file mode 100644 index 0000000000..0c2eef7fac --- /dev/null +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryCleaner.java @@ -0,0 +1,14 @@ +package org.thingsboard.server.common.transport.quota.tenant; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.transport.quota.inmemory.IntervalRegistryCleaner; + +@Component +public class TenantIntervalRegistryCleaner extends IntervalRegistryCleaner { + + public TenantIntervalRegistryCleaner(TenantMsgsIntervalRegistry intervalRegistry, + @Value("${quota.rule.tenant.cleanPeriodMs}") long cleanPeriodMs) { + super(intervalRegistry, cleanPeriodMs); + } +} diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryLogger.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryLogger.java new file mode 100644 index 0000000000..14729a6859 --- /dev/null +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryLogger.java @@ -0,0 +1,37 @@ +package org.thingsboard.server.common.transport.quota.tenant; + +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.transport.quota.inmemory.IntervalRegistryLogger; + +import java.util.Map; +import java.util.concurrent.TimeUnit; + +@Slf4j +@Component +public class TenantIntervalRegistryLogger extends IntervalRegistryLogger { + + private final long logIntervalMin; + + public TenantIntervalRegistryLogger(@Value("${quota.rule.tenant.log.topSize}") int topSize, + @Value("${quota.rule.tenant.log.intervalMin}") long logIntervalMin, + TenantMsgsIntervalRegistry intervalRegistry) { + super(topSize, logIntervalMin, intervalRegistry); + this.logIntervalMin = logIntervalMin; + } + + protected void log(Map top, int uniqHosts, long requestsCount) { + long rps = requestsCount / TimeUnit.MINUTES.toSeconds(logIntervalMin); + StringBuilder builder = new StringBuilder("Tenant Quota Statistic : "); + builder.append("uniqTenants : ").append(uniqHosts).append("; "); + builder.append("requestsCount : ").append(requestsCount).append("; "); + builder.append("RPS : ").append(rps).append(" "); + builder.append("top -> "); + for (Map.Entry host : top.entrySet()) { + builder.append(host.getKey()).append(" : ").append(host.getValue()).append("; "); + } + + log.info(builder.toString()); + } +} diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantMsgsIntervalRegistry.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantMsgsIntervalRegistry.java new file mode 100644 index 0000000000..2c875dec26 --- /dev/null +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantMsgsIntervalRegistry.java @@ -0,0 +1,16 @@ +package org.thingsboard.server.common.transport.quota.tenant; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.transport.quota.inmemory.KeyBasedIntervalRegistry; + +@Component +public class TenantMsgsIntervalRegistry extends KeyBasedIntervalRegistry { + + public TenantMsgsIntervalRegistry(@Value("${quota.rule.tenant.intervalMs}") long intervalDurationMs, + @Value("${quota.rule.tenant.ttlMs}") long ttlMs, + @Value("${quota.rule.tenant.whitelist}") String whiteList, + @Value("${quota.rule.tenant.blacklist}") String blackList) { + super(intervalDurationMs, ttlMs, whiteList, blackList, "Rule Tenant"); + } +} diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantQuotaService.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantQuotaService.java new file mode 100644 index 0000000000..fb06c70043 --- /dev/null +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantQuotaService.java @@ -0,0 +1,15 @@ +package org.thingsboard.server.common.transport.quota.tenant; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.transport.quota.AbstractQuotaService; + +@Component +public class TenantQuotaService extends AbstractQuotaService { + + public TenantQuotaService(TenantMsgsIntervalRegistry requestRegistry, TenantRequestLimitPolicy requestsPolicy, + TenantIntervalRegistryCleaner registryCleaner, TenantIntervalRegistryLogger registryLogger, + @Value("${quota.rule.tenant.enabled}") boolean enabled) { + super(requestRegistry, requestsPolicy, registryCleaner, registryLogger, enabled); + } +} diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantRequestLimitPolicy.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantRequestLimitPolicy.java new file mode 100644 index 0000000000..646906f9fe --- /dev/null +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantRequestLimitPolicy.java @@ -0,0 +1,13 @@ +package org.thingsboard.server.common.transport.quota.tenant; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.transport.quota.RequestLimitPolicy; + +@Component +public class TenantRequestLimitPolicy extends RequestLimitPolicy { + + public TenantRequestLimitPolicy(@Value("${quota.rule.tenant.limit}") long limit) { + super(limit); + } +} diff --git a/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/HostRequestLimitPolicyTest.java b/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/HostRequestLimitPolicyTest.java index 174d18235d..07e03ef0ac 100644 --- a/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/HostRequestLimitPolicyTest.java +++ b/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/HostRequestLimitPolicyTest.java @@ -16,6 +16,7 @@ package org.thingsboard.server.common.transport.quota; import org.junit.Test; +import org.thingsboard.server.common.transport.quota.host.HostRequestLimitPolicy; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; diff --git a/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/HostRequestsQuotaServiceTest.java b/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/HostRequestsQuotaServiceTest.java index 547f0cf0e4..20f8a55d0d 100644 --- a/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/HostRequestsQuotaServiceTest.java +++ b/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/HostRequestsQuotaServiceTest.java @@ -17,9 +17,7 @@ package org.thingsboard.server.common.transport.quota; import org.junit.Before; import org.junit.Test; -import org.thingsboard.server.common.transport.quota.inmemory.HostRequestIntervalRegistry; -import org.thingsboard.server.common.transport.quota.inmemory.IntervalRegistryCleaner; -import org.thingsboard.server.common.transport.quota.inmemory.IntervalRegistryLogger; +import org.thingsboard.server.common.transport.quota.host.*; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; @@ -35,8 +33,8 @@ public class HostRequestsQuotaServiceTest { private HostRequestIntervalRegistry requestRegistry = mock(HostRequestIntervalRegistry.class); private HostRequestLimitPolicy requestsPolicy = mock(HostRequestLimitPolicy.class); - private IntervalRegistryCleaner registryCleaner = mock(IntervalRegistryCleaner.class); - private IntervalRegistryLogger registryLogger = mock(IntervalRegistryLogger.class); + private HostIntervalRegistryCleaner registryCleaner = mock(HostIntervalRegistryCleaner.class); + private HostIntervalRegistryLogger registryLogger = mock(HostIntervalRegistryLogger.class); @Before public void init() { diff --git a/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/inmemory/HostRequestIntervalRegistryTest.java b/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/inmemory/HostRequestIntervalRegistryTest.java index 78b82eec08..b49dd00bc9 100644 --- a/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/inmemory/HostRequestIntervalRegistryTest.java +++ b/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/inmemory/HostRequestIntervalRegistryTest.java @@ -15,11 +15,9 @@ */ package org.thingsboard.server.common.transport.quota.inmemory; -import com.google.common.collect.Sets; import org.junit.Before; import org.junit.Test; - -import java.util.Collections; +import org.thingsboard.server.common.transport.quota.host.HostRequestIntervalRegistry; import static org.junit.Assert.assertEquals; diff --git a/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLoggerTest.java b/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLoggerTest.java index c9139aeeda..6e51420cd5 100644 --- a/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLoggerTest.java +++ b/common/transport/src/test/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLoggerTest.java @@ -18,6 +18,8 @@ package org.thingsboard.server.common.transport.quota.inmemory; import com.google.common.collect.ImmutableMap; import org.junit.Before; import org.junit.Test; +import org.thingsboard.server.common.transport.quota.host.HostIntervalRegistryLogger; +import org.thingsboard.server.common.transport.quota.host.HostRequestIntervalRegistry; import java.util.Collections; import java.util.Map; @@ -37,7 +39,7 @@ public class IntervalRegistryLoggerTest { @Before public void init() { - logger = new IntervalRegistryLogger(3, 10, requestRegistry); + logger = new HostIntervalRegistryLogger(3, 10, requestRegistry); } @Test diff --git a/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportService.java b/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportService.java index 15706d4d6f..6c8437cb03 100644 --- a/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportService.java +++ b/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportService.java @@ -27,6 +27,7 @@ import org.springframework.stereotype.Service; import org.thingsboard.server.common.transport.SessionMsgProcessor; import org.thingsboard.server.common.transport.auth.DeviceAuthService; import org.thingsboard.server.common.transport.quota.QuotaService; +import org.thingsboard.server.common.transport.quota.host.HostRequestsQuotaService; import org.thingsboard.server.transport.coap.adaptors.CoapTransportAdaptor; import javax.annotation.PostConstruct; @@ -55,7 +56,7 @@ public class CoapTransportService { private DeviceAuthService authService; @Autowired(required = false) - private QuotaService quotaService; + private HostRequestsQuotaService quotaService; @Value("${coap.bind_address}") diff --git a/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java b/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java index 4d90b5f21c..4ac9799e0b 100644 --- a/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java +++ b/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java @@ -36,6 +36,7 @@ import org.thingsboard.server.common.transport.SessionMsgProcessor; import org.thingsboard.server.common.transport.adaptor.JsonConverter; import org.thingsboard.server.common.transport.auth.DeviceAuthService; import org.thingsboard.server.common.transport.quota.QuotaService; +import org.thingsboard.server.common.transport.quota.host.HostRequestsQuotaService; import org.thingsboard.server.transport.http.session.HttpSessionCtx; import javax.servlet.http.HttpServletRequest; @@ -61,7 +62,7 @@ public class DeviceApiController { private DeviceAuthService authService; @Autowired(required = false) - private QuotaService quotaService; + private HostRequestsQuotaService quotaService; @RequestMapping(value = "/{deviceToken}/attributes", method = RequestMethod.GET, produces = "application/json") public DeferredResult getDeviceAttributes(@PathVariable("deviceToken") String deviceToken, diff --git a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java index f0129e1b6c..1b37ed4513 100644 --- a/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java +++ b/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportService.java @@ -29,7 +29,7 @@ import org.springframework.context.ApplicationContext; import org.springframework.stereotype.Service; import org.thingsboard.server.common.transport.SessionMsgProcessor; import org.thingsboard.server.common.transport.auth.DeviceAuthService; -import org.thingsboard.server.common.transport.quota.QuotaService; +import org.thingsboard.server.common.transport.quota.host.HostRequestsQuotaService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.transport.mqtt.adaptors.MqttTransportAdaptor; @@ -67,7 +67,7 @@ public class MqttTransportService { private MqttSslHandlerProvider sslHandlerProvider; @Autowired(required = false) - private QuotaService quotaService; + private HostRequestsQuotaService quotaService; @Value("${mqtt.bind_address}") private String host; From 94d07fb7b1b5fd03a410dd913230a55d3671bf91 Mon Sep 17 00:00:00 2001 From: vparomskiy Date: Thu, 17 May 2018 09:59:03 +0300 Subject: [PATCH 2/5] fix lincense headers --- .../transport/quota/AbstractQuotaService.java | 15 +++++++++++++++ .../transport/quota/RequestLimitPolicy.java | 15 +++++++++++++++ .../quota/host/HostIntervalRegistryCleaner.java | 15 +++++++++++++++ .../quota/host/HostIntervalRegistryLogger.java | 15 +++++++++++++++ .../quota/host/HostRequestIntervalRegistry.java | 8 ++++---- .../quota/host/HostRequestLimitPolicy.java | 8 ++++---- .../quota/host/HostRequestsQuotaService.java | 8 ++++---- .../quota/inmemory/IntervalRegistryLogger.java | 8 ++++---- .../quota/inmemory/KeyBasedIntervalRegistry.java | 15 +++++++++++++++ .../tenant/TenantIntervalRegistryCleaner.java | 15 +++++++++++++++ .../tenant/TenantIntervalRegistryLogger.java | 15 +++++++++++++++ .../quota/tenant/TenantMsgsIntervalRegistry.java | 15 +++++++++++++++ .../quota/tenant/TenantQuotaService.java | 15 +++++++++++++++ .../quota/tenant/TenantRequestLimitPolicy.java | 15 +++++++++++++++ 14 files changed, 166 insertions(+), 16 deletions(-) diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/AbstractQuotaService.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/AbstractQuotaService.java index b2e32bacae..9d7bd1720c 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/AbstractQuotaService.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/AbstractQuotaService.java @@ -1,3 +1,18 @@ +/** + * 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.common.transport.quota; import org.thingsboard.server.common.transport.quota.inmemory.IntervalRegistryCleaner; diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/RequestLimitPolicy.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/RequestLimitPolicy.java index 72b6b493a8..0ff1230b1c 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/RequestLimitPolicy.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/RequestLimitPolicy.java @@ -1,3 +1,18 @@ +/** + * 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.common.transport.quota; diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryCleaner.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryCleaner.java index 1afd88e529..ea75f5dc09 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryCleaner.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryCleaner.java @@ -1,3 +1,18 @@ +/** + * 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.common.transport.quota.host; import org.springframework.beans.factory.annotation.Value; diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryLogger.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryLogger.java index 09dee68d18..65767f1648 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryLogger.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostIntervalRegistryLogger.java @@ -1,3 +1,18 @@ +/** + * 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.common.transport.quota.host; import lombok.extern.slf4j.Slf4j; diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestIntervalRegistry.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestIntervalRegistry.java index 3d588d5948..9b3b4614e4 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestIntervalRegistry.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestIntervalRegistry.java @@ -1,12 +1,12 @@ /** * 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 - *

+ * + * 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. diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestLimitPolicy.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestLimitPolicy.java index 849e78fbbb..eeef9247e6 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestLimitPolicy.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestLimitPolicy.java @@ -1,12 +1,12 @@ /** * 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 - *

+ * + * 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. diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestsQuotaService.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestsQuotaService.java index 0acd7b2b7d..69342b538b 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestsQuotaService.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/host/HostRequestsQuotaService.java @@ -1,12 +1,12 @@ /** * 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 - *

+ * + * 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. diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLogger.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLogger.java index 1b02f02100..30399a140f 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLogger.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/IntervalRegistryLogger.java @@ -1,12 +1,12 @@ /** * 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 - *

+ * + * 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. diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/KeyBasedIntervalRegistry.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/KeyBasedIntervalRegistry.java index 9884316ec5..0b0fef85e5 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/KeyBasedIntervalRegistry.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/inmemory/KeyBasedIntervalRegistry.java @@ -1,3 +1,18 @@ +/** + * 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.common.transport.quota.inmemory; import com.google.common.collect.Sets; diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryCleaner.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryCleaner.java index 0c2eef7fac..c48117069a 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryCleaner.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryCleaner.java @@ -1,3 +1,18 @@ +/** + * 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.common.transport.quota.tenant; import org.springframework.beans.factory.annotation.Value; diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryLogger.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryLogger.java index 14729a6859..c56f457423 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryLogger.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantIntervalRegistryLogger.java @@ -1,3 +1,18 @@ +/** + * 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.common.transport.quota.tenant; import lombok.extern.slf4j.Slf4j; diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantMsgsIntervalRegistry.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantMsgsIntervalRegistry.java index 2c875dec26..6e8402c701 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantMsgsIntervalRegistry.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantMsgsIntervalRegistry.java @@ -1,3 +1,18 @@ +/** + * 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.common.transport.quota.tenant; import org.springframework.beans.factory.annotation.Value; diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantQuotaService.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantQuotaService.java index fb06c70043..a68860a2e0 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantQuotaService.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantQuotaService.java @@ -1,3 +1,18 @@ +/** + * 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.common.transport.quota.tenant; import org.springframework.beans.factory.annotation.Value; diff --git a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantRequestLimitPolicy.java b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantRequestLimitPolicy.java index 646906f9fe..cc32c81ad7 100644 --- a/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantRequestLimitPolicy.java +++ b/common/transport/src/main/java/org/thingsboard/server/common/transport/quota/tenant/TenantRequestLimitPolicy.java @@ -1,3 +1,18 @@ +/** + * 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.common.transport.quota.tenant; import org.springframework.beans.factory.annotation.Value; From 084907dfc4944f218c8f54c3e5c5efbf6c0bc13d Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Thu, 17 May 2018 10:20:08 +0300 Subject: [PATCH 3/5] DB Msg queue refactor. --- .../thingsboard/server/dao/queue/db/{nosql => }/MsgAck.java | 2 +- .../server/dao/queue/db/{nosql => }/UnprocessedMsgFilter.java | 3 ++- .../server/dao/queue/db/nosql/CassandraMsgQueue.java | 2 ++ .../dao/queue/db/nosql/repository/CassandraAckRepository.java | 2 +- .../server/dao/queue/db/repository/AckRepository.java | 2 +- .../server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java | 4 ++-- .../queue/db/nosql/repository/CassandraAckRepositoryTest.java | 2 +- dao/src/test/resources/application-test.properties | 1 - 8 files changed, 10 insertions(+), 8 deletions(-) rename dao/src/main/java/org/thingsboard/server/dao/queue/db/{nosql => }/MsgAck.java (94%) rename dao/src/main/java/org/thingsboard/server/dao/queue/db/{nosql => }/UnprocessedMsgFilter.java (92%) diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/MsgAck.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java similarity index 94% rename from dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/MsgAck.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java index 1b1cd3f467..a1b039a1a7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/MsgAck.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/MsgAck.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.queue.db.nosql; +package org.thingsboard.server.dao.queue.db; import lombok.Data; import lombok.EqualsAndHashCode; diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilter.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java similarity index 92% rename from dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilter.java rename to dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java index c912e8e114..66eaa6d46b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilter.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/UnprocessedMsgFilter.java @@ -13,10 +13,11 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.dao.queue.db.nosql; +package org.thingsboard.server.dao.queue.db; import org.springframework.stereotype.Component; import org.thingsboard.server.common.msg.TbMsg; +import org.thingsboard.server.dao.queue.db.MsgAck; import java.util.Collection; import java.util.List; diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java index ceaab586cf..ce481b78b8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/CassandraMsgQueue.java @@ -26,6 +26,8 @@ import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.dao.queue.MsgQueue; +import org.thingsboard.server.dao.queue.db.MsgAck; +import org.thingsboard.server.dao.queue.db.UnprocessedMsgFilter; import org.thingsboard.server.dao.queue.db.repository.AckRepository; import org.thingsboard.server.dao.queue.db.repository.MsgRepository; import org.thingsboard.server.dao.util.NoSqlDao; diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java index 1ffbec3d81..ab7de0415b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepository.java @@ -22,7 +22,7 @@ import com.google.common.util.concurrent.ListenableFuture; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import org.thingsboard.server.dao.nosql.CassandraAbstractDao; -import org.thingsboard.server.dao.queue.db.nosql.MsgAck; +import org.thingsboard.server.dao.queue.db.MsgAck; import org.thingsboard.server.dao.queue.db.repository.AckRepository; import org.thingsboard.server.dao.util.NoSqlDao; diff --git a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java index 458dba81cb..6fbd2da57e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/queue/db/repository/AckRepository.java @@ -16,7 +16,7 @@ package org.thingsboard.server.dao.queue.db.repository; import com.google.common.util.concurrent.ListenableFuture; -import org.thingsboard.server.dao.queue.db.nosql.MsgAck; +import org.thingsboard.server.dao.queue.db.MsgAck; import java.util.List; import java.util.UUID; diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java index 39a432d244..fd9bf21164 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/UnprocessedMsgFilterTest.java @@ -18,8 +18,8 @@ package org.thingsboard.server.dao.queue.db.nosql; import com.google.common.collect.Lists; import org.junit.Test; import org.thingsboard.server.common.msg.TbMsg; -import org.thingsboard.server.dao.queue.db.nosql.MsgAck; -import org.thingsboard.server.dao.queue.db.nosql.UnprocessedMsgFilter; +import org.thingsboard.server.dao.queue.db.MsgAck; +import org.thingsboard.server.dao.queue.db.UnprocessedMsgFilter; import java.util.Collection; import java.util.List; diff --git a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java index f2b8c88f01..b2f38dc539 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/queue/db/nosql/repository/CassandraAckRepositoryTest.java @@ -23,7 +23,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.test.util.ReflectionTestUtils; import org.thingsboard.server.dao.service.AbstractServiceTest; import org.thingsboard.server.dao.service.DaoNoSqlTest; -import org.thingsboard.server.dao.queue.db.nosql.MsgAck; +import org.thingsboard.server.dao.queue.db.MsgAck; import java.util.List; import java.util.UUID; diff --git a/dao/src/test/resources/application-test.properties b/dao/src/test/resources/application-test.properties index dbd8b84a4b..f2dab45491 100644 --- a/dao/src/test/resources/application-test.properties +++ b/dao/src/test/resources/application-test.properties @@ -30,4 +30,3 @@ redis.connection.db=0 redis.connection.password= rule.queue.type=memory -rule.queue.max_size=10000 \ No newline at end of file From 199f108905e86c84a6ce5a1f02f0c55f3a9d15a8 Mon Sep 17 00:00:00 2001 From: vparomskiy Date: Thu, 17 May 2018 10:36:28 +0300 Subject: [PATCH 4/5] fix test --- .../thingsboard/server/transport/coap/CoapServerTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/transport/coap/src/test/java/org/thingsboard/server/transport/coap/CoapServerTest.java b/transport/coap/src/test/java/org/thingsboard/server/transport/coap/CoapServerTest.java index 320f06e7f1..48150567a0 100644 --- a/transport/coap/src/test/java/org/thingsboard/server/transport/coap/CoapServerTest.java +++ b/transport/coap/src/test/java/org/thingsboard/server/transport/coap/CoapServerTest.java @@ -50,7 +50,7 @@ import org.thingsboard.server.common.msg.session.*; import org.thingsboard.server.common.transport.SessionMsgProcessor; import org.thingsboard.server.common.transport.auth.DeviceAuthResult; import org.thingsboard.server.common.transport.auth.DeviceAuthService; -import org.thingsboard.server.common.transport.quota.QuotaService; +import org.thingsboard.server.common.transport.quota.host.HostRequestsQuotaService; import java.util.ArrayList; import java.util.List; @@ -134,8 +134,8 @@ public class CoapServerTest { } @Bean - public static QuotaService quotaService() { - return key -> false; + public static HostRequestsQuotaService quotaService() { + return new HostRequestsQuotaService(null, null, null, null, false); } } From 9175aba24fa2ffded90de35e78939e889a5c1a07 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Thu, 17 May 2018 13:45:06 +0300 Subject: [PATCH 5/5] Improve Relation Service --- .../dao/relation/BaseRelationService.java | 96 +++++++------------ .../server/dao/relation/RelationService.java | 5 +- .../dao/sql/relation/JpaRelationDao.java | 36 +++---- 3 files changed, 53 insertions(+), 84 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java b/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java index 7b2b391b10..ac5bfb29b5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.relation; import com.google.common.base.Function; +import com.google.common.util.concurrent.AsyncFunction; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; @@ -176,97 +177,64 @@ public class BaseRelationService implements RelationService { } @Override - public boolean deleteEntityRelations(EntityId entity) { + public void deleteEntityRelations(EntityId entityId) throws ExecutionException, InterruptedException { + deleteEntityRelationsAsync(entityId).get(); + } + + @Override + public ListenableFuture deleteEntityRelationsAsync(EntityId entityId) { Cache cache = cacheManager.getCache(RELATIONS_CACHE); - log.trace("Executing deleteEntityRelations [{}]", entity); - validate(entity); + log.trace("Executing deleteEntityRelationsAsync [{}]", entityId); + validate(entityId); List>> inboundRelationsList = new ArrayList<>(); for (RelationTypeGroup typeGroup : RelationTypeGroup.values()) { - inboundRelationsList.add(relationDao.findAllByTo(entity, typeGroup)); + inboundRelationsList.add(relationDao.findAllByTo(entityId, typeGroup)); } - ListenableFuture>> inboundRelations = Futures.allAsList(inboundRelationsList); - ListenableFuture> inboundDeletions = Futures.transform(inboundRelations, relations -> - getBooleans(relations, cache, true)); - ListenableFuture inboundFuture = Futures.transform(inboundDeletions, getListToBooleanFunction()); - boolean inboundDeleteResult = false; - try { - inboundDeleteResult = inboundFuture.get(); - } catch (InterruptedException | ExecutionException e) { - log.error("Error deleting entity inbound relations", e); - } + ListenableFuture>> inboundRelations = Futures.allAsList(inboundRelationsList); List>> outboundRelationsList = new ArrayList<>(); for (RelationTypeGroup typeGroup : RelationTypeGroup.values()) { - outboundRelationsList.add(relationDao.findAllByFrom(entity, typeGroup)); + outboundRelationsList.add(relationDao.findAllByFrom(entityId, typeGroup)); } - ListenableFuture>> outboundRelations = Futures.allAsList(outboundRelationsList); - Futures.transform(outboundRelations, relations -> getBooleans(relations, cache, false)); - boolean outboundDeleteResult = relationDao.deleteOutboundRelations(entity); - return inboundDeleteResult && outboundDeleteResult; - } - - private List getBooleans(List> relations, Cache cache, boolean isRemove) { - List results = new ArrayList<>(); - for (List relationList : relations) { - relationList.forEach(relation -> checkFromDeleteSync(cache, results, relation, isRemove)); - } - return results; - } - - private void checkFromDeleteSync(Cache cache, List results, EntityRelation relation, boolean isRemove) { - if (isRemove) { - results.add(relationDao.deleteRelation(relation)); - } - cacheEviction(relation, cache); - } + ListenableFuture>> outboundRelations = Futures.allAsList(outboundRelationsList); - @Override - public ListenableFuture deleteEntityRelationsAsync(EntityId entity) { - Cache cache = cacheManager.getCache(RELATIONS_CACHE); - log.trace("Executing deleteEntityRelationsAsync [{}]", entity); - validate(entity); - List>> inboundRelationsList = new ArrayList<>(); - for (RelationTypeGroup typeGroup : RelationTypeGroup.values()) { - inboundRelationsList.add(relationDao.findAllByTo(entity, typeGroup)); - } - ListenableFuture>> inboundRelations = Futures.allAsList(inboundRelationsList); ListenableFuture> inboundDeletions = Futures.transformAsync(inboundRelations, relations -> { - List> results = getListenableFutures(relations, cache, true); + List> results = deleteRelationGroupsAsync(relations, cache, true); return Futures.allAsList(results); }); - ListenableFuture inboundFuture = Futures.transform(inboundDeletions, getListToBooleanFunction()); + ListenableFuture> outboundDeletions = Futures.transformAsync(outboundRelations, + relations -> { + List> results = deleteRelationGroupsAsync(relations, cache, false); + return Futures.allAsList(results); + }); - List>> outboundRelationsList = new ArrayList<>(); - for (RelationTypeGroup typeGroup : RelationTypeGroup.values()) { - outboundRelationsList.add(relationDao.findAllByFrom(entity, typeGroup)); - } - ListenableFuture>> outboundRelations = Futures.allAsList(outboundRelationsList); - Futures.transformAsync(outboundRelations, relations -> { - List> results = getListenableFutures(relations, cache, false); - return Futures.allAsList(results); - }); + ListenableFuture>> deletionsFuture = Futures.allAsList(inboundDeletions, outboundDeletions); - ListenableFuture outboundFuture = relationDao.deleteOutboundRelationsAsync(entity); - return Futures.transform(Futures.allAsList(Arrays.asList(inboundFuture, outboundFuture)), getListToBooleanFunction()); + return Futures.transformAsync(deletionsFuture, (deletions) -> { + relationDao.deleteOutboundRelationsAsync(entityId); + return null; + }); } - private List> getListenableFutures(List> relations, Cache cache, boolean isRemove) { + private List> deleteRelationGroupsAsync(List> relations, Cache cache, boolean deleteFromDb) { List> results = new ArrayList<>(); for (List relationList : relations) { - relationList.forEach(relation -> checkFromDeleteAsync(cache, results, relation, isRemove)); + relationList.forEach(relation -> results.add(deleteAsync(cache, relation, deleteFromDb))); } return results; } - private void checkFromDeleteAsync(Cache cache, List> results, EntityRelation relation, boolean isRemove) { - if (isRemove) { - results.add(relationDao.deleteRelationAsync(relation)); - } + private ListenableFuture deleteAsync(Cache cache, EntityRelation relation, boolean deleteFromDb) { cacheEviction(relation, cache); + if (deleteFromDb) { + return relationDao.deleteRelationAsync(relation); + } else { + return Futures.immediateFuture(false); + } } private void cacheEviction(EntityRelation relation, Cache cache) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/relation/RelationService.java b/dao/src/main/java/org/thingsboard/server/dao/relation/RelationService.java index ca1b959c95..bd945f1586 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/relation/RelationService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/relation/RelationService.java @@ -23,6 +23,7 @@ import org.thingsboard.server.common.data.relation.EntityRelationsQuery; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import java.util.List; +import java.util.concurrent.ExecutionException; /** * Created by ashvayka on 27.04.17. @@ -47,9 +48,9 @@ public interface RelationService { ListenableFuture deleteRelationAsync(EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup); - boolean deleteEntityRelations(EntityId entity); + void deleteEntityRelations(EntityId entity) throws ExecutionException, InterruptedException; - ListenableFuture deleteEntityRelationsAsync(EntityId entity); + ListenableFuture deleteEntityRelationsAsync(EntityId entity); List findByFrom(EntityId from, RelationTypeGroup typeGroup); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java index 6776f562bc..2a25bb0eb4 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java @@ -132,39 +132,35 @@ public class JpaRelationDao extends JpaAbstractDaoListeningExecutorService imple @Override public boolean deleteRelation(EntityRelation relation) { RelationCompositeKey key = new RelationCompositeKey(relation); - boolean relationExistsBeforeDelete = relationRepository.exists(key); - relationRepository.delete(key); - return relationExistsBeforeDelete; + return deleteRelationIfExists(key); } @Override public ListenableFuture deleteRelationAsync(EntityRelation relation) { RelationCompositeKey key = new RelationCompositeKey(relation); return service.submit( - () -> { - boolean relationExistsBeforeDelete = relationRepository.exists(key); - relationRepository.delete(key); - return relationExistsBeforeDelete; - }); + () -> deleteRelationIfExists(key)); } @Override public boolean deleteRelation(EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup) { RelationCompositeKey key = getRelationCompositeKey(from, to, relationType, typeGroup); - boolean relationExistsBeforeDelete = relationRepository.exists(key); - relationRepository.delete(key); - return relationExistsBeforeDelete; + return deleteRelationIfExists(key); } @Override public ListenableFuture deleteRelationAsync(EntityId from, EntityId to, String relationType, RelationTypeGroup typeGroup) { RelationCompositeKey key = getRelationCompositeKey(from, to, relationType, typeGroup); return service.submit( - () -> { - boolean relationExistsBeforeDelete = relationRepository.exists(key); - relationRepository.delete(key); - return relationExistsBeforeDelete; - }); + () -> deleteRelationIfExists(key)); + } + + private boolean deleteRelationIfExists(RelationCompositeKey key) { + boolean relationExistsBeforeDelete = relationRepository.exists(key); + if (relationExistsBeforeDelete) { + relationRepository.delete(key); + } + return relationExistsBeforeDelete; } @Override @@ -172,7 +168,9 @@ public class JpaRelationDao extends JpaAbstractDaoListeningExecutorService imple boolean relationExistsBeforeDelete = relationRepository .findAllByFromIdAndFromType(UUIDConverter.fromTimeUUID(entity.getId()), entity.getEntityType().name()) .size() > 0; - relationRepository.deleteByFromIdAndFromType(UUIDConverter.fromTimeUUID(entity.getId()), entity.getEntityType().name()); + if (relationExistsBeforeDelete) { + relationRepository.deleteByFromIdAndFromType(UUIDConverter.fromTimeUUID(entity.getId()), entity.getEntityType().name()); + } return relationExistsBeforeDelete; } @@ -183,7 +181,9 @@ public class JpaRelationDao extends JpaAbstractDaoListeningExecutorService imple boolean relationExistsBeforeDelete = relationRepository .findAllByFromIdAndFromType(UUIDConverter.fromTimeUUID(entity.getId()), entity.getEntityType().name()) .size() > 0; - relationRepository.deleteByFromIdAndFromType(UUIDConverter.fromTimeUUID(entity.getId()), entity.getEntityType().name()); + if (relationExistsBeforeDelete) { + relationRepository.deleteByFromIdAndFromType(UUIDConverter.fromTimeUUID(entity.getId()), entity.getEntityType().name()); + } return relationExistsBeforeDelete; }); }