diff --git a/application/pom.xml b/application/pom.xml
index 55009709d1..4792dbf93a 100644
--- a/application/pom.xml
+++ b/application/pom.xml
@@ -193,8 +193,8 @@
logback-classic
- javax.mail
- mail
+ com.sun.mail
+ javax.mail
org.apache.curator
diff --git a/application/src/main/data/json/demo/dashboards/gateways.json b/application/src/main/data/json/demo/dashboards/gateways.json
index f6370696c1..7d75908229 100644
--- a/application/src/main/data/json/demo/dashboards/gateways.json
+++ b/application/src/main/data/json/demo/dashboards/gateways.json
@@ -5,7 +5,7 @@
"94715984-ae74-76e4-20b7-2f956b01ed80": {
"isSystemType": true,
"bundleAlias": "entity_admin_widgets",
- "typeAlias": "device_admin_table2",
+ "typeAlias": "device_admin_table",
"type": "latest",
"title": "New widget",
"sizeX": 24,
@@ -1271,4 +1271,4 @@
}
},
"name": "Gateways"
-}
\ No newline at end of file
+}
diff --git a/application/src/main/data/json/system/widget_bundles/entity_admin_widgets.json b/application/src/main/data/json/system/widget_bundles/entity_admin_widgets.json
index c0d0fa7cdf..b65445f3ca 100644
--- a/application/src/main/data/json/system/widget_bundles/entity_admin_widgets.json
+++ b/application/src/main/data/json/system/widget_bundles/entity_admin_widgets.json
@@ -6,7 +6,7 @@
},
"widgetTypes": [
{
- "alias": "device_admin_table2",
+ "alias": "device_admin_table",
"name": "Device admin table",
"descriptor": {
"type": "latest",
@@ -22,7 +22,7 @@
}
},
{
- "alias": "device_admin_table",
+ "alias": "asset_admin_table",
"name": "Asset admin table",
"descriptor": {
"type": "latest",
diff --git a/application/src/main/java/org/thingsboard/server/controller/AuthController.java b/application/src/main/java/org/thingsboard/server/controller/AuthController.java
index 3097b9c2b1..49cfe482df 100644
--- a/application/src/main/java/org/thingsboard/server/controller/AuthController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/AuthController.java
@@ -338,7 +338,7 @@ public class AuthController extends BaseController {
@RequestMapping(value = "/noauth/oauth2Clients", method = RequestMethod.POST)
@ResponseBody
- public List getOath2Clients() throws ThingsboardException {
+ public List getOAuth2Clients() throws ThingsboardException {
try {
return oauth2Service.getOAuth2Clients();
} catch (Exception e) {
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
index a5abd77ad5..b7f6bd3dd2 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
@@ -23,6 +23,7 @@ import org.springframework.stereotype.Service;
import org.thingsboard.rule.engine.api.RpcError;
import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.common.data.id.TenantId;
+import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.TbActorMsg;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback;
@@ -45,6 +46,7 @@ import org.thingsboard.server.service.encoding.DataDecodingEncodingService;
import org.thingsboard.server.service.queue.processing.AbstractConsumerService;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponse;
import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService;
+import org.thingsboard.server.service.rpc.ToDeviceRpcRequestActorMsg;
import org.thingsboard.server.service.state.DeviceStateService;
import org.thingsboard.server.service.subscription.SubscriptionManagerService;
import org.thingsboard.server.service.subscription.TbLocalSubscriptionService;
@@ -100,7 +102,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService actorMsg = encodingService.decode(toCoreMsg.getToDeviceActorNotificationMsg().toByteArray());
if (actorMsg.isPresent()) {
- log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get());
- actorContext.tell(actorMsg.get(), ActorRef.noSender());
+ TbActorMsg tbActorMsg = actorMsg.get();
+ if (tbActorMsg.getMsgType().equals(MsgType.DEVICE_RPC_REQUEST_TO_DEVICE_ACTOR_MSG)) {
+ tbCoreDeviceRpcService.forwardRpcRequestToDeviceActor((ToDeviceRpcRequestActorMsg) tbActorMsg);
+ } else {
+ log.trace("[{}] Forwarding message to App Actor {}", id, actorMsg.get());
+ actorContext.tell(actorMsg.get(), ActorRef.noSender());
+ }
}
callback.onSuccess();
}
diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
index 0c48c54412..182a62d397 100644
--- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
+++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
@@ -21,6 +21,7 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service;
+import org.thingsboard.rule.engine.api.RpcError;
import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.TbActorMsg;
@@ -31,6 +32,7 @@ import org.thingsboard.server.common.msg.queue.ServiceQueue;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TbMsgCallback;
+import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
import org.thingsboard.server.queue.TbQueueConsumer;
@@ -48,6 +50,9 @@ import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingStr
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingStrategyFactory;
import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategy;
import org.thingsboard.server.service.queue.processing.TbRuleEngineSubmitStrategyFactory;
+import org.thingsboard.server.service.rpc.FromDeviceRpcResponse;
+import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService;
+import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService;
import org.thingsboard.server.service.stats.RuleEngineStatisticsService;
import javax.annotation.PostConstruct;
@@ -81,6 +86,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
private final TbRuleEngineQueueFactory tbRuleEngineQueueFactory;
private final TbQueueRuleEngineSettings ruleEngineSettings;
private final RuleEngineStatisticsService statisticsService;
+ private final TbRuleEngineDeviceRpcService tbDeviceRpcService;
private final ConcurrentMap>> consumers = new ConcurrentHashMap<>();
private final ConcurrentMap consumerConfigurations = new ConcurrentHashMap<>();
private final ConcurrentMap consumerStats = new ConcurrentHashMap<>();
@@ -90,13 +96,15 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
TbRuleEngineSubmitStrategyFactory submitStrategyFactory,
TbQueueRuleEngineSettings ruleEngineSettings,
TbRuleEngineQueueFactory tbRuleEngineQueueFactory, RuleEngineStatisticsService statisticsService,
- ActorSystemContext actorContext, DataDecodingEncodingService encodingService) {
+ ActorSystemContext actorContext, DataDecodingEncodingService encodingService,
+ TbRuleEngineDeviceRpcService tbDeviceRpcService) {
super(actorContext, encodingService, tbRuleEngineQueueFactory.createToRuleEngineNotificationsMsgConsumer());
this.statisticsService = statisticsService;
this.ruleEngineSettings = ruleEngineSettings;
this.tbRuleEngineQueueFactory = tbRuleEngineQueueFactory;
this.submitStrategyFactory = submitStrategyFactory;
this.processingStrategyFactory = processingStrategyFactory;
+ this.tbDeviceRpcService = tbDeviceRpcService;
}
@PostConstruct
@@ -227,7 +235,15 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
actorContext.tell(actorMsg.get(), ActorRef.noSender());
}
callback.onSuccess();
+ } else if (nfMsg.hasFromDeviceRpcResponse()) {
+ TransportProtos.FromDeviceRPCResponseProto proto = nfMsg.getFromDeviceRpcResponse();
+ RpcError error = proto.getError() > 0 ? RpcError.values()[proto.getError()] : null;
+ FromDeviceRpcResponse response = new FromDeviceRpcResponse(new UUID(proto.getRequestIdMSB(), proto.getRequestIdLSB())
+ , proto.getResponse(), error);
+ tbDeviceRpcService.processRpcResponseFromDevice(response);
+ callback.onSuccess();
} else {
+ log.trace("Received notification with missing handler");
callback.onSuccess();
}
}
diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml
index 254162b56b..03e461015f 100644
--- a/application/src/main/resources/thingsboard.yml
+++ b/application/src/main/resources/thingsboard.yml
@@ -621,8 +621,7 @@ queue:
notifications: "${TB_QUEUE_RABBIT_MQ_NOTIFICATIONS_QUEUE_PROPERTIES:x-max-length-bytes:1048576000;x-message-ttl:604800000}"
js-executor: "${TB_QUEUE_RABBIT_MQ_JE_QUEUE_PROPERTIES:x-max-length-bytes:1048576000;x-message-ttl:604800000}"
partitions:
- hash_function_name: "${TB_QUEUE_PARTITIONS_HASH_FUNCTION_NAME:murmur3_128}"
- virtual_nodes_size: "${TB_QUEUE_PARTITIONS_VIRTUAL_NODES_SIZE:16}"
+ hash_function_name: "${TB_QUEUE_PARTITIONS_HASH_FUNCTION_NAME:murmur3_128}" # murmur3_32, murmur3_128 or sha256
transport_api:
requests_topic: "${TB_QUEUE_TRANSPORT_API_REQUEST_TOPIC:tb_transport.api.requests}"
responses_topic: "${TB_QUEUE_TRANSPORT_API_RESPONSE_TOPIC:tb_transport.api.responses}"
@@ -638,7 +637,7 @@ queue:
pack-processing-timeout: "${TB_QUEUE_CORE_PACK_PROCESSING_TIMEOUT_MS:60000}"
stats:
enabled: "${TB_QUEUE_CORE_STATS_ENABLED:true}"
- print-interval-ms: "${TB_QUEUE_CORE_STATS_PRINT_INTERVAL_MS:10000}"
+ print-interval-ms: "${TB_QUEUE_CORE_STATS_PRINT_INTERVAL_MS:60000}"
js:
# JS Eval request topic
request_topic: "${REMOTE_JS_EVAL_REQUEST_TOPIC:js_eval.requests}"
@@ -658,7 +657,7 @@ queue:
pack-processing-timeout: "${TB_QUEUE_RULE_ENGINE_PACK_PROCESSING_TIMEOUT_MS:60000}"
stats:
enabled: "${TB_QUEUE_RULE_ENGINE_STATS_ENABLED:true}"
- print-interval-ms: "${TB_QUEUE_RULE_ENGINE_STATS_PRINT_INTERVAL_MS:10000}"
+ print-interval-ms: "${TB_QUEUE_RULE_ENGINE_STATS_PRINT_INTERVAL_MS:60000}"
queues:
- name: "${TB_QUEUE_RE_MAIN_QUEUE_NAME:Main}"
topic: "${TB_QUEUE_RE_MAIN_TOPIC:tb_rule_engine.main}"
diff --git a/application/src/test/java/org/thingsboard/server/service/cluster/routing/ConsistentHashParitionServiceTest.java b/application/src/test/java/org/thingsboard/server/service/cluster/routing/HashPartitionServiceTest.java
similarity index 83%
rename from application/src/test/java/org/thingsboard/server/service/cluster/routing/ConsistentHashParitionServiceTest.java
rename to application/src/test/java/org/thingsboard/server/service/cluster/routing/HashPartitionServiceTest.java
index 6b03482c0e..dcf98c1443 100644
--- a/application/src/test/java/org/thingsboard/server/service/cluster/routing/ConsistentHashParitionServiceTest.java
+++ b/application/src/test/java/org/thingsboard/server/service/cluster/routing/HashPartitionServiceTest.java
@@ -26,14 +26,14 @@ import org.springframework.context.ApplicationEventPublisher;
import org.springframework.test.util.ReflectionTestUtils;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
-import org.thingsboard.server.queue.discovery.ConsistentHashPartitionService;
+import org.thingsboard.server.common.msg.queue.ServiceQueue;
+import org.thingsboard.server.queue.discovery.HashPartitionService;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.discovery.TenantRoutingInfoService;
import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
-import org.thingsboard.server.queue.settings.TbRuleEngineQueueConfiguration;
import java.util.ArrayList;
import java.util.Collections;
@@ -48,18 +48,18 @@ import static org.mockito.Mockito.when;
@Slf4j
@RunWith(MockitoJUnitRunner.class)
-public class ConsistentHashParitionServiceTest {
+public class HashPartitionServiceTest {
public static final int ITERATIONS = 1000000;
- private ConsistentHashPartitionService clusterRoutingService;
+ public static final int SERVER_COUNT = 3;
+ private HashPartitionService clusterRoutingService;
private TbServiceInfoProvider discoveryService;
private TenantRoutingInfoService routingInfoService;
private ApplicationEventPublisher applicationEventPublisher;
private TbQueueRuleEngineSettings ruleEngineSettings;
- private String hashFunctionName = "murmur3_128";
- private Integer virtualNodesSize = 16;
+ private String hashFunctionName = "sha256";
@Before
@@ -68,25 +68,28 @@ public class ConsistentHashParitionServiceTest {
applicationEventPublisher = mock(ApplicationEventPublisher.class);
routingInfoService = mock(TenantRoutingInfoService.class);
ruleEngineSettings = mock(TbQueueRuleEngineSettings.class);
- clusterRoutingService = new ConsistentHashPartitionService(discoveryService,
+ clusterRoutingService = new HashPartitionService(discoveryService,
routingInfoService,
applicationEventPublisher,
ruleEngineSettings
);
when(ruleEngineSettings.getQueues()).thenReturn(Collections.emptyList());
ReflectionTestUtils.setField(clusterRoutingService, "coreTopic", "tb.core");
- ReflectionTestUtils.setField(clusterRoutingService, "corePartitions", 3);
+ ReflectionTestUtils.setField(clusterRoutingService, "corePartitions", 10);
ReflectionTestUtils.setField(clusterRoutingService, "hashFunctionName", hashFunctionName);
- ReflectionTestUtils.setField(clusterRoutingService, "virtualNodesSize", virtualNodesSize);
TransportProtos.ServiceInfo currentServer = TransportProtos.ServiceInfo.newBuilder()
- .setServiceId("100.96.1.1")
+ .setServiceId("tb-core-0")
+ .setTenantIdMSB(TenantId.NULL_UUID.getMostSignificantBits())
+ .setTenantIdLSB(TenantId.NULL_UUID.getLeastSignificantBits())
.addAllServiceTypes(Collections.singletonList(ServiceType.TB_CORE.name()))
.build();
// when(discoveryService.getServiceInfo()).thenReturn(currentServer);
List otherServers = new ArrayList<>();
- for (int i = 1; i < 30; i++) {
+ for (int i = 1; i < SERVER_COUNT; i++) {
otherServers.add(TransportProtos.ServiceInfo.newBuilder()
- .setServiceId("100.96." + i * 2 + "." + i)
+ .setServiceId("tb-rule-" + i)
+ .setTenantIdMSB(TenantId.NULL_UUID.getMostSignificantBits())
+ .setTenantIdLSB(TenantId.NULL_UUID.getLeastSignificantBits())
.addAllServiceTypes(Collections.singletonList(ServiceType.TB_CORE.name()))
.build());
}
@@ -116,12 +119,11 @@ public class ConsistentHashParitionServiceTest {
long end = System.currentTimeMillis();
double diff = (data.get(data.size() - 1).getValue() - data.get(0).getValue());
double diffPercent = (diff / ITERATIONS) * 100.0;
- System.out.println("Size: " + virtualNodesSize + " Time: " + (end - start) + " Diff: " + diff + "(" + String.format("%f", diffPercent) + "%)");
+ System.out.println("Time: " + (end - start) + " Diff: " + diff + "(" + String.format("%f", diffPercent) + "%)");
Assert.assertTrue(diffPercent < 0.5);
for (Map.Entry entry : data) {
System.out.println(entry.getKey() + ": " + entry.getValue());
}
-
}
}
diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ConsistentHashPartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java
similarity index 86%
rename from common/queue/src/main/java/org/thingsboard/server/queue/discovery/ConsistentHashPartitionService.java
rename to common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java
index 918570bc56..f3f11fa7a1 100644
--- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/ConsistentHashPartitionService.java
+++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java
@@ -18,6 +18,7 @@ package org.thingsboard.server.queue.discovery;
import com.google.common.hash.HashCode;
import com.google.common.hash.HashFunction;
import com.google.common.hash.Hashing;
+import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.ApplicationEventPublisher;
@@ -35,6 +36,7 @@ import org.thingsboard.server.queue.settings.TbQueueRuleEngineSettings;
import javax.annotation.PostConstruct;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
+import java.util.Comparator;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
@@ -48,7 +50,7 @@ import java.util.stream.Collectors;
@Service
@Slf4j
-public class ConsistentHashPartitionService implements PartitionService {
+public class HashPartitionService implements PartitionService {
@Value("${queue.core.topic}")
private String coreTopic;
@@ -56,8 +58,6 @@ public class ConsistentHashPartitionService implements PartitionService {
private Integer corePartitions;
@Value("${queue.partitions.hash_function_name:murmur3_128}")
private String hashFunctionName;
- @Value("${queue.partitions.virtual_nodes_size:16}")
- private Integer virtualNodesSize;
private final ApplicationEventPublisher applicationEventPublisher;
private final TbServiceInfoProvider serviceInfoProvider;
@@ -76,10 +76,10 @@ public class ConsistentHashPartitionService implements PartitionService {
private HashFunction hashFunction;
- public ConsistentHashPartitionService(TbServiceInfoProvider serviceInfoProvider,
- TenantRoutingInfoService tenantRoutingInfoService,
- ApplicationEventPublisher applicationEventPublisher,
- TbQueueRuleEngineSettings tbQueueRuleEngineSettings) {
+ public HashPartitionService(TbServiceInfoProvider serviceInfoProvider,
+ TenantRoutingInfoService tenantRoutingInfoService,
+ ApplicationEventPublisher applicationEventPublisher,
+ TbQueueRuleEngineSettings tbQueueRuleEngineSettings) {
this.serviceInfoProvider = serviceInfoProvider;
this.tenantRoutingInfoService = tenantRoutingInfoService;
this.applicationEventPublisher = applicationEventPublisher;
@@ -128,20 +128,22 @@ public class ConsistentHashPartitionService implements PartitionService {
public void recalculatePartitions(ServiceInfo currentService, List otherServices) {
logServiceInfo(currentService);
otherServices.forEach(this::logServiceInfo);
- Map> circles = new HashMap<>();
- addNode(circles, currentService);
+ Map> queueServicesMap = new HashMap<>();
+ addNode(queueServicesMap, currentService);
for (ServiceInfo other : otherServices) {
- addNode(circles, other);
+ addNode(queueServicesMap, other);
}
+ queueServicesMap.values().forEach(list -> list.sort((a, b) -> a.getServiceId().compareTo(b.getServiceId())));
+
ConcurrentMap> oldPartitions = myPartitions;
TenantId myIsolatedOrSystemTenantId = getSystemOrIsolatedTenantId(currentService);
myPartitions = new ConcurrentHashMap<>();
- partitionSizes.forEach((type, size) -> {
- ServiceQueueKey myServiceQueueKey = new ServiceQueueKey(type, myIsolatedOrSystemTenantId);
+ partitionSizes.forEach((serviceQueue, size) -> {
+ ServiceQueueKey myServiceQueueKey = new ServiceQueueKey(serviceQueue, myIsolatedOrSystemTenantId);
for (int i = 0; i < size; i++) {
- ServiceInfo serviceInfo = resolveByPartitionIdx(circles.get(myServiceQueueKey), i);
+ ServiceInfo serviceInfo = resolveByPartitionIdx(queueServicesMap.get(myServiceQueueKey), i);
if (currentService.equals(serviceInfo)) {
- ServiceQueueKey serviceQueueKey = new ServiceQueueKey(type, getSystemOrIsolatedTenantId(serviceInfo));
+ ServiceQueueKey serviceQueueKey = new ServiceQueueKey(serviceQueue, getSystemOrIsolatedTenantId(serviceInfo));
myPartitions.computeIfAbsent(serviceQueueKey, key -> new ArrayList<>()).add(i);
}
}
@@ -293,7 +295,7 @@ public class ConsistentHashPartitionService implements PartitionService {
return new TenantId(new UUID(serviceInfo.getTenantIdMSB(), serviceInfo.getTenantIdLSB()));
}
- private void addNode(Map> circles, ServiceInfo instance) {
+ private void addNode(Map> queueServiceList, ServiceInfo instance) {
TenantId tenantId = getSystemOrIsolatedTenantId(instance);
for (String serviceTypeStr : instance.getServiceTypesList()) {
ServiceType serviceType = ServiceType.valueOf(serviceTypeStr.toUpperCase());
@@ -302,34 +304,20 @@ public class ConsistentHashPartitionService implements PartitionService {
ServiceQueueKey serviceQueueKey = new ServiceQueueKey(new ServiceQueue(serviceType, queue.getName()), tenantId);
partitionSizes.put(new ServiceQueue(ServiceType.TB_RULE_ENGINE, queue.getName()), queue.getPartitions());
partitionTopics.put(new ServiceQueue(ServiceType.TB_RULE_ENGINE, queue.getName()), queue.getTopic());
- for (int i = 0; i < virtualNodesSize; i++) {
- circles.computeIfAbsent(serviceQueueKey, key -> new ConsistentHashCircle<>()).put(hash(instance, i).asLong(), instance);
- }
+ queueServiceList.computeIfAbsent(serviceQueueKey, key -> new ArrayList<>()).add(instance);
}
} else {
ServiceQueueKey serviceQueueKey = new ServiceQueueKey(new ServiceQueue(serviceType), tenantId);
- for (int i = 0; i < virtualNodesSize; i++) {
- circles.computeIfAbsent(serviceQueueKey, key -> new ConsistentHashCircle<>()).put(hash(instance, i).asLong(), instance);
- }
+ queueServiceList.computeIfAbsent(serviceQueueKey, key -> new ArrayList<>()).add(instance);
}
}
}
- private ServiceInfo resolveByPartitionIdx(ConsistentHashCircle circle, Integer partitionIdx) {
- if (circle == null || circle.isEmpty()) {
+ private ServiceInfo resolveByPartitionIdx(List servers, Integer partitionIdx) {
+ if (servers == null || servers.isEmpty()) {
return null;
}
- Long hash = hashFunction.newHasher().putInt(partitionIdx).hash().asLong();
- if (!circle.containsKey(hash)) {
- ConcurrentNavigableMap tailMap = circle.tailMap(hash);
- hash = tailMap.isEmpty() ?
- circle.firstKey() : tailMap.firstKey();
- }
- return circle.get(hash);
- }
-
- private HashCode hash(ServiceInfo instance, int i) {
- return hashFunction.newHasher().putString(instance.getServiceId(), StandardCharsets.UTF_8).putInt(i).hash();
+ return servers.get(partitionIdx % servers.size());
}
public static HashFunction forName(String name) {
@@ -338,12 +326,11 @@ public class ConsistentHashPartitionService implements PartitionService {
return Hashing.murmur3_32();
case "murmur3_128":
return Hashing.murmur3_128();
- case "crc32":
- return Hashing.crc32();
- case "md5":
- return Hashing.md5();
+ case "sha256":
+ return Hashing.sha256();
default:
throw new IllegalArgumentException("Can't find hash function with name " + name);
}
}
+
}
diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java
index 82678fde5d..3b4b26aa47 100644
--- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java
+++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java
@@ -37,6 +37,7 @@ import org.thingsboard.server.gen.transport.TransportProtos;
import java.lang.reflect.Field;
import java.util.List;
import java.util.Optional;
+import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
@@ -55,6 +56,8 @@ public class CoapTransportResource extends CoapResource {
private final Field observerField;
private final long timeout;
private final ConcurrentMap tokenToSessionIdMap = new ConcurrentHashMap<>();
+ private final Set rpcSubscriptions = ConcurrentHashMap.newKeySet();
+ private final Set attributeSubscriptions = ConcurrentHashMap.newKeySet();
public CoapTransportResource(CoapTransportContext context, String name) {
super(name);
@@ -149,11 +152,13 @@ public class CoapTransportResource extends CoapResource {
transportService.process(sessionInfo,
transportContext.getAdaptor().convertToPostAttributes(sessionId, request),
new CoapOkCallback(exchange));
+ reportActivity(sessionId, sessionInfo);
break;
case POST_TELEMETRY_REQUEST:
transportService.process(sessionInfo,
transportContext.getAdaptor().convertToPostTelemetry(sessionId, request),
new CoapOkCallback(exchange));
+ reportActivity(sessionId, sessionInfo);
break;
case CLAIM_REQUEST:
transportService.process(sessionInfo,
@@ -161,6 +166,7 @@ public class CoapTransportResource extends CoapResource {
new CoapOkCallback(exchange));
break;
case SUBSCRIBE_ATTRIBUTES_REQUEST:
+ attributeSubscriptions.add(sessionId);
advanced.setObserver(new CoapExchangeObserverProxy((ExchangeObserver) observerField.get(advanced),
registerAsyncCoapSession(exchange, request, sessionInfo, sessionId)));
transportService.process(sessionInfo,
@@ -168,6 +174,7 @@ public class CoapTransportResource extends CoapResource {
new CoapNoOpCallback(exchange));
break;
case UNSUBSCRIBE_ATTRIBUTES_REQUEST:
+ attributeSubscriptions.remove(sessionId);
TransportProtos.SessionInfoProto attrSession = lookupAsyncSessionInfo(request);
if (attrSession != null) {
transportService.process(attrSession,
@@ -177,6 +184,7 @@ public class CoapTransportResource extends CoapResource {
}
break;
case SUBSCRIBE_RPC_COMMANDS_REQUEST:
+ rpcSubscriptions.add(sessionId);
advanced.setObserver(new CoapExchangeObserverProxy((ExchangeObserver) observerField.get(advanced),
registerAsyncCoapSession(exchange, request, sessionInfo, sessionId)));
transportService.process(sessionInfo,
@@ -184,13 +192,13 @@ public class CoapTransportResource extends CoapResource {
new CoapNoOpCallback(exchange));
break;
case UNSUBSCRIBE_RPC_COMMANDS_REQUEST:
+ rpcSubscriptions.remove(sessionId);
TransportProtos.SessionInfoProto rpcSession = lookupAsyncSessionInfo(request);
if (rpcSession != null) {
transportService.process(rpcSession,
TransportProtos.SubscribeToRPCMsg.newBuilder().setUnsubscribe(true).build(),
new CoapOkCallback(exchange));
- transportService.process(sessionInfo, getSessionEventMsg(TransportProtos.SessionEvent.CLOSED), null);
- transportService.deregisterSession(rpcSession);
+ closeAndDeregister(sessionInfo);
}
break;
case TO_DEVICE_RPC_RESPONSE:
@@ -221,6 +229,14 @@ public class CoapTransportResource extends CoapResource {
}));
}
+ private void reportActivity(UUID sessionId, TransportProtos.SessionInfoProto sessionInfo) {
+ transportContext.getTransportService().process(sessionInfo, TransportProtos.SubscriptionInfoProto.newBuilder()
+ .setAttributeSubscription(attributeSubscriptions.contains(sessionId))
+ .setRpcSubscription(rpcSubscriptions.contains(sessionId))
+ .setLastActivityTime(System.currentTimeMillis())
+ .build(), TransportServiceCallback.EMPTY);
+ }
+
private TransportProtos.SessionInfoProto lookupAsyncSessionInfo(Request request) {
String token = request.getSource().getHostAddress() + ":" + request.getSourcePort() + ":" + request.getTokenString();
return tokenToSessionIdMap.remove(token);
@@ -438,6 +454,9 @@ public class CoapTransportResource extends CoapResource {
private void closeAndDeregister(TransportProtos.SessionInfoProto session) {
transportService.process(session, getSessionEventMsg(TransportProtos.SessionEvent.CLOSED), null);
transportService.deregisterSession(session);
+ UUID sessionId = new UUID(session.getSessionIdMSB(), session.getSessionIdLSB());
+ rpcSubscriptions.remove(sessionId);
+ attributeSubscriptions.remove(sessionId);
}
}
diff --git a/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java b/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java
index c2f171e39b..727600a626 100644
--- a/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java
+++ b/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java
@@ -36,6 +36,7 @@ import org.thingsboard.server.common.transport.TransportContext;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.adaptor.JsonConverter;
+import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.DeviceInfoProto;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg;
@@ -102,6 +103,7 @@ public class DeviceApiController {
TransportService transportService = transportContext.getTransportService();
transportService.process(sessionInfo, JsonConverter.convertToAttributesProto(new JsonParser().parse(json)),
new HttpOkCallback(responseWriter));
+ reportActivity(sessionInfo);
}));
return responseWriter;
}
@@ -115,6 +117,7 @@ public class DeviceApiController {
TransportService transportService = transportContext.getTransportService();
transportService.process(sessionInfo, JsonConverter.convertToTelemetryProto(new JsonParser().parse(json)),
new HttpOkCallback(responseWriter));
+ reportActivity(sessionInfo);
}));
return responseWriter;
}
@@ -274,7 +277,6 @@ public class DeviceApiController {
}
}
-
private static class HttpSessionListener implements SessionMsgListener {
private final DeferredResult responseWriter;
@@ -308,4 +310,13 @@ public class DeviceApiController {
responseWriter.setResult(new ResponseEntity<>(JsonConverter.toJson(msg).toString(), HttpStatus.OK));
}
}
+
+ private void reportActivity(SessionInfoProto sessionInfo) {
+ transportContext.getTransportService().process(sessionInfo, TransportProtos.SubscriptionInfoProto.newBuilder()
+ .setAttributeSubscription(false)
+ .setRpcSubscription(false)
+ .setLastActivityTime(System.currentTimeMillis())
+ .build(), TransportServiceCallback.EMPTY);
+ }
+
}
diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportServiceCallback.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportServiceCallback.java
index 1f0d6c4c2c..f5b9ff22dd 100644
--- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportServiceCallback.java
+++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportServiceCallback.java
@@ -20,7 +20,20 @@ package org.thingsboard.server.common.transport;
*/
public interface TransportServiceCallback {
+ TransportServiceCallback EMPTY = new TransportServiceCallback() {
+ @Override
+ public void onSuccess(Void msg) {
+
+ }
+
+ @Override
+ public void onError(Throwable e) {
+
+ }
+ };
+
void onSuccess(T msg);
+
void onError(Throwable e);
}
diff --git a/pom.xml b/pom.xml
index 4a8fcad0ee..524fb8ec3e 100755
--- a/pom.xml
+++ b/pom.xml
@@ -30,10 +30,10 @@
${basedir}
thingsboard
- 2.2.4.RELEASE
+ 2.2.6.RELEASE
2.1.2.RELEASE
- 5.2.2.RELEASE
- 5.2.2.RELEASE
+ 5.2.6.RELEASE
+ 5.2.3.RELEASE
2.2.4.RELEASE
3.1.0
0.7.0
@@ -62,14 +62,14 @@
2.6.2
1.7
2.0
- 1.4.3
+ 1.6.2
4.2.0
3.5.5
3.11.4
1.22.1
1.16.18
1.1.0
- 4.1.45.Final
+ 4.1.49.Final
1.5.0
4.8.0
2.19.1
@@ -467,12 +467,12 @@
org.springframework.security
spring-security-oauth2-client
- ${spring.version}
+ ${spring-security.version}
org.springframework.security
spring-security-oauth2-jose
- ${spring.version}
+ ${spring-security.version}
org.springframework.boot
@@ -587,8 +587,8 @@
${rabbitmq.version}
- javax.mail
- mail
+ com.sun.mail
+ javax.mail
${mail.version}
@@ -606,6 +606,12 @@
org.apache.zookeeper
zookeeper
${zookeeper.version}
+
+
+ log4j
+ log4j
+
+
com.jayway.jsonpath
@@ -688,6 +694,12 @@
com.github.fge
json-schema-validator
${json-schema-validator.version}
+
+
+ javax.mail
+ mailapi
+
+
com.typesafe.akka
@@ -924,12 +936,6 @@
com.microsoft.azure
azure-servicebus
${azure-servicebus.version}
-
-
- com.microsoft.azure
- adal4j
-
-
org.passay
diff --git a/rule-engine/rule-engine-api/pom.xml b/rule-engine/rule-engine-api/pom.xml
index 2cc6bf8d64..3978251987 100644
--- a/rule-engine/rule-engine-api/pom.xml
+++ b/rule-engine/rule-engine-api/pom.xml
@@ -93,5 +93,10 @@
spring-data-redis
provided
+
+ com.sun.mail
+ javax.mail
+ provided
+