From 5862b417aa48859c9f8038c927f27ed36f86aae8 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 28 Jul 2023 12:06:04 +0300 Subject: [PATCH] Add custom topic properties configuration --- .../entitiy/queue/DefaultTbQueueService.java | 8 ++++++-- .../org/thingsboard/server/queue/TbQueueAdmin.java | 6 +++++- .../server/common/data/queue/Queue.java | 14 +++++++++++++- .../queue/RuleEngineTbQueueAdminFactory.java | 2 +- .../queue/azure/servicebus/TbServiceBusAdmin.java | 7 ++++--- .../server/queue/kafka/TbKafkaAdmin.java | 6 +++--- .../provider/InMemoryTbTransportQueueFactory.java | 2 +- .../server/queue/pubsub/TbPubSubAdmin.java | 2 +- .../server/queue/rabbitmq/TbRabbitMqAdmin.java | 9 ++++++++- .../queue/rabbitmq/TbRabbitMqQueueArguments.java | 8 ++++---- .../server/queue/sqs/TbAwsSqsAdmin.java | 4 +++- .../server/queue/sqs/TbAwsSqsQueueAttributes.java | 8 +++++++- .../server/queue/util/PropertyUtils.java | 14 ++++++++++++++ .../queue/tenant-profile-queues.component.ts | 3 ++- .../components/profile/tenant-profile.component.ts | 9 ++++++--- .../components/queue/queue-form.component.html | 5 +++++ .../home/components/queue/queue-form.component.ts | 3 ++- ui-ngx/src/app/shared/models/queue.models.ts | 1 + .../src/assets/locale/locale.constant-en_US.json | 2 ++ 19 files changed, 88 insertions(+), 25 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java index 0166a425d3..9e3d38cf32 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java @@ -96,7 +96,9 @@ public class DefaultTbQueueService extends AbstractTbEntityService implements Tb private void onQueueCreated(Queue queue) { for (int i = 0; i < queue.getPartitions(); i++) { tbQueueAdmin.createTopicIfNotExists( - new TopicPartitionInfo(queue.getTopic(), queue.getTenantId(), i, false).getFullTopicName()); + new TopicPartitionInfo(queue.getTopic(), queue.getTenantId(), i, false).getFullTopicName(), + queue.getCustomProperties() + ); } tbClusterService.onQueueChange(queue); @@ -111,7 +113,9 @@ public class DefaultTbQueueService extends AbstractTbEntityService implements Tb 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()); + new TopicPartitionInfo(queue.getTopic(), queue.getTenantId(), i, false).getFullTopicName(), + queue.getCustomProperties() + ); } tbClusterService.onQueueChange(queue); } else { diff --git a/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueAdmin.java b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueAdmin.java index 4b2bde733e..19aa0284ea 100644 --- a/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueAdmin.java +++ b/common/cluster-api/src/main/java/org/thingsboard/server/queue/TbQueueAdmin.java @@ -17,7 +17,11 @@ package org.thingsboard.server.queue; public interface TbQueueAdmin { - void createTopicIfNotExists(String topic); + default void createTopicIfNotExists(String topic) { + createTopicIfNotExists(topic, null); + } + + void createTopicIfNotExists(String topic, String properties); void destroy(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/queue/Queue.java b/common/data/src/main/java/org/thingsboard/server/common/data/queue/Queue.java index f757998d04..deb0e57f6b 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/queue/Queue.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/queue/Queue.java @@ -15,6 +15,8 @@ */ package org.thingsboard.server.common.data.queue; +import com.fasterxml.jackson.annotation.JsonIgnore; +import com.fasterxml.jackson.databind.JsonNode; import lombok.Data; import org.thingsboard.server.common.data.HasName; import org.thingsboard.server.common.data.HasTenantId; @@ -25,6 +27,8 @@ import org.thingsboard.server.common.data.tenant.profile.TenantProfileQueueConfi import org.thingsboard.server.common.data.validation.Length; import org.thingsboard.server.common.data.validation.NoXss; +import java.util.Optional; + @Data public class Queue extends SearchTextBasedWithAdditionalInfo implements HasName, HasTenantId { private TenantId tenantId; @@ -65,4 +69,12 @@ public class Queue extends SearchTextBasedWithAdditionalInfo implements public String getSearchText() { return getName(); } -} \ No newline at end of file + + @JsonIgnore + public String getCustomProperties() { + return Optional.ofNullable(getAdditionalInfo()) + .map(info -> info.get("customProperties")) + .filter(JsonNode::isTextual).map(JsonNode::asText).orElse(null); + } + +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/RuleEngineTbQueueAdminFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/RuleEngineTbQueueAdminFactory.java index fe29c3a04f..7a8764325f 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/RuleEngineTbQueueAdminFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/RuleEngineTbQueueAdminFactory.java @@ -99,7 +99,7 @@ public class RuleEngineTbQueueAdminFactory { return new TbQueueAdmin() { @Override - public void createTopicIfNotExists(String topic) { + public void createTopicIfNotExists(String topic, String properties) { } @Override diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java index e171bb7a31..d95d2064a4 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java @@ -22,6 +22,7 @@ import com.microsoft.azure.servicebus.primitives.MessagingEntityAlreadyExistsExc import com.microsoft.azure.servicebus.primitives.ServiceBusException; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.queue.TbQueueAdmin; +import org.thingsboard.server.queue.util.PropertyUtils; import java.io.IOException; import java.time.Duration; @@ -60,7 +61,7 @@ public class TbServiceBusAdmin implements TbQueueAdmin { } @Override - public void createTopicIfNotExists(String topic) { + public void createTopicIfNotExists(String topic, String properties) { if (queues.contains(topic)) { return; } @@ -68,7 +69,7 @@ public class TbServiceBusAdmin implements TbQueueAdmin { try { QueueDescription queueDescription = new QueueDescription(topic); queueDescription.setRequiresDuplicateDetection(false); - setQueueConfigs(queueDescription); + setQueueConfigs(queueDescription, PropertyUtils.getProps(queueConfigs, properties)); client.createQueue(queueDescription); queues.add(topic); @@ -107,7 +108,7 @@ public class TbServiceBusAdmin implements TbQueueAdmin { } } - private void setQueueConfigs(QueueDescription queueDescription) { + private void setQueueConfigs(QueueDescription queueDescription, Map queueConfigs) { queueConfigs.forEach((confKey, confValue) -> { switch (confKey) { case MAX_SIZE: diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java index f15b9258e8..d486d04783 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java @@ -21,6 +21,7 @@ import org.apache.kafka.clients.admin.CreateTopicsResult; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.common.errors.TopicExistsException; import org.thingsboard.server.queue.TbQueueAdmin; +import org.thingsboard.server.queue.util.PropertyUtils; import java.util.Collections; import java.util.Map; @@ -62,12 +63,12 @@ public class TbKafkaAdmin implements TbQueueAdmin { } @Override - public void createTopicIfNotExists(String topic) { + public void createTopicIfNotExists(String topic, String properties) { if (topics.contains(topic)) { return; } try { - NewTopic newTopic = new NewTopic(topic, numPartitions, replicationFactor).configs(topicConfigs); + NewTopic newTopic = new NewTopic(topic, numPartitions, replicationFactor).configs(PropertyUtils.getProps(topicConfigs, properties)); createTopic(newTopic).values().get(topic).get(); topics.add(topic); } catch (ExecutionException ee) { @@ -81,7 +82,6 @@ public class TbKafkaAdmin implements TbQueueAdmin { log.warn("[{}] Failed to create topic", topic, e); throw new RuntimeException(e); } - } @Override diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java index 60c464ecab..4d0089457d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/provider/InMemoryTbTransportQueueFactory.java @@ -73,7 +73,7 @@ public class InMemoryTbTransportQueueFactory implements TbTransportQueueFactory templateBuilder.queueAdmin(new TbQueueAdmin() { @Override - public void createTopicIfNotExists(String topic) {} + public void createTopicIfNotExists(String topic, String properties) {} @Override public void destroy() {} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java index d1a4942ad3..f9f20c2448 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java @@ -103,7 +103,7 @@ public class TbPubSubAdmin implements TbQueueAdmin { } @Override - public void createTopicIfNotExists(String partition) { + public void createTopicIfNotExists(String partition, String properties) { TopicName topicName = TopicName.newBuilder() .setTopic(partition) .setProject(pubSubSettings.getProjectId()) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java index 00a2ee4c6c..fb646f383a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java @@ -18,9 +18,11 @@ package org.thingsboard.server.queue.rabbitmq; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.StringUtils; import org.thingsboard.server.queue.TbQueueAdmin; import java.io.IOException; +import java.util.HashMap; import java.util.Map; import java.util.concurrent.TimeoutException; @@ -50,7 +52,12 @@ public class TbRabbitMqAdmin implements TbQueueAdmin { } @Override - public void createTopicIfNotExists(String topic) { + public void createTopicIfNotExists(String topic, String properties) { + Map arguments = this.arguments; + if (StringUtils.isNotBlank(properties)) { + arguments = new HashMap<>(arguments); + arguments.putAll(TbRabbitMqQueueArguments.getArgs(properties)); + } try { channel.queueDeclare(topic, false, false, false, arguments); } catch (IOException e) { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqQueueArguments.java b/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqQueueArguments.java index cb96abdf3c..8fa8c537e6 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqQueueArguments.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqQueueArguments.java @@ -65,7 +65,7 @@ public class TbRabbitMqQueueArguments { vcArgs = getArgs(vcProperties); } - private Map getArgs(String properties) { + public static Map getArgs(String properties) { Map configs = new HashMap<>(); if (StringUtils.isNotEmpty(properties)) { for (String property : properties.split(";")) { @@ -78,7 +78,7 @@ public class TbRabbitMqQueueArguments { return configs; } - private Object getObjectValue(String str) { + private static Object getObjectValue(String str) { if (str.equalsIgnoreCase("true") || str.equalsIgnoreCase("false")) { return Boolean.valueOf(str); } else if (isNumeric(str)) { @@ -87,7 +87,7 @@ public class TbRabbitMqQueueArguments { return str; } - private Object getNumericValue(String str) { + private static Object getNumericValue(String str) { if (str.contains(".")) { return Double.valueOf(str); } else { @@ -97,7 +97,7 @@ public class TbRabbitMqQueueArguments { private static final Pattern PATTERN = Pattern.compile("-?\\d+(\\.\\d+)?"); - public boolean isNumeric(String strNum) { + private static boolean isNumeric(String strNum) { if (strNum == null) { return false; } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java index f88a34941a..ba4eeb6ca4 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java @@ -26,6 +26,7 @@ import com.amazonaws.services.sqs.model.CreateQueueRequest; import com.amazonaws.services.sqs.model.GetQueueUrlResult; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.queue.TbQueueAdmin; +import org.thingsboard.server.queue.util.PropertyUtils; import java.util.Map; import java.util.function.Function; @@ -63,11 +64,12 @@ public class TbAwsSqsAdmin implements TbQueueAdmin { } @Override - public void createTopicIfNotExists(String topic) { + public void createTopicIfNotExists(String topic, String properties) { String queueName = convertTopicToQueueName(topic); if (queues.containsKey(queueName)) { return; } + Map attributes = PropertyUtils.getProps(this.attributes, properties, TbAwsSqsQueueAttributes::toConfigs); final CreateQueueRequest createQueueRequest = new CreateQueueRequest(queueName).withAttributes(attributes); String queueUrl = sqsClient.createQueue(createQueueRequest).getQueueUrl(); queues.put(getQueueNameFromUrl(queueUrl), queueUrl); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsQueueAttributes.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsQueueAttributes.java index 66110ade74..faa8eccc90 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsQueueAttributes.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsQueueAttributes.java @@ -76,6 +76,12 @@ public class TbAwsSqsQueueAttributes { private Map getConfigs(String properties) { Map configs = new HashMap<>(defaultAttributes); + configs.putAll(toConfigs(properties)); + return configs; + } + + public static Map toConfigs(String properties) { + Map configs = new HashMap<>(); if (StringUtils.isNotEmpty(properties)) { for (String property : properties.split(";")) { int delimiterPosition = property.indexOf(":"); @@ -88,7 +94,7 @@ public class TbAwsSqsQueueAttributes { return configs; } - private void validateAttributeName(String key) { + private static void validateAttributeName(String key) { QueueAttributeName.fromValue(key); } } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/util/PropertyUtils.java b/common/queue/src/main/java/org/thingsboard/server/queue/util/PropertyUtils.java index 089d7f2219..afee64f382 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/util/PropertyUtils.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/util/PropertyUtils.java @@ -19,6 +19,7 @@ import org.thingsboard.server.common.data.StringUtils; import java.util.HashMap; import java.util.Map; +import java.util.function.Function; public class PropertyUtils { @@ -37,4 +38,17 @@ public class PropertyUtils { return configs; } + public static Map getProps(Map defaultProperties, String propertiesStr) { + return getProps(defaultProperties, propertiesStr, PropertyUtils::getProps); + } + + public static Map getProps(Map defaultProperties, String propertiesStr, Function> parser) { + Map properties = defaultProperties; + if (StringUtils.isNotBlank(propertiesStr)) { + properties = new HashMap<>(properties); + properties.putAll(parser.apply(propertiesStr)); + } + return properties; + } + } diff --git a/ui-ngx/src/app/modules/home/components/profile/queue/tenant-profile-queues.component.ts b/ui-ngx/src/app/modules/home/components/profile/queue/tenant-profile-queues.component.ts index ac128b3ec5..298b4fe575 100644 --- a/ui-ngx/src/app/modules/home/components/profile/queue/tenant-profile-queues.component.ts +++ b/ui-ngx/src/app/modules/home/components/profile/queue/tenant-profile-queues.component.ts @@ -173,7 +173,8 @@ export class TenantProfileQueuesComponent implements ControlValueAccessor, Valid }, topic: '', additionalInfo: { - description: '' + description: '', + customProperties: '' } }; this.idMap.push(queue.id); diff --git a/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.ts b/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.ts index fa8b1636ff..bed7531495 100644 --- a/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.ts +++ b/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.ts @@ -74,7 +74,8 @@ export class TenantProfileComponent extends EntityComponent { }, topic: 'tb_rule_engine.main', additionalInfo: { - description: '' + description: '', + customProperties: '' } }, { @@ -97,7 +98,8 @@ export class TenantProfileComponent extends EntityComponent { maxPauseBetweenRetries: 5 }, additionalInfo: { - description: '' + description: '', + customProperties: '' } }, { @@ -120,7 +122,8 @@ export class TenantProfileComponent extends EntityComponent { maxPauseBetweenRetries: 5 }, additionalInfo: { - description: '' + description: '', + customProperties: '' } } ]; diff --git a/ui-ngx/src/app/modules/home/components/queue/queue-form.component.html b/ui-ngx/src/app/modules/home/components/queue/queue-form.component.html index 56b20d8aa8..4845f7f25a 100644 --- a/ui-ngx/src/app/modules/home/components/queue/queue-form.component.html +++ b/ui-ngx/src/app/modules/home/components/queue/queue-form.component.html @@ -203,6 +203,11 @@ + + queue.custom-properties + + queue.custom-properties-hint + queue.description diff --git a/ui-ngx/src/app/modules/home/components/queue/queue-form.component.ts b/ui-ngx/src/app/modules/home/components/queue/queue-form.component.ts index b123cf2ad5..e4fbcd031b 100644 --- a/ui-ngx/src/app/modules/home/components/queue/queue-form.component.ts +++ b/ui-ngx/src/app/modules/home/components/queue/queue-form.component.ts @@ -117,7 +117,8 @@ export class QueueFormComponent implements ControlValueAccessor, OnInit, OnDestr }), topic: [''], additionalInfo: this.fb.group({ - description: [''] + description: [''], + customProperties: [''] }) }); this.valueChange$ = this.queueFormGroup.valueChanges.subscribe(() => { diff --git a/ui-ngx/src/app/shared/models/queue.models.ts b/ui-ngx/src/app/shared/models/queue.models.ts index 76ee0bf022..07a57e68c2 100644 --- a/ui-ngx/src/app/shared/models/queue.models.ts +++ b/ui-ngx/src/app/shared/models/queue.models.ts @@ -121,5 +121,6 @@ export interface QueueInfo extends BaseData { topic: string; additionalInfo: { description?: string; + customProperties?: string; }; } diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index c9d236bf2a..c1092f9dad 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -3465,6 +3465,8 @@ "description": "Description", "description-hint": "This text will be displayed in the Queue description instead of the selected strategy", "alt-description": "Submit Strategy: {{submitStrategy}}, Processing Strategy: {{processingStrategy}}", + "custom-properties": "Custom properties", + "custom-properties-hint": "Custom queue (topic) creation properties, e.g. 'retention.ms:604800000;retention.bytes:1048576000'", "strategies": { "sequential-by-originator-label": "Sequential by originator", "sequential-by-originator-hint": "New message for e.g. device A is not submitted until previous message for device A is acknowledged",