Browse Source

Fix conflicts

pull/10184/head
Igor Kulikov 3 years ago
parent
commit
2fb9b74bae
  1. 115
      application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java
  2. 64
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  3. 10
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  4. 71
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
  5. 4
      application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerStats.java
  6. 12
      application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractTbRuleEngineSubmitStrategy.java
  7. 6
      application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java
  8. 9
      application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java
  9. 8
      common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueClusterService.java
  10. 12
      common/proto/src/main/proto/queue.proto
  11. 57
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java
  12. 4
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java
  13. 8
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  14. 3
      dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueService.java
  15. 1
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java
  16. 37
      ui-ngx/src/app/core/http/image.service.ts
  17. 23
      ui-ngx/src/app/modules/home/components/widget/lib/maps/leaflet-map.ts
  18. 3
      ui-ngx/src/app/modules/home/components/widget/lib/maps/markers.ts

115
application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java

@ -50,22 +50,15 @@ public class DefaultTbQueueService extends AbstractTbEntityService implements Tb
public Queue saveQueue(Queue queue) {
boolean create = queue.getId() == null;
Queue oldQueue;
if (create) {
oldQueue = null;
} else {
oldQueue = queueService.findQueueById(queue.getTenantId(), queue.getId());
}
//TODO: add checkNotNull
Queue savedQueue = queueService.saveQueue(queue);
if (create) {
onQueueCreated(savedQueue);
} else {
onQueueUpdated(savedQueue, oldQueue);
}
createTopicsIfNeeded(savedQueue, oldQueue);
tbClusterService.onQueuesUpdate(List.of(savedQueue));
return savedQueue;
}
@ -73,54 +66,14 @@ public class DefaultTbQueueService extends AbstractTbEntityService implements Tb
public void deleteQueue(TenantId tenantId, QueueId queueId) {
Queue queue = queueService.findQueueById(tenantId, queueId);
queueService.deleteQueue(tenantId, queueId);
onQueueDeleted(queue);
tbClusterService.onQueuesDelete(List.of(queue));
}
@Override
public void deleteQueueByQueueName(TenantId tenantId, String queueName) {
Queue queue = queueService.findQueueByTenantIdAndNameInternal(tenantId, queueName);
queueService.deleteQueue(tenantId, queue.getId());
onQueueDeleted(queue);
}
private void onQueueCreated(Queue queue) {
for (int i = 0; i < queue.getPartitions(); i++) {
tbQueueAdmin.createTopicIfNotExists(
new TopicPartitionInfo(queue.getTopic(), queue.getTenantId(), i, false).getFullTopicName(),
queue.getCustomProperties()
);
}
tbClusterService.onQueueChange(queue);
}
private void onQueueUpdated(Queue queue, Queue oldQueue) {
int oldPartitions = oldQueue.getPartitions();
int currentPartitions = queue.getPartitions();
if (currentPartitions != oldPartitions) {
if (currentPartitions > oldPartitions) {
log.info("Added [{}] new partitions to [{}] queue", currentPartitions - oldPartitions, queue.getName());
for (int i = oldPartitions; i < currentPartitions; i++) {
tbQueueAdmin.createTopicIfNotExists(
new TopicPartitionInfo(queue.getTopic(), queue.getTenantId(), i, false).getFullTopicName(),
queue.getCustomProperties()
);
}
tbClusterService.onQueueChange(queue);
} else {
log.info("Removed [{}] partitions from [{}] queue", oldPartitions - currentPartitions, queue.getName());
tbClusterService.onQueueChange(queue);
// TODO: move all the messages left in old partitions and delete topics
}
} else if (!oldQueue.equals(queue)) {
tbClusterService.onQueueChange(queue);
}
}
private void onQueueDeleted(Queue queue) {
tbClusterService.onQueueDelete(queue);
// queueStatsService.deleteQueueStatsByQueueId(tenantId, queueId);
tbClusterService.onQueuesDelete(List.of(queue));
}
@Override
@ -176,26 +129,56 @@ public class DefaultTbQueueService extends AbstractTbEntityService implements Tb
log.debug("[{}] Handling profile queue config update: creating queues {}, updating {}, deleting {}. Affected tenants: {}",
newTenantProfile.getUuidId(), toCreate, toUpdate, toRemove, tenantIds);
}
tenantIds.forEach(tenantId -> {
toCreate.forEach(key -> saveQueue(new Queue(tenantId, newQueues.get(key))));
toUpdate.forEach(key -> {
Queue queueToUpdate = new Queue(tenantId, newQueues.get(key));
Queue foundQueue = queueService.findQueueByTenantIdAndName(tenantId, key);
queueToUpdate.setId(foundQueue.getId());
queueToUpdate.setCreatedTime(foundQueue.getCreatedTime());
List<Queue> updated = new ArrayList<>();
List<Queue> deleted = new ArrayList<>();
for (TenantId tenantId : tenantIds) {
for (String name : toCreate) {
updated.add(new Queue(tenantId, newQueues.get(name)));
}
if (!queueToUpdate.equals(foundQueue)) {
saveQueue(queueToUpdate);
for (String name : toUpdate) {
Queue queue = new Queue(tenantId, newQueues.get(name));
Queue foundQueue = queueService.findQueueByTenantIdAndName(tenantId, name);
if (foundQueue != null) {
queue.setId(foundQueue.getId());
queue.setCreatedTime(foundQueue.getCreatedTime());
}
});
if (!queue.equals(foundQueue)) {
updated.add(queue);
createTopicsIfNeeded(queue, foundQueue);
}
}
for (String name : toRemove) {
Queue queue = queueService.findQueueByTenantIdAndNameInternal(tenantId, name);
deleted.add(queue);
}
}
toRemove.forEach(q -> {
Queue queue = queueService.findQueueByTenantIdAndNameInternal(tenantId, q);
QueueId queueIdForRemove = queue.getId();
deleteQueue(tenantId, queueIdForRemove);
if (!updated.isEmpty()) {
updated = updated.stream()
.map(queueService::saveQueue)
.collect(Collectors.toList());
tbClusterService.onQueuesUpdate(updated);
}
if (!deleted.isEmpty()) {
deleted.forEach(queue -> {
queueService.deleteQueue(queue.getTenantId(), queue.getId());
});
});
tbClusterService.onQueuesDelete(deleted);
}
}
private void createTopicsIfNeeded(Queue queue, Queue oldQueue) {
int newPartitions = queue.getPartitions();
int oldPartitions = oldQueue != null ? oldQueue.getPartitions() : 0;
for (int i = oldPartitions; i < newPartitions; i++) {
tbQueueAdmin.createTopicIfNotExists(
new TopicPartitionInfo(queue.getTopic(), queue.getTenantId(), i, false).getFullTopicName(),
queue.getCustomProperties()
);
}
}
}

64
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

@ -62,6 +62,8 @@ import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg;
import org.thingsboard.server.dao.edge.EdgeService;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.FromDeviceRPCResponseProto;
import org.thingsboard.server.gen.transport.TransportProtos.QueueDeleteMsg;
import org.thingsboard.server.gen.transport.TransportProtos.QueueUpdateMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToCoreNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineMsg;
@ -81,9 +83,11 @@ import org.thingsboard.server.service.profile.TbAssetProfileCache;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import java.util.Optional;
import java.util.List;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.util.ProtoUtils.toProto;
@ -563,40 +567,40 @@ public class DefaultTbClusterService implements TbClusterService {
}
@Override
public void onQueueChange(Queue queue) {
log.trace("[{}][{}] Processing queue change [{}] event", queue.getTenantId(), queue.getId(), queue.getName());
TransportProtos.QueueUpdateMsg queueUpdateMsg = TransportProtos.QueueUpdateMsg.newBuilder()
.setTenantIdMSB(queue.getTenantId().getId().getMostSignificantBits())
.setTenantIdLSB(queue.getTenantId().getId().getLeastSignificantBits())
.setQueueIdMSB(queue.getId().getId().getMostSignificantBits())
.setQueueIdLSB(queue.getId().getId().getLeastSignificantBits())
.setQueueName(queue.getName())
.setQueueTopic(queue.getTopic())
.setPartitions(queue.getPartitions())
.build();
ToRuleEngineNotificationMsg ruleEngineMsg = ToRuleEngineNotificationMsg.newBuilder().setQueueUpdateMsg(queueUpdateMsg).build();
ToCoreNotificationMsg coreMsg = ToCoreNotificationMsg.newBuilder().setQueueUpdateMsg(queueUpdateMsg).build();
ToTransportMsg transportMsg = ToTransportMsg.newBuilder().setQueueUpdateMsg(queueUpdateMsg).build();
public void onQueuesUpdate(List<Queue> queues) {
List<QueueUpdateMsg> queueUpdateMsgs = queues.stream()
.map(queue -> QueueUpdateMsg.newBuilder()
.setTenantIdMSB(queue.getTenantId().getId().getMostSignificantBits())
.setTenantIdLSB(queue.getTenantId().getId().getLeastSignificantBits())
.setQueueIdMSB(queue.getId().getId().getMostSignificantBits())
.setQueueIdLSB(queue.getId().getId().getLeastSignificantBits())
.setQueueName(queue.getName())
.setQueueTopic(queue.getTopic())
.setPartitions(queue.getPartitions())
.build())
.collect(Collectors.toList());
ToRuleEngineNotificationMsg ruleEngineMsg = ToRuleEngineNotificationMsg.newBuilder().addAllQueueUpdateMsgs(queueUpdateMsgs).build();
ToCoreNotificationMsg coreMsg = ToCoreNotificationMsg.newBuilder().addAllQueueUpdateMsgs(queueUpdateMsgs).build();
ToTransportMsg transportMsg = ToTransportMsg.newBuilder().addAllQueueUpdateMsgs(queueUpdateMsgs).build();
doSendQueueNotifications(ruleEngineMsg, coreMsg, transportMsg);
}
@Override
public void onQueueDelete(Queue queue) {
log.trace("[{}][{}] Processing queue delete [{}] event", queue.getTenantId(), queue.getId(), queue.getName());
TransportProtos.QueueDeleteMsg queueDeleteMsg = TransportProtos.QueueDeleteMsg.newBuilder()
.setTenantIdMSB(queue.getTenantId().getId().getMostSignificantBits())
.setTenantIdLSB(queue.getTenantId().getId().getLeastSignificantBits())
.setQueueIdMSB(queue.getId().getId().getMostSignificantBits())
.setQueueIdLSB(queue.getId().getId().getLeastSignificantBits())
.setQueueName(queue.getName())
.build();
ToRuleEngineNotificationMsg ruleEngineMsg = ToRuleEngineNotificationMsg.newBuilder().setQueueDeleteMsg(queueDeleteMsg).build();
ToCoreNotificationMsg coreMsg = ToCoreNotificationMsg.newBuilder().setQueueDeleteMsg(queueDeleteMsg).build();
ToTransportMsg transportMsg = ToTransportMsg.newBuilder().setQueueDeleteMsg(queueDeleteMsg).build();
public void onQueuesDelete(List<Queue> queues) {
List<QueueDeleteMsg> queueDeleteMsgs = queues.stream()
.map(queue -> QueueDeleteMsg.newBuilder()
.setTenantIdMSB(queue.getTenantId().getId().getMostSignificantBits())
.setTenantIdLSB(queue.getTenantId().getId().getLeastSignificantBits())
.setQueueIdMSB(queue.getId().getId().getMostSignificantBits())
.setQueueIdLSB(queue.getId().getId().getLeastSignificantBits())
.setQueueName(queue.getName())
.build())
.collect(Collectors.toList());
ToRuleEngineNotificationMsg ruleEngineMsg = ToRuleEngineNotificationMsg.newBuilder().addAllQueueDeleteMsgs(queueDeleteMsgs).build();
ToCoreNotificationMsg coreMsg = ToCoreNotificationMsg.newBuilder().addAllQueueDeleteMsgs(queueDeleteMsgs).build();
ToTransportMsg transportMsg = ToTransportMsg.newBuilder().addAllQueueDeleteMsgs(queueDeleteMsgs).build();
doSendQueueNotifications(ruleEngineMsg, coreMsg, transportMsg);
}

10
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -391,13 +391,11 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
} else if (!toCoreNotification.getFromEdgeSyncResponseMsg().isEmpty()) {
//will be removed in 3.6.1 in favour of hasFromEdgeSyncResponse()
forwardToAppActor(id, encodingService.decode(toCoreNotification.getFromEdgeSyncResponseMsg().toByteArray()), callback);
} else if (toCoreNotification.hasQueueUpdateMsg()) {
TransportProtos.QueueUpdateMsg queue = toCoreNotification.getQueueUpdateMsg();
partitionService.updateQueue(queue);
} else if (toCoreNotification.getQueueUpdateMsgsCount() > 0) {
partitionService.updateQueues(toCoreNotification.getQueueUpdateMsgsList());
callback.onSuccess();
} else if (toCoreNotification.hasQueueDeleteMsg()) {
TransportProtos.QueueDeleteMsg queue = toCoreNotification.getQueueDeleteMsg();
partitionService.removeQueue(queue);
} else if (toCoreNotification.getQueueDeleteMsgsCount() > 0) {
partitionService.removeQueues(toCoreNotification.getQueueDeleteMsgsList());
callback.onSuccess();
} else if (toCoreNotification.hasVcResponseMsg()) {
vcQueueService.processResponse(toCoreNotification.getVcResponseMsg());

71
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java

@ -36,6 +36,8 @@ import org.thingsboard.server.common.util.ProtoUtils;
import org.thingsboard.server.dao.queue.QueueService;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.QueueDeleteMsg;
import org.thingsboard.server.gen.transport.TransportProtos.QueueUpdateMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToRuleEngineNotificationMsg;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
import org.thingsboard.server.queue.discovery.PartitionService;
@ -164,11 +166,11 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
, proto.getResponse(), error);
tbDeviceRpcService.processRpcResponseFromDevice(response);
callback.onSuccess();
} else if (nfMsg.hasQueueUpdateMsg()) {
updateQueue(nfMsg.getQueueUpdateMsg());
} else if (nfMsg.getQueueUpdateMsgsCount() > 0) {
updateQueues(nfMsg.getQueueUpdateMsgsList());
callback.onSuccess();
} else if (nfMsg.hasQueueDeleteMsg()) {
deleteQueue(nfMsg.getQueueDeleteMsg());
} else if (nfMsg.getQueueDeleteMsgsCount() > 0) {
deleteQueues(nfMsg.getQueueDeleteMsgsList());
callback.onSuccess();
} else {
log.trace("Received notification with missing handler");
@ -176,39 +178,48 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService<
}
}
private void updateQueue(TransportProtos.QueueUpdateMsg queueUpdateMsg) {
log.info("Received queue update msg: [{}]", queueUpdateMsg);
TenantId tenantId = new TenantId(new UUID(queueUpdateMsg.getTenantIdMSB(), queueUpdateMsg.getTenantIdLSB()));
if (partitionService.isManagedByCurrentService(tenantId)) {
QueueId queueId = new QueueId(new UUID(queueUpdateMsg.getQueueIdMSB(), queueUpdateMsg.getQueueIdLSB()));
String queueName = queueUpdateMsg.getQueueName();
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queueName, tenantId);
Queue queue = queueService.findQueueById(tenantId, queueId);
TbRuleEngineQueueConsumerManager consumerManager = getOrCreateConsumer(queueKey);
Queue oldQueue = consumerManager.getQueue();
consumerManager.update(queue);
if (oldQueue != null && queue.getPartitions() == oldQueue.getPartitions()) {
return;
private void updateQueues(List<QueueUpdateMsg> queueUpdateMsgs) {
boolean partitionsChanged = false;
for (QueueUpdateMsg queueUpdateMsg : queueUpdateMsgs) {
log.info("Received queue update msg: [{}]", queueUpdateMsg);
TenantId tenantId = new TenantId(new UUID(queueUpdateMsg.getTenantIdMSB(), queueUpdateMsg.getTenantIdLSB()));
if (partitionService.isManagedByCurrentService(tenantId)) {
QueueId queueId = new QueueId(new UUID(queueUpdateMsg.getQueueIdMSB(), queueUpdateMsg.getQueueIdLSB()));
String queueName = queueUpdateMsg.getQueueName();
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queueName, tenantId);
Queue queue = queueService.findQueueById(tenantId, queueId);
TbRuleEngineQueueConsumerManager consumerManager = getOrCreateConsumer(queueKey);
Queue oldQueue = consumerManager.getQueue();
consumerManager.update(queue);
if (oldQueue == null || queue.getPartitions() != oldQueue.getPartitions()) {
partitionsChanged = true;
}
} else {
partitionsChanged = true;
}
}
partitionService.updateQueue(queueUpdateMsg);
partitionService.recalculatePartitions(ctx.getServiceInfoProvider().getServiceInfo(),
new ArrayList<>(partitionService.getOtherServices(ServiceType.TB_RULE_ENGINE)));
if (partitionsChanged) {
partitionService.updateQueues(queueUpdateMsgs);
partitionService.recalculatePartitions(ctx.getServiceInfoProvider().getServiceInfo(),
new ArrayList<>(partitionService.getOtherServices(ServiceType.TB_RULE_ENGINE)));
}
}
private void deleteQueue(TransportProtos.QueueDeleteMsg queueDeleteMsg) {
log.info("Received queue delete msg: [{}]", queueDeleteMsg);
TenantId tenantId = new TenantId(new UUID(queueDeleteMsg.getTenantIdMSB(), queueDeleteMsg.getTenantIdLSB()));
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queueDeleteMsg.getQueueName(), tenantId);
var consumerManager = consumers.remove(queueKey);
if (consumerManager != null) {
consumerManager.delete(true);
private void deleteQueues(List<QueueDeleteMsg> queueDeleteMsgs) {
for (QueueDeleteMsg queueDeleteMsg : queueDeleteMsgs) {
log.info("Received queue delete msg: [{}]", queueDeleteMsg);
TenantId tenantId = new TenantId(new UUID(queueDeleteMsg.getTenantIdMSB(), queueDeleteMsg.getTenantIdLSB()));
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queueDeleteMsg.getQueueName(), tenantId);
var consumerManager = consumers.remove(queueKey);
if (consumerManager != null) {
consumerManager.delete(true);
}
}
partitionService.removeQueue(queueDeleteMsg);
partitionService.removeQueues(queueDeleteMsgs);
partitionService.recalculatePartitions(ctx.getServiceInfoProvider().getServiceInfo(), new ArrayList<>(partitionService.getOtherServices(ServiceType.TB_RULE_ENGINE)));
}

4
application/src/main/java/org/thingsboard/server/service/queue/TbCoreConsumerStats.java

@ -184,9 +184,9 @@ public class TbCoreConsumerStats {
toCoreNfEdgeSyncResponseCounter.increment();
} else if (!msg.getFromEdgeSyncResponseMsg().isEmpty()) {
toCoreNfEdgeSyncResponseCounter.increment();
} else if (msg.hasQueueUpdateMsg()) {
} else if (msg.getQueueUpdateMsgsCount() > 0) {
toCoreNfQueueUpdateCounter.increment();
} else if (msg.hasQueueDeleteMsg()) {
} else if (msg.getQueueDeleteMsgsCount() > 0) {
toCoreNfQueueDeleteCounter.increment();
} else if (msg.hasVcResponseMsg()) {
toCoreNfVersionControlResponseCounter.increment();

12
application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractTbRuleEngineSubmitStrategy.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.service.queue.processing;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.common.TbProtoQueueMsg;
@ -51,7 +52,16 @@ public abstract class AbstractTbRuleEngineSubmitStrategy implements TbRuleEngine
List<IdMsgPair<TransportProtos.ToRuleEngineMsg>> newOrderedMsgList = new ArrayList<>(reprocessMap.size());
for (IdMsgPair<TransportProtos.ToRuleEngineMsg> pair : orderedMsgList) {
if (reprocessMap.containsKey(pair.uuid)) {
newOrderedMsgList.add(pair);
if (StringUtils.isNotEmpty(pair.getMsg().getValue().getFailureMessage())) {
var toRuleEngineMsg = TransportProtos.ToRuleEngineMsg.newBuilder(pair.getMsg().getValue())
.clearFailureMessage()
.clearRelationTypes()
.build();
var newMsg = new TbProtoQueueMsg<>(pair.getMsg().getKey(), toRuleEngineMsg, pair.getMsg().getHeaders());
newOrderedMsgList.add(new IdMsgPair<>(pair.getUuid(), newMsg));
} else {
newOrderedMsgList.add(pair);
}
}
}
orderedMsgList = newOrderedMsgList;

6
application/src/test/java/org/thingsboard/server/queue/discovery/HashPartitionServiceTest.java

@ -315,7 +315,7 @@ public class HashPartitionServiceTest {
.setPartitions(isolatedQueue.getPartitions())
.build();
partitionService_common.updateQueue(queueUpdateMsg);
partitionService_common.updateQueues(List.of(queueUpdateMsg));
partitionService_common.recalculatePartitions(commonRuleEngine, List.of(dedicatedRuleEngine));
// expecting event about no partitions for isolated queue key
verifyPartitionChangeEvent(event -> {
@ -323,7 +323,7 @@ public class HashPartitionServiceTest {
return event.getPartitionsMap().get(queueKey).isEmpty();
});
partitionService_dedicated.updateQueue(queueUpdateMsg);
partitionService_dedicated.updateQueues(List.of(queueUpdateMsg));
partitionService_dedicated.recalculatePartitions(dedicatedRuleEngine, List.of(commonRuleEngine));
verifyPartitionChangeEvent(event -> {
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.MAIN_QUEUE_NAME, tenantId);
@ -342,7 +342,7 @@ public class HashPartitionServiceTest {
.setQueueIdLSB(isolatedQueue.getUuidId().getLeastSignificantBits())
.setQueueName(isolatedQueue.getName())
.build();
partitionService_dedicated.removeQueue(queueDeleteMsg);
partitionService_dedicated.removeQueues(List.of(queueDeleteMsg));
verifyPartitionChangeEvent(event -> {
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, DataConstants.MAIN_QUEUE_NAME, tenantId);
return event.getPartitionsMap().get(queueKey).isEmpty();

9
application/src/test/java/org/thingsboard/server/service/queue/DefaultTbClusterServiceTest.java

@ -41,6 +41,7 @@ import org.thingsboard.server.service.gateway_device.GatewayNotificationsService
import org.thingsboard.server.service.profile.TbAssetProfileCache;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import java.util.List;
import java.util.UUID;
import static org.mockito.ArgumentMatchers.any;
@ -95,7 +96,7 @@ public class DefaultTbClusterServiceTest {
when(producerProvider.getRuleEngineNotificationsMsgProducer()).thenReturn(tbQueueProducer);
clusterService.onQueueChange(createTestQueue());
clusterService.onQueuesUpdate(List.of(createTestQueue()));
verify(topicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, MONOLITH);
verify(topicService, never()).getNotificationsTopic(eq(ServiceType.TB_CORE), any());
@ -120,7 +121,7 @@ public class DefaultTbClusterServiceTest {
when(producerProvider.getRuleEngineNotificationsMsgProducer()).thenReturn(tbQueueProducer);
clusterService.onQueueChange(createTestQueue());
clusterService.onQueuesUpdate(List.of(createTestQueue()));
verify(topicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith1);
verify(topicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith2);
@ -148,7 +149,7 @@ public class DefaultTbClusterServiceTest {
when(producerProvider.getRuleEngineNotificationsMsgProducer()).thenReturn(tbREQueueProducer);
when(producerProvider.getTransportNotificationsMsgProducer()).thenReturn(tbTransportQueueProducer);
clusterService.onQueueChange(createTestQueue());
clusterService.onQueuesUpdate(List.of(createTestQueue()));
verify(topicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, MONOLITH);
verify(topicService, times(1)).getNotificationsTopic(ServiceType.TB_TRANSPORT, TRANSPORT);
@ -194,7 +195,7 @@ public class DefaultTbClusterServiceTest {
when(producerProvider.getTbCoreNotificationsMsgProducer()).thenReturn(tbCoreQueueProducer);
when(producerProvider.getTransportNotificationsMsgProducer()).thenReturn(tbTransportQueueProducer);
clusterService.onQueueChange(createTestQueue());
clusterService.onQueuesUpdate(List.of(createTestQueue()));
verify(topicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith1);
verify(topicService, times(1)).getNotificationsTopic(ServiceType.TB_RULE_ENGINE, monolith2);

8
common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueClusterService.java

@ -17,8 +17,12 @@ package org.thingsboard.server.queue;
import org.thingsboard.server.common.data.queue.Queue;
import java.util.List;
public interface TbQueueClusterService {
void onQueueChange(Queue queue);
void onQueueDelete(Queue queue);
void onQueuesUpdate(List<Queue> queues);
void onQueuesDelete(List<Queue> queues);
}

12
common/proto/src/main/proto/queue.proto

@ -1280,8 +1280,8 @@ message ToCoreNotificationMsg {
FromDeviceRPCResponseProto fromDeviceRpcResponse = 2;
bytes componentLifecycleMsg = 3 [deprecated = true];
bytes edgeEventUpdateMsg = 4 [deprecated = true];
QueueUpdateMsg queueUpdateMsg = 5;
QueueDeleteMsg queueDeleteMsg = 6;
repeated QueueUpdateMsg queueUpdateMsgs = 5;
repeated QueueDeleteMsg queueDeleteMsgs = 6;
VersionControlResponseMsg vcResponseMsg = 7;
bytes toEdgeSyncRequestMsg = 8 [deprecated = true];
bytes fromEdgeSyncResponseMsg = 9 [deprecated = true];
@ -1307,8 +1307,8 @@ message ToRuleEngineMsg {
message ToRuleEngineNotificationMsg {
bytes componentLifecycleMsg = 1 [deprecated = true];
FromDeviceRPCResponseProto fromDeviceRpcResponse = 2;
QueueUpdateMsg queueUpdateMsg = 3;
QueueDeleteMsg queueDeleteMsg = 4;
repeated QueueUpdateMsg queueUpdateMsgs = 3;
repeated QueueDeleteMsg queueDeleteMsgs = 4;
ComponentLifecycleMsgProto componentLifecycle = 5;
}
@ -1328,8 +1328,8 @@ message ToTransportMsg {
ResourceUpdateMsg resourceUpdateMsg = 12;
ResourceDeleteMsg resourceDeleteMsg = 13;
UplinkNotificationMsg uplinkNotificationMsg = 14;
QueueUpdateMsg queueUpdateMsg = 15;
QueueDeleteMsg queueDeleteMsg = 16;
repeated QueueUpdateMsg queueUpdateMsgs = 15;
repeated QueueDeleteMsg queueDeleteMsgs = 16;
}
message UsageStatsKVProto{

57
common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java

@ -171,27 +171,35 @@ public class HashPartitionService implements PartitionService {
}
@Override
public void updateQueue(TransportProtos.QueueUpdateMsg queueUpdateMsg) {
TenantId tenantId = new TenantId(new UUID(queueUpdateMsg.getTenantIdMSB(), queueUpdateMsg.getTenantIdLSB()));
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queueUpdateMsg.getQueueName(), tenantId);
partitionTopicsMap.put(queueKey, queueUpdateMsg.getQueueTopic());
partitionSizesMap.put(queueKey, queueUpdateMsg.getPartitions());
myPartitions.remove(queueKey);
if (!tenantId.isSysTenantId()) {
tenantRoutingInfoMap.remove(tenantId);
public void updateQueues(List<TransportProtos.QueueUpdateMsg> queueUpdateMsgs) {
for (TransportProtos.QueueUpdateMsg queueUpdateMsg : queueUpdateMsgs) {
TenantId tenantId = TenantId.fromUUID(new UUID(queueUpdateMsg.getTenantIdMSB(), queueUpdateMsg.getTenantIdLSB()));
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queueUpdateMsg.getQueueName(), tenantId);
partitionTopicsMap.put(queueKey, queueUpdateMsg.getQueueTopic());
partitionSizesMap.put(queueKey, queueUpdateMsg.getPartitions());
if (!tenantId.isSysTenantId()) {
tenantRoutingInfoMap.remove(tenantId);
}
}
}
@Override
public void removeQueue(TransportProtos.QueueDeleteMsg queueDeleteMsg) {
TenantId tenantId = new TenantId(new UUID(queueDeleteMsg.getTenantIdMSB(), queueDeleteMsg.getTenantIdLSB()));
QueueKey queueKey = new QueueKey(ServiceType.TB_RULE_ENGINE, queueDeleteMsg.getQueueName(), tenantId);
myPartitions.remove(queueKey);
partitionTopicsMap.remove(queueKey);
partitionSizesMap.remove(queueKey);
evictTenantInfo(tenantId);
public void removeQueues(List<TransportProtos.QueueDeleteMsg> queueDeleteMsgs) {
List<QueueKey> queueKeys = queueDeleteMsgs.stream()
.map(queueDeleteMsg -> {
TenantId tenantId = TenantId.fromUUID(new UUID(queueDeleteMsg.getTenantIdMSB(), queueDeleteMsg.getTenantIdLSB()));
return new QueueKey(ServiceType.TB_RULE_ENGINE, queueDeleteMsg.getQueueName(), tenantId);
})
.collect(Collectors.toList());
queueKeys.forEach(queueKey -> {
myPartitions.remove(queueKey);
partitionTopicsMap.remove(queueKey);
partitionSizesMap.remove(queueKey);
evictTenantInfo(queueKey.getTenantId());
});
if (serviceInfoProvider.isService(ServiceType.TB_RULE_ENGINE)) {
publishPartitionChangeEvent(ServiceType.TB_RULE_ENGINE, Map.of(queueKey, Collections.emptySet()));
publishPartitionChangeEvent(ServiceType.TB_RULE_ENGINE, queueKeys.stream()
.collect(Collectors.toMap(k -> k, k -> Collections.emptySet())));
}
}
@ -321,13 +329,11 @@ public class HashPartitionService implements PartitionService {
.forEach(removed::add);
}
removed.forEach(queueKey -> {
log.info("[{}] NO MORE PARTITIONS FOR CURRENT KEY", queueKey);
changedPartitionsMap.put(queueKey, Collections.emptySet());
});
myPartitions.forEach((queueKey, partitions) -> {
if (!partitions.equals(oldPartitions.get(queueKey))) {
log.info("[{}] NEW PARTITIONS: {}", queueKey, partitions);
Set<TopicPartitionInfo> tpiList = partitions.stream()
.map(partition -> buildTopicPartitionInfo(queueKey, partition))
.collect(Collectors.toSet());
@ -373,14 +379,11 @@ public class HashPartitionService implements PartitionService {
}
private void publishPartitionChangeEvent(ServiceType serviceType, Map<QueueKey, Set<TopicPartitionInfo>> partitionsMap) {
if (log.isDebugEnabled()) {
log.debug("Publishing partition change event for service type " + serviceType + ":" + System.lineSeparator() +
partitionsMap.entrySet().stream()
.map(entry -> entry.getKey() + " - " + entry.getValue().stream()
.map(TopicPartitionInfo::getFullTopicName).sorted()
.collect(Collectors.toList()))
.collect(Collectors.joining(System.lineSeparator())));
}
log.info("Partitions changed: {}", System.lineSeparator() + partitionsMap.entrySet().stream()
.map(entry -> "[" + entry.getKey() + "] - [" + entry.getValue().stream()
.map(tpi -> tpi.getPartition().orElse(-1).toString()).sorted()
.collect(Collectors.joining(", ")) + "]")
.collect(Collectors.joining(System.lineSeparator())));
PartitionChangeEvent event = new PartitionChangeEvent(this, serviceType, partitionsMap);
try {
applicationEventPublisher.publishEvent(event);
@ -490,7 +493,7 @@ public class HashPartitionService implements PartitionService {
}
private void logServiceInfo(TransportProtos.ServiceInfo server) {
log.info("[{}] Found common server: [{}]", server.getServiceId(), server.getServiceTypesList());
log.info("[{}] Found common server: {}", server.getServiceId(), server.getServiceTypesList());
}
private void addNode(Map<QueueKey, List<ServiceInfo>> queueServiceList, ServiceInfo instance) {

4
common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java

@ -63,9 +63,9 @@ public interface PartitionService {
int countTransportsByType(String type);
void updateQueue(TransportProtos.QueueUpdateMsg queueUpdateMsg);
void updateQueues(List<TransportProtos.QueueUpdateMsg> queueUpdateMsgs);
void removeQueue(TransportProtos.QueueDeleteMsg queueDeleteMsg);
void removeQueues(List<TransportProtos.QueueDeleteMsg> queueDeleteMsgs);
void removeTenant(TenantId tenantId);

8
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

@ -993,10 +993,10 @@ public class DefaultTransportService extends TransportActivityManager implements
log.warn("ResourceDelete - [{}] [{}]", id, mdRez);
transportCallbackExecutor.submit(() -> mdRez.getListener().onResourceDelete(msg));
});
} else if (toSessionMsg.hasQueueUpdateMsg()) {
partitionService.updateQueue(toSessionMsg.getQueueUpdateMsg());
} else if (toSessionMsg.hasQueueDeleteMsg()) {
partitionService.removeQueue(toSessionMsg.getQueueDeleteMsg());
} else if (toSessionMsg.getQueueUpdateMsgsCount() > 0) {
partitionService.updateQueues(toSessionMsg.getQueueUpdateMsgsList());
} else if (toSessionMsg.getQueueDeleteMsgsCount() > 0) {
partitionService.removeQueues(toSessionMsg.getQueueDeleteMsgsList());
} else {
//TODO: should we notify the device actor about missed session?
log.debug("[{}] Missing session.", sessionId);

3
dao/src/main/java/org/thingsboard/server/dao/queue/BaseQueueService.java

@ -57,9 +57,6 @@ public class BaseQueueService extends AbstractEntityService implements QueueServ
@Autowired
private DataValidator<Queue> queueValidator;
// @Autowired
// private QueueStatsService queueStatsService;
@Override
public Queue saveQueue(Queue queue) {
log.trace("Executing createOrUpdateQueue [{}]", queue);

1
dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java

@ -149,6 +149,7 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
}
@Override
@Transactional
public RuleChainUpdateResult saveRuleChainMetaData(TenantId tenantId, RuleChainMetaData ruleChainMetaData, Function<RuleNode, RuleNode> ruleNodeUpdater) {
Validator.validateId(ruleChainMetaData.getRuleChainId(), "Incorrect rule chain id.");
RuleChain ruleChain = findRuleChainById(tenantId, ruleChainMetaData.getRuleChainId());

37
ui-ngx/src/app/core/http/image.service.ts

@ -18,7 +18,7 @@ import { Injectable } from '@angular/core';
import { HttpClient } from '@angular/common/http';
import { PageLink } from '@shared/models/page/page-link';
import { defaultHttpOptionsFromConfig, defaultHttpUploadOptions, RequestConfig } from '@core/http/http-utils';
import { Observable, of } from 'rxjs';
import { Observable, of, ReplaySubject } from 'rxjs';
import { PageData } from '@shared/models/page/page-data';
import {
NO_IMAGE_DATA_URI,
@ -36,6 +36,9 @@ import { ResourcesService } from '@core/services/resources.service';
providedIn: 'root'
})
export class ImageService {
private imagesLoading: { [url: string]: ReplaySubject<Blob> } = {};
constructor(
private http: HttpClient,
private sanitizer: DomSanitizer,
@ -95,12 +98,34 @@ export class ImageService {
parts[parts.length - 1] = encodeURIComponent(key);
const encodedUrl = parts.join('/');
const imageLink = preview ? (encodedUrl + '/preview') : encodedUrl;
const options = defaultHttpOptionsFromConfig({ignoreLoading: true, ignoreErrors: true});
return this.http
.get(imageLink, {...options, ...{ responseType: 'blob' } }).pipe(
return this.loadImageDataUrl(imageLink, asString, emptyUrl);
}
private loadImageDataUrl(imageLink: string, asString = false, emptyUrl = NO_IMAGE_DATA_URI): Observable<SafeUrl | string> {
let request: ReplaySubject<Blob>;
if (this.imagesLoading[imageLink]) {
request = this.imagesLoading[imageLink];
} else {
request = new ReplaySubject<Blob>(1);
this.imagesLoading[imageLink] = request;
const options = defaultHttpOptionsFromConfig({ignoreLoading: true, ignoreErrors: true});
this.http.get(imageLink, {...options, ...{ responseType: 'blob' } }).subscribe({
next: (value) => {
request.next(value);
request.complete();
},
error: err => {
request.error(err);
},
complete: () => {
delete this.imagesLoading[imageLink];
}
});
}
return request.pipe(
switchMap(val => blobToBase64(val).pipe(
map((dataUrl) => asString ? dataUrl : this.sanitizer.bypassSecurityTrustUrl(dataUrl))
)),
map((dataUrl) => asString ? dataUrl : this.sanitizer.bypassSecurityTrustUrl(dataUrl))
)),
catchError(() => of(asString ? emptyUrl : this.sanitizer.bypassSecurityTrustUrl(emptyUrl)))
);
}

23
ui-ngx/src/app/modules/home/components/widget/lib/maps/leaflet-map.ts

@ -32,7 +32,7 @@ import {
WidgetUnitedMapSettings
} from './map-models';
import { Marker } from './markers';
import { map, Observable, of, switchMap } from 'rxjs';
import { map, Observable, of } from 'rxjs';
import { Polyline } from './polyline';
import { Polygon } from './polygon';
import { Circle } from './circle';
@ -64,6 +64,7 @@ import { MatDialog } from '@angular/material/dialog';
import { FormattedData, ReplaceInfo } from '@shared/models/widget.models';
import ITooltipsterInstance = JQueryTooltipster.ITooltipsterInstance;
import { ImagePipe } from '@shared/pipe/image.pipe';
import { take, tap } from 'rxjs/operators';
export default abstract class LeafletMap {
@ -940,7 +941,12 @@ export default abstract class LeafletMap {
this.markersData = markersData;
if (this.options.useClusterMarkers) {
if (createdMarkers.length) {
this.markersCluster.addLayers(createdMarkers.map(marker => marker.leafletMarker));
createdMarkers.forEach((marker) => {
marker.createMarkerIconSubject.pipe(
tap(() => this.markersCluster.addLayer(marker.leafletMarker)),
take(1)
).subscribe();
});
}
if (updatedMarkers.length) {
this.markersCluster.refreshClusters(updatedMarkers.map(marker => marker.leafletMarker));
@ -971,10 +977,15 @@ export default abstract class LeafletMap {
}
this.markers.set(key, newMarker);
if (!this.options.useClusterMarkers) {
this.map.addLayer(newMarker.leafletMarker);
if (this.map.pm.globalDragModeEnabled() && newMarker.leafletMarker.pm) {
newMarker.leafletMarker.pm.enableLayerDrag();
}
newMarker.createMarkerIconSubject.pipe(
tap(() => {
this.map.addLayer(newMarker.leafletMarker);
if (this.map.pm.globalDragModeEnabled() && newMarker.leafletMarker.pm) {
newMarker.leafletMarker.pm.enableLayerDrag();
}
}),
take(1)
).subscribe();
}
return newMarker;
}

3
ui-ngx/src/app/modules/home/components/widget/lib/maps/markers.ts

@ -23,6 +23,7 @@ import { fillDataPattern, isDefined, isDefinedAndNotNull, processDataPattern, sa
import LeafletMap from './leaflet-map';
import { FormattedData } from '@shared/models/widget.models';
import { ImagePipe } from '@shared/pipe/image.pipe';
import { ReplaySubject } from 'rxjs';
export class Marker {
@ -33,6 +34,7 @@ export class Marker {
tooltipOffset: L.LatLngTuple;
markerOffset: L.LatLngTuple;
tooltip: L.Popup;
createMarkerIconSubject = new ReplaySubject<MarkerIconInfo>();
constructor(private map: LeafletMap,
private location: L.LatLng,
@ -148,6 +150,7 @@ export class Marker {
this.labelOffset = [0, -iconInfo.size[1] * this.markerOffset[1] + 10];
}
this.updateMarkerLabel(settings);
this.createMarkerIconSubject.next(iconInfo);
});
}

Loading…
Cancel
Save