Browse Source

Merge remote-tracking branch 'upstream/master'

pull/2733/head
Volodymyr Babak 6 years ago
parent
commit
4240ac10f2
  1. 4
      application/pom.xml
  2. 4
      application/src/main/data/json/demo/dashboards/gateways.json
  3. 4
      application/src/main/data/json/system/widget_bundles/entity_admin_widgets.json
  4. 2
      application/src/main/java/org/thingsboard/server/controller/AuthController.java
  5. 13
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  6. 18
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
  7. 7
      application/src/main/resources/thingsboard.yml
  8. 30
      application/src/test/java/org/thingsboard/server/service/cluster/routing/HashPartitionServiceTest.java
  9. 63
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java
  10. 23
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java
  11. 13
      common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java
  12. 13
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportServiceCallback.java
  13. 36
      pom.xml
  14. 5
      rule-engine/rule-engine-api/pom.xml

4
application/pom.xml

@ -193,8 +193,8 @@
<artifactId>logback-classic</artifactId>
</dependency>
<dependency>
<groupId>javax.mail</groupId>
<artifactId>mail</artifactId>
<groupId>com.sun.mail</groupId>
<artifactId>javax.mail</artifactId>
</dependency>
<dependency>
<groupId>org.apache.curator</groupId>

4
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"
}
}

4
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",

2
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<OAuth2ClientInfo> getOath2Clients() throws ThingsboardException {
public List<OAuth2ClientInfo> getOAuth2Clients() throws ThingsboardException {
try {
return oauth2Service.getOAuth2Clients();
} catch (Exception e) {

13
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<ToCore
}
@PreDestroy
public void destroy(){
public void destroy() {
super.destroy();
}
@ -143,8 +145,13 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
} else if (toCoreMsg.getToDeviceActorNotificationMsg() != null && !toCoreMsg.getToDeviceActorNotificationMsg().isEmpty()) {
Optional<TbActorMsg> 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();
}

18
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<String, TbQueueConsumer<TbProtoQueueMsg<ToRuleEngineMsg>>> consumers = new ConcurrentHashMap<>();
private final ConcurrentMap<String, TbRuleEngineQueueConfiguration> consumerConfigurations = new ConcurrentHashMap<>();
private final ConcurrentMap<String, TbRuleEngineConsumerStats> 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();
}
}

7
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}"

30
application/src/test/java/org/thingsboard/server/service/cluster/routing/ConsistentHashParitionServiceTest.java → 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<TransportProtos.ServiceInfo> 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<Integer, Integer> entry : data) {
System.out.println(entry.getKey() + ": " + entry.getValue());
}
}
}

63
common/queue/src/main/java/org/thingsboard/server/queue/discovery/ConsistentHashPartitionService.java → 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<ServiceInfo> otherServices) {
logServiceInfo(currentService);
otherServices.forEach(this::logServiceInfo);
Map<ServiceQueueKey, ConsistentHashCircle<ServiceInfo>> circles = new HashMap<>();
addNode(circles, currentService);
Map<ServiceQueueKey, List<ServiceInfo>> 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<ServiceQueueKey, List<Integer>> 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<ServiceQueueKey, ConsistentHashCircle<ServiceInfo>> circles, ServiceInfo instance) {
private void addNode(Map<ServiceQueueKey, List<ServiceInfo>> 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<ServiceInfo> circle, Integer partitionIdx) {
if (circle == null || circle.isEmpty()) {
private ServiceInfo resolveByPartitionIdx(List<ServiceInfo> servers, Integer partitionIdx) {
if (servers == null || servers.isEmpty()) {
return null;
}
Long hash = hashFunction.newHasher().putInt(partitionIdx).hash().asLong();
if (!circle.containsKey(hash)) {
ConcurrentNavigableMap<Long, ServiceInfo> 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);
}
}
}

23
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<String, TransportProtos.SessionInfoProto> tokenToSessionIdMap = new ConcurrentHashMap<>();
private final Set<UUID> rpcSubscriptions = ConcurrentHashMap.newKeySet();
private final Set<UUID> 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);
}
}

13
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<ResponseEntity> 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);
}
}

13
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<T> {
TransportServiceCallback<Void> EMPTY = new TransportServiceCallback<Void>() {
@Override
public void onSuccess(Void msg) {
}
@Override
public void onError(Throwable e) {
}
};
void onSuccess(T msg);
void onError(Throwable e);
}

36
pom.xml

@ -30,10 +30,10 @@
<properties>
<main.dir>${basedir}</main.dir>
<pkg.user>thingsboard</pkg.user>
<spring-boot.version>2.2.4.RELEASE</spring-boot.version>
<spring-boot.version>2.2.6.RELEASE</spring-boot.version>
<spring-oauth2.version>2.1.2.RELEASE</spring-oauth2.version>
<spring.version>5.2.2.RELEASE</spring.version>
<spring-security.version>5.2.2.RELEASE</spring-security.version>
<spring.version>5.2.6.RELEASE</spring.version>
<spring-security.version>5.2.3.RELEASE</spring-security.version>
<spring-data-redis.version>2.2.4.RELEASE</spring-data-redis.version>
<jedis.version>3.1.0</jedis.version>
<jjwt.version>0.7.0</jjwt.version>
@ -62,14 +62,14 @@
<gson.version>2.6.2</gson.version>
<velocity.version>1.7</velocity.version>
<velocity-tools.version>2.0</velocity-tools.version>
<mail.version>1.4.3</mail.version>
<mail.version>1.6.2</mail.version>
<curator.version>4.2.0</curator.version>
<zookeeper.version>3.5.5</zookeeper.version>
<protobuf.version>3.11.4</protobuf.version>
<grpc.version>1.22.1</grpc.version>
<lombok.version>1.16.18</lombok.version>
<paho.client.version>1.1.0</paho.client.version>
<netty.version>4.1.45.Final</netty.version>
<netty.version>4.1.49.Final</netty.version>
<os-maven-plugin.version>1.5.0</os-maven-plugin.version>
<rabbitmq.version>4.8.0</rabbitmq.version>
<surfire.version>2.19.1</surfire.version>
@ -467,12 +467,12 @@
<dependency>
<groupId>org.springframework.security</groupId>
<artifactId>spring-security-oauth2-client</artifactId>
<version>${spring.version}</version>
<version>${spring-security.version}</version>
</dependency>
<dependency>
<groupId>org.springframework.security</groupId>
<artifactId>spring-security-oauth2-jose</artifactId>
<version>${spring.version}</version>
<version>${spring-security.version}</version>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
@ -587,8 +587,8 @@
<version>${rabbitmq.version}</version>
</dependency>
<dependency>
<groupId>javax.mail</groupId>
<artifactId>mail</artifactId>
<groupId>com.sun.mail</groupId>
<artifactId>javax.mail</artifactId>
<version>${mail.version}</version>
</dependency>
<dependency>
@ -606,6 +606,12 @@
<groupId>org.apache.zookeeper</groupId>
<artifactId>zookeeper</artifactId>
<version>${zookeeper.version}</version>
<exclusions>
<exclusion>
<groupId>log4j</groupId>
<artifactId>log4j</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>com.jayway.jsonpath</groupId>
@ -688,6 +694,12 @@
<groupId>com.github.fge</groupId>
<artifactId>json-schema-validator</artifactId>
<version>${json-schema-validator.version}</version>
<exclusions>
<exclusion>
<groupId>javax.mail</groupId>
<artifactId>mailapi</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>com.typesafe.akka</groupId>
@ -924,12 +936,6 @@
<groupId>com.microsoft.azure</groupId>
<artifactId>azure-servicebus</artifactId>
<version>${azure-servicebus.version}</version>
<exclusions>
<exclusion>
<groupId>com.microsoft.azure</groupId>
<artifactId>adal4j</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>org.passay</groupId>

5
rule-engine/rule-engine-api/pom.xml

@ -93,5 +93,10 @@
<artifactId>spring-data-redis</artifactId>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>com.sun.mail</groupId>
<artifactId>javax.mail</artifactId>
<scope>provided</scope>
</dependency>
</dependencies>
</project>

Loading…
Cancel
Save