diff --git a/application/src/main/java/org/thingsboard/server/controller/QueueController.java b/application/src/main/java/org/thingsboard/server/controller/QueueController.java index 9b83aa43fb..39a971e6ed 100644 --- a/application/src/main/java/org/thingsboard/server/controller/QueueController.java +++ b/application/src/main/java/org/thingsboard/server/controller/QueueController.java @@ -15,10 +15,7 @@ */ package org.thingsboard.server.controller; -import io.swagger.annotations.ApiOperation; -import io.swagger.annotations.ApiParam; import lombok.RequiredArgsConstructor; -import org.springframework.http.MediaType; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.RequestBody; @@ -38,14 +35,7 @@ import org.thingsboard.server.service.entitiy.queue.TbQueueService; import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Resource; -import java.util.Collections; -import java.util.Set; import java.util.UUID; -import java.util.stream.Collectors; - -import static org.thingsboard.server.controller.ControllerConstants.QUEUE_SERVICE_TYPE_ALLOWABLE_VALUES; -import static org.thingsboard.server.controller.ControllerConstants.QUEUE_SERVICE_TYPE_DESCRIPTION; -import static org.thingsboard.server.controller.ControllerConstants.TENANT_AUTHORITY_PARAGRAPH; @RestController @TbCoreComponent @@ -55,27 +45,6 @@ public class QueueController extends BaseController { private final TbQueueService tbQueueService; - @ApiOperation(value = "Get queue names (getTenantQueuesByServiceType)", - notes = "Returns a set of unique queue names based on service type. " + TENANT_AUTHORITY_PARAGRAPH) - @PreAuthorize("hasAuthority('TENANT_ADMIN')") - @RequestMapping(value = "/queues", params = {"serviceType"}, produces = MediaType.APPLICATION_JSON_VALUE, method = RequestMethod.GET) - @ResponseBody() - public Set getTenantQueuesByServiceType(@ApiParam(value = QUEUE_SERVICE_TYPE_DESCRIPTION, allowableValues = QUEUE_SERVICE_TYPE_ALLOWABLE_VALUES) - @RequestParam String serviceType) throws ThingsboardException { - checkParameter("serviceType", serviceType); - try { - ServiceType type = ServiceType.valueOf(serviceType); - switch (type) { - case TB_RULE_ENGINE: - return queueService.findQueuesByTenantId(getTenantId()).stream().map(Queue::getName).collect(Collectors.toSet()); - default: - return Collections.emptySet(); - } - } catch (Exception e) { - throw handleException(e); - } - } - @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") @RequestMapping(value = "/queues", params = {"serviceType", "pageSize", "page"}, method = RequestMethod.GET) @ResponseBody diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseTenantControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseTenantControllerTest.java index 1ae3f497e9..2cc8bccc42 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseTenantControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseTenantControllerTest.java @@ -31,11 +31,26 @@ import org.springframework.test.web.servlet.ResultActions; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.TenantInfo; +import org.thingsboard.server.common.data.TenantProfile; +import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; +import org.thingsboard.server.common.data.queue.ProcessingStrategy; +import org.thingsboard.server.common.data.queue.ProcessingStrategyType; +import org.thingsboard.server.common.data.queue.Queue; +import org.thingsboard.server.common.data.queue.SubmitStrategy; +import org.thingsboard.server.common.data.queue.SubmitStrategyType; +import org.thingsboard.server.common.data.security.Authority; +import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; +import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; +import org.thingsboard.server.common.data.tenant.profile.TenantProfileQueueConfiguration; import java.util.ArrayList; +import java.util.Comparator; +import java.util.HashMap; import java.util.List; +import java.util.Map; +import java.util.Random; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; @@ -328,4 +343,180 @@ public abstract class BaseTenantControllerTest extends AbstractControllerTest { return Futures.allAsList(futures); } + @Test + public void testUpdateQueueConfigForIsolatedTenant() throws Exception { + Comparator queueComparator = Comparator.comparing(Queue::getName); + final String username = "isolatedtenant@thingsboard.org"; + final String password = "123456"; + loginSysAdmin(); + + List sysAdminQueues; + PageLink pageLink = new PageLink(10); + PageData pageData; + pageData = doGetTypedWithPageLink("/api/queues?serviceType=TB_RULE_ENGINE&", new TypeReference<>() { + }, pageLink); + sysAdminQueues = pageData.getData(); + + Tenant tenant = new Tenant(); + tenant.setTitle("Isolated tenant"); + tenant = doPost("/api/tenant", tenant, Tenant.class); + + User tenantUser = new User(); + tenantUser.setAuthority(Authority.TENANT_ADMIN); + tenantUser.setTenantId(tenant.getId()); + tenantUser.setEmail(username); + createUserAndLogin(tenantUser, password); + + List foundTenantQueues; + + pageLink = new PageLink(10); + pageData = doGetTypedWithPageLink("/api/queues?serviceType=TB_RULE_ENGINE&", new TypeReference<>() {}, pageLink); + foundTenantQueues = pageData.getData(); + + Assert.assertEquals(sysAdminQueues, foundTenantQueues); + + loginSysAdmin(); + + TenantProfile tenantProfile = new TenantProfile(); + tenantProfile.setName("isolated-tb-rule-engine"); + TenantProfileData tenantProfileData = new TenantProfileData(); + tenantProfileData.setConfiguration(new DefaultTenantProfileConfiguration()); + tenantProfile.setProfileData(tenantProfileData); + tenantProfile.setIsolatedTbRuleEngine(true); + addQueueConfig(tenantProfile, "Main"); + addQueueConfig(tenantProfile, "Test"); + tenantProfile = doPost("/api/tenantProfile", tenantProfile, TenantProfile.class); + + tenant.setTenantProfileId(tenantProfile.getId()); + doPost("/api/tenant", tenant, Tenant.class); + + login(username, password); + + pageLink = new PageLink(10); + pageData = doGetTypedWithPageLink("/api/queues?serviceType=TB_RULE_ENGINE&", new TypeReference<>() {}, pageLink); + foundTenantQueues = pageData.getData(); + + Assert.assertEquals(2, foundTenantQueues.size()); + + List queuesFromConfig = getQueuesFromConfig(tenantProfile.getProfileData().getQueueConfiguration(), foundTenantQueues); + queuesFromConfig.sort(queueComparator); + foundTenantQueues.sort(queueComparator); + + Assert.assertEquals(queuesFromConfig, foundTenantQueues); + + loginSysAdmin(); + + TenantProfile tenantProfile2 = new TenantProfile(); + tenantProfile2.setName("isolated-tb-rule-engine2"); + TenantProfileData tenantProfileData2 = new TenantProfileData(); + tenantProfileData2.setConfiguration(new DefaultTenantProfileConfiguration()); + tenantProfile2.setProfileData(tenantProfileData2); + tenantProfile2.setIsolatedTbRuleEngine(true); + addQueueConfig(tenantProfile2, "Main"); + addQueueConfig(tenantProfile2, "Test"); + addQueueConfig(tenantProfile2, "Test2"); + tenantProfile2 = doPost("/api/tenantProfile", tenantProfile2, TenantProfile.class); + + tenant.setTenantProfileId(tenantProfile2.getId()); + doPost("/api/tenant", tenant, Tenant.class); + + login(username, password); + + pageLink = new PageLink(10); + pageData = doGetTypedWithPageLink("/api/queues?serviceType=TB_RULE_ENGINE&", new TypeReference<>() {}, pageLink); + foundTenantQueues = pageData.getData(); + + Assert.assertEquals(3, foundTenantQueues.size()); + + queuesFromConfig = getQueuesFromConfig(tenantProfile2.getProfileData().getQueueConfiguration(), foundTenantQueues); + queuesFromConfig.sort(queueComparator); + foundTenantQueues.sort(queueComparator); + + Assert.assertEquals(queuesFromConfig, foundTenantQueues); + + loginSysAdmin(); + + tenantProfile2.getProfileData().getQueueConfiguration().removeIf(q -> q.getName().equals("Test")); + tenantProfile2.getProfileData().getQueueConfiguration().removeIf(q -> q.getName().equals("Test2")); + addQueueConfig(tenantProfile2, "Test2"); + addQueueConfig(tenantProfile2, "Test3"); + + tenantProfile2 = doPost("/api/tenantProfile", tenantProfile2, TenantProfile.class); + + login(username, password); + + pageLink = new PageLink(10); + pageData = doGetTypedWithPageLink("/api/queues?serviceType=TB_RULE_ENGINE&", new TypeReference<>() {}, pageLink); + foundTenantQueues = pageData.getData(); + + Assert.assertEquals(3, foundTenantQueues.size()); + + queuesFromConfig = getQueuesFromConfig(tenantProfile2.getProfileData().getQueueConfiguration(), foundTenantQueues); + queuesFromConfig.sort(queueComparator); + foundTenantQueues.sort(queueComparator); + + Assert.assertEquals(queuesFromConfig, foundTenantQueues); + + loginSysAdmin(); + + tenant.setTenantProfileId(null); + doPost("/api/tenant", tenant, Tenant.class); + + login(username, password); + for (Queue queue : foundTenantQueues) { + doGet("/api/queues/" + queue.getId()).andExpect(status().isNotFound()); + } + + loginSysAdmin(); + doDelete("/api/tenant/" + tenant.getId().getId().toString()).andExpect(status().isOk()); + } + + private void addQueueConfig(TenantProfile tenantProfile, String queueName) { + TenantProfileQueueConfiguration queueConfiguration = new TenantProfileQueueConfiguration(); + queueConfiguration.setName(queueName); + queueConfiguration.setTopic("tb_rule_engine." + queueName.toLowerCase()); + queueConfiguration.setPollInterval(25); + queueConfiguration.setPartitions(new Random().nextInt(100)); + queueConfiguration.setConsumerPerPartition(true); + queueConfiguration.setPackProcessingTimeout(2000); + SubmitStrategy submitStrategy = new SubmitStrategy(); + submitStrategy.setType(SubmitStrategyType.BURST); + submitStrategy.setBatchSize(1000); + queueConfiguration.setSubmitStrategy(submitStrategy); + ProcessingStrategy processingStrategy = new ProcessingStrategy(); + processingStrategy.setType(ProcessingStrategyType.SKIP_ALL_FAILURES); + processingStrategy.setRetries(3); + processingStrategy.setFailurePercentage(0); + processingStrategy.setPauseBetweenRetries(3); + processingStrategy.setMaxPauseBetweenRetries(3); + queueConfiguration.setProcessingStrategy(processingStrategy); + TenantProfileData profileData = tenantProfile.getProfileData(); + + List configs = profileData.getQueueConfiguration(); + if (configs == null) { + configs = new ArrayList<>(); + } + configs.add(queueConfiguration); + profileData.setQueueConfiguration(configs); + tenantProfile.setProfileData(profileData); + } + + private List getQueuesFromConfig(List queueConfiguration, List queues) { + List result = new ArrayList<>(); + Map queueMap = new HashMap<>(); + for (Queue queue : queues) { + queueMap.put(queue.getName(), queue); + } + + for (TenantProfileQueueConfiguration config : queueConfiguration) { + Queue queue = queueMap.get(config.getName()); + if (queue != null) { + Queue expectedQueue = new Queue(queue.getTenantId(), config); + expectedQueue.setId(queue.getId()); + expectedQueue.setCreatedTime(queue.getCreatedTime()); + result.add(queue); + } + } + return result; + } } 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 f65609beae..b4d65706a2 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 @@ -52,6 +52,7 @@ public class Queue extends SearchTextBasedWithAdditionalInfo implements this.packProcessingTimeout = queueConfiguration.getPackProcessingTimeout(); this.submitStrategy = queueConfiguration.getSubmitStrategy(); this.processingStrategy = queueConfiguration.getProcessingStrategy(); + setAdditionalInfo(queueConfiguration.getAdditionalInfo()); } @Override diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/TenantProfileQueueConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/TenantProfileQueueConfiguration.java index 8d3389af3a..300aab3096 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/TenantProfileQueueConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/TenantProfileQueueConfiguration.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.common.data.tenant.profile; +import com.fasterxml.jackson.databind.JsonNode; import lombok.Data; import org.thingsboard.server.common.data.queue.ProcessingStrategy; import org.thingsboard.server.common.data.queue.SubmitStrategy; @@ -29,4 +30,5 @@ public class TenantProfileQueueConfiguration { private long packProcessingTimeout; private SubmitStrategy submitStrategy; private ProcessingStrategy processingStrategy; + private JsonNode additionalInfo; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/QueueValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/QueueValidator.java index 8257a6e659..c8ce639d40 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/QueueValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/QueueValidator.java @@ -40,12 +40,12 @@ public class QueueValidator extends DataValidator { @Override protected void validateCreate(TenantId tenantId, Queue queue) { - if (queueDao.findQueueByTenantIdAndTopic(tenantId, queue.getTopic()) != null) { - throw new DataValidationException(String.format("Queue with topic: %s already exists!", queue.getTopic())); - } if (queueDao.findQueueByTenantIdAndName(tenantId, queue.getName()) != null) { throw new DataValidationException(String.format("Queue with name: %s already exists!", queue.getName())); } + if (queueDao.findQueueByTenantIdAndTopic(tenantId, queue.getTopic()) != null) { + throw new DataValidationException(String.format("Queue with topic: %s already exists!", queue.getTopic())); + } } @Override diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/BaseTenantProfileServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/BaseTenantProfileServiceTest.java index 210105c73c..575fab55df 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/BaseTenantProfileServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/BaseTenantProfileServiceTest.java @@ -15,6 +15,7 @@ */ package org.thingsboard.server.dao.service; +import com.fasterxml.jackson.databind.node.NullNode; import org.junit.After; import org.junit.Assert; import org.junit.Test; @@ -73,6 +74,7 @@ public abstract class BaseTenantProfileServiceTest extends AbstractServiceTest { mainQueueProcessingStrategy.setPauseBetweenRetries(3); mainQueueProcessingStrategy.setMaxPauseBetweenRetries(3); mainQueueConfiguration.setProcessingStrategy(mainQueueProcessingStrategy); + mainQueueConfiguration.setAdditionalInfo(NullNode.getInstance()); tenantProfile.getProfileData().setQueueConfiguration(Collections.singletonList(mainQueueConfiguration)); TenantProfile savedTenantProfile = tenantProfileService.saveTenantProfile(TenantId.SYS_TENANT_ID, tenantProfile); diff --git a/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java b/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java index fdbbb82464..fe78879125 100644 --- a/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java +++ b/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java @@ -37,7 +37,6 @@ import org.springframework.util.MultiValueMap; import org.springframework.util.StringUtils; import org.springframework.web.client.HttpClientErrorException; import org.springframework.web.client.RestTemplate; -import org.springframework.web.multipart.MultipartFile; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.rest.client.utils.RestJsonConverter; import org.thingsboard.server.common.data.AdminSettings; @@ -91,6 +90,7 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityViewId; import org.thingsboard.server.common.data.id.OAuth2ClientRegistrationTemplateId; import org.thingsboard.server.common.data.id.OtaPackageId; +import org.thingsboard.server.common.data.id.QueueId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleNodeId; import org.thingsboard.server.common.data.id.TbResourceId; @@ -119,6 +119,7 @@ import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.query.EntityCountQuery; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityDataQuery; +import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelationInfo; import org.thingsboard.server.common.data.relation.EntityRelationsQuery; @@ -146,7 +147,6 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Optional; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; import java.util.stream.Collectors; @@ -328,11 +328,11 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { params.put("entityType", entityId.getEntityType().name()); params.put("entityId", entityId.getId().toString()); params.put("fetchOriginator", String.valueOf(fetchOriginator)); - if(searchStatus != null) { + if (searchStatus != null) { params.put("searchStatus", searchStatus.name()); urlSecondPart += "&searchStatus={searchStatus}"; } - if(status != null) { + if (status != null) { params.put("status", status.name()); urlSecondPart += "&status={status}"; } @@ -340,7 +340,7 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { addTimePageLinkToParam(params, pageLink); return restTemplate.exchange( - baseURL + urlSecondPart + "&" + getTimeUrlParams(pageLink), + baseURL + urlSecondPart + "&" + getTimeUrlParams(pageLink), HttpMethod.GET, HttpEntity.EMPTY, new ParameterizedTypeReference>() { @@ -523,12 +523,12 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { public List getAssetsByIds(List assetIds) { return restTemplate.exchange( - baseURL + "/api/assets?assetIds={assetIds}", - HttpMethod.GET, - HttpEntity.EMPTY, - new ParameterizedTypeReference>() { - }, - listIdsToString(assetIds)) + baseURL + "/api/assets?assetIds={assetIds}", + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference>() { + }, + listIdsToString(assetIds)) .getBody(); } @@ -543,7 +543,7 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { public List getAssetTypes() { return restTemplate.exchange(URI.create( - baseURL + "/api/asset/types"), + baseURL + "/api/asset/types"), HttpMethod.GET, HttpEntity.EMPTY, new ParameterizedTypeReference>() { @@ -746,13 +746,13 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { public List getComponentDescriptorsByTypes(List componentTypes, RuleChainType ruleChainType) { return restTemplate.exchange( - baseURL + "/api/components?componentTypes={componentTypes}&ruleChainType={ruleChainType}", - HttpMethod.GET, - HttpEntity.EMPTY, - new ParameterizedTypeReference>() { - }, - listEnumToString(componentTypes), - ruleChainType) + baseURL + "/api/components?componentTypes={componentTypes}&ruleChainType={ruleChainType}", + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference>() { + }, + listEnumToString(componentTypes), + ruleChainType) .getBody(); } @@ -2904,7 +2904,8 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { baseURL + "/api/resource/{resourceId}/download", HttpMethod.GET, HttpEntity.EMPTY, - new ParameterizedTypeReference<>() {}, + new ParameterizedTypeReference<>() { + }, params ); } @@ -2917,7 +2918,8 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { baseURL + "/api/resource/info/{resourceId}", HttpMethod.GET, HttpEntity.EMPTY, - new ParameterizedTypeReference() {}, + new ParameterizedTypeReference() { + }, params ).getBody(); } @@ -2930,7 +2932,8 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { baseURL + "/api/resource/{resourceId}", HttpMethod.GET, HttpEntity.EMPTY, - new ParameterizedTypeReference() {}, + new ParameterizedTypeReference() { + }, params ).getBody(); } @@ -2950,7 +2953,8 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { baseURL + "/api/resource?" + getUrlParams(pageLink), HttpMethod.GET, HttpEntity.EMPTY, - new ParameterizedTypeReference>() {}, + new ParameterizedTypeReference>() { + }, params ).getBody(); } @@ -2967,7 +2971,8 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { baseURL + "/api/otaPackage/{otaPackageId}/download", HttpMethod.GET, HttpEntity.EMPTY, - new ParameterizedTypeReference<>() {}, + new ParameterizedTypeReference<>() { + }, params ); } @@ -2980,7 +2985,8 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { baseURL + "/api/otaPackage/info/{otaPackageId}", HttpMethod.GET, HttpEntity.EMPTY, - new ParameterizedTypeReference() {}, + new ParameterizedTypeReference() { + }, params ).getBody(); } @@ -2993,7 +2999,8 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { baseURL + "/api/otaPackage/{otaPackageId}", HttpMethod.GET, HttpEntity.EMPTY, - new ParameterizedTypeReference() {}, + new ParameterizedTypeReference() { + }, params ).getBody(); } @@ -3004,13 +3011,13 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { return restTemplate.postForEntity(baseURL + "/api/otaPackage?isUrl={isUrl}", otaPackageInfo, OtaPackageInfo.class, params).getBody(); } - public OtaPackageInfo saveOtaPackageData(OtaPackageId otaPackageId, String checkSum, ChecksumAlgorithm checksumAlgorithm, MultipartFile file) throws Exception { + public OtaPackageInfo saveOtaPackageData(OtaPackageId otaPackageId, String checkSum, ChecksumAlgorithm checksumAlgorithm, String fileName, byte[] fileBytes) throws Exception { HttpHeaders header = new HttpHeaders(); header.setContentType(MediaType.MULTIPART_FORM_DATA); MultiValueMap fileMap = new LinkedMultiValueMap<>(); - fileMap.add(HttpHeaders.CONTENT_DISPOSITION, "form-data; name=file; filename=" + file.getName()); - HttpEntity fileEntity = new HttpEntity<>(new ByteArrayResource(file.getBytes()), fileMap); + fileMap.add(HttpHeaders.CONTENT_DISPOSITION, "form-data; name=file; filename=" + fileName); + HttpEntity fileEntity = new HttpEntity<>(new ByteArrayResource(fileBytes), fileMap); MultiValueMap body = new LinkedMultiValueMap<>(); body.add("file", fileEntity); @@ -3021,7 +3028,7 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { params.put("checksumAlgorithm", checksumAlgorithm.name()); String url = "/api/otaPackage/{otaPackageId}?checksumAlgorithm={checksumAlgorithm}"; - if(checkSum != null) { + if (checkSum != null) { url += "&checkSum={checkSum}"; } @@ -3068,6 +3075,39 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { restTemplate.delete(baseURL + "/api/otaPackage/{otaPackageId}", otaPackageId.getId().toString()); } + public PageData getQueuesByServiceType(String serviceType, PageLink pageLink) { + Map params = new HashMap<>(); + params.put("serviceType", serviceType); + addPageLinkToParam(params, pageLink); + + return restTemplate.exchange( + baseURL + "/api/queues?{serviceType}&" + getUrlParams(pageLink), + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference>() { + }, + params + ).getBody(); + } + + public Queue getQueueById(QueueId queueId) { + return restTemplate.exchange( + baseURL + "/api/queue/" + queueId, + HttpMethod.GET, + HttpEntity.EMPTY, + new ParameterizedTypeReference() { + } + ).getBody(); + } + + public Queue saveQueue(Queue queue, String serviceType) { + return restTemplate.postForEntity(baseURL + "/api/queues?serviceType=" + serviceType, queue, Queue.class).getBody(); + } + + public void deleteQueue(QueueId queueId) { + restTemplate.delete(baseURL + "/api/queues/" + queueId); + } + @Deprecated public Optional getAttributes(String accessToken, String clientKeys, String sharedKeys) { Map params = new HashMap<>(); diff --git a/ui-ngx/src/app/modules/home/components/profile/queue/tenant-profile-queues.component.html b/ui-ngx/src/app/modules/home/components/profile/queue/tenant-profile-queues.component.html index 34c43f240c..9339438726 100644 --- a/ui-ngx/src/app/modules/home/components/profile/queue/tenant-profile-queues.component.html +++ b/ui-ngx/src/app/modules/home/components/profile/queue/tenant-profile-queues.component.html @@ -16,33 +16,16 @@ -->
-
- - - -
- - {{ getName(queuesControl.value.name) }} - - - -
-
- - - - -
-
+
+ +
{ }; constructor(private store: Store, @@ -64,11 +66,17 @@ export class TenantProfileDataComponent implements ControlValueAccessor, OnInit this.tenantProfileDataFormGroup = this.fb.group({ configuration: [null, Validators.required] }); - this.tenantProfileDataFormGroup.valueChanges.subscribe(() => { + this.valueChange$ = this.tenantProfileDataFormGroup.valueChanges.subscribe(() => { this.updateModel(); }); } + ngOnDestroy() { + if (this.valueChange$) { + this.valueChange$.unsubscribe(); + } + } + setDisabledState(isDisabled: boolean): void { this.disabled = isDisabled; if (this.disabled) { @@ -87,7 +95,7 @@ export class TenantProfileDataComponent implements ControlValueAccessor, OnInit if (this.tenantProfileDataFormGroup.valid) { tenantProfileData = this.tenantProfileDataFormGroup.getRawValue(); } - this.propagateChange(tenantProfileData); + this.propagateChange(tenantProfileData.configuration); } } diff --git a/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.scss b/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.scss index c76509e9e3..f9940b3dae 100644 --- a/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.scss +++ b/ui-ngx/src/app/modules/home/components/profile/tenant-profile.component.scss @@ -35,6 +35,10 @@ width: fit-content; } } + + .mat-expansion-panel-header { + height: 48px; + } .expansion-panel-block { padding-bottom: 16px; } 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 6a57ee57d0..ed88fe9fba 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 @@ -23,6 +23,7 @@ import { ActionNotificationShow } from '@app/core/notification/notification.acti import { TranslateService } from '@ngx-translate/core'; import { EntityTableConfig } from '@home/models/entity/entities-table-config.models'; import { EntityComponent } from '../entity/entity.component'; +import { guid } from '@core/utils'; @Component({ selector: 'tb-tenant-profile', @@ -54,6 +55,7 @@ export class TenantProfileComponent extends EntityComponent { buildForm(entity: TenantProfile): FormGroup { const mainQueue = [ { + id: guid(), consumerPerPartition: true, name: 'Main', packProcessingTimeout: 2000, @@ -70,7 +72,10 @@ export class TenantProfileComponent extends EntityComponent { batchSize: 1000, type: 'BURST' }, - topic: 'tb_rule_engine.main' + topic: 'tb_rule_engine.main', + additionalInfo: { + description: '' + } } ]; const formGroup = this.fb.group( 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 e13faf0f43..920ad2f162 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 @@ -15,161 +15,182 @@ limitations under the License. --> - -
- - admin.queue-name - - - {{ 'queue.name-required' | translate }} - - - - queue.poll-interval - - - {{ 'queue.poll-interval-required' | translate }} - - - {{ 'queue.poll-interval-min-value' | translate }} - - - - queue.partitions - - - {{ 'queue.partitions-required' | translate }} - - - {{ 'queue.partitions-min-value' | translate }} - - - -
{{ 'queue.consumer-per-partition' | translate }}
-
{{'queue.consumer-per-partition-hint' | translate}}
-
- - queue.processing-timeout - - - {{ 'queue.pack-processing-timeout-required' | translate }} - - - {{ 'queue.pack-processing-timeout-min-value' | translate }} - - - - - - - queue.submit-strategy - - - -
- - queue.submit-strategy - - - {{ strategy }} - - - - {{ 'queue.submit-strategy-type-required' | translate }} - - - - queue.batch-size - - - {{ 'queue.batch-size-required' | translate }} - - - {{ 'queue.batch-size-min-value' | translate }} - - -
-
-
- - - - queue.processing-strategy - - - -
- - queue.processing-strategy - - - {{ strategy }} - - - - {{ 'queue.processing-strategy-type-required' | translate }} - - - - queue.retries - - - {{ 'queue.retries-required' | translate }} - - - {{ 'queue.retries-min-value' | translate }} - - - - queue.failure-percentage - - - {{ 'queue.failure-percentage-required' | translate }} - - - {{ 'queue.failure-percentage-min-value' | translate }} - - - {{ 'queue.failure-percentage-max-value' | translate }} - - - - queue.pause-between-retries - - - {{ 'queue.pause-between-retries-required' | translate }} - - - {{ 'queue.pause-between-retries-min-value' | translate }} - - - - queue.max-pause-between-retries - - - {{ 'queue.max-pause-between-retries-required' | translate }} - - - {{ 'queue.max-pause-between-retries-min-value' | translate }} - - -
-
-
-
- - queue.description - - -
+ + +
+ + {{ queueTitle }} + + + +
+
+ +
+ + admin.queue-name + + + {{ 'queue.name-required' | translate }} + + + {{ 'queue.name-unique' | translate }} + + + + queue.poll-interval + + + {{ 'queue.poll-interval-required' | translate }} + + + {{ 'queue.poll-interval-min-value' | translate }} + + + + queue.partitions + + + {{ 'queue.partitions-required' | translate }} + + + {{ 'queue.partitions-min-value' | translate }} + + + +
{{ 'queue.consumer-per-partition' | translate }}
+
{{'queue.consumer-per-partition-hint' | translate}}
+
+ + queue.processing-timeout + + + {{ 'queue.pack-processing-timeout-required' | translate }} + + + {{ 'queue.pack-processing-timeout-min-value' | translate }} + + + + + + + queue.submit-strategy + + + +
+ + queue.submit-strategy + + + {{ strategy }} + + + + {{ 'queue.submit-strategy-type-required' | translate }} + + + + queue.batch-size + + + {{ 'queue.batch-size-required' | translate }} + + + {{ 'queue.batch-size-min-value' | translate }} + + +
+
+
+ + + + queue.processing-strategy + + + +
+ + queue.processing-strategy + + + {{ strategy }} + + + + {{ 'queue.processing-strategy-type-required' | translate }} + + + + queue.retries + + + {{ 'queue.retries-required' | translate }} + + + {{ 'queue.retries-min-value' | translate }} + + + + queue.failure-percentage + + + {{ 'queue.failure-percentage-required' | translate }} + + + {{ 'queue.failure-percentage-min-value' | translate }} + + + {{ 'queue.failure-percentage-max-value' | translate }} + + + + queue.pause-between-retries + + + {{ 'queue.pause-between-retries-required' | translate }} + + + {{ 'queue.pause-between-retries-min-value' | translate }} + + + + queue.max-pause-between-retries + + + {{ 'queue.max-pause-between-retries-required' | translate }} + + + {{ 'queue.max-pause-between-retries-min-value' | translate }} + + +
+
+
+
+ + 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 38552bff25..6a437799b5 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 @@ -14,7 +14,7 @@ /// limitations under the License. /// -import { Component, forwardRef, Input, OnInit } from '@angular/core'; +import { Component, forwardRef, Input, OnInit, Output, EventEmitter, OnDestroy } from '@angular/core'; import { ControlValueAccessor, FormBuilder, @@ -29,6 +29,7 @@ import { MatDialog } from '@angular/material/dialog'; import { UtilsService } from '@core/services/utils.service'; import { QueueInfo, QueueProcessingStrategyTypes, QueueSubmitStrategyTypes } from '@shared/models/queue.models'; import { isDefinedAndNotNull } from '@core/utils'; +import { Subscription } from 'rxjs'; @Component({ selector: 'tb-queue-form', @@ -47,7 +48,7 @@ import { isDefinedAndNotNull } from '@core/utils'; } ] }) -export class QueueFormComponent implements ControlValueAccessor, OnInit, Validator { +export class QueueFormComponent implements ControlValueAccessor, OnInit, OnDestroy, Validator { @Input() disabled: boolean; @@ -55,20 +56,28 @@ export class QueueFormComponent implements ControlValueAccessor, OnInit, Validat @Input() newQueue = false; + @Input() + mainQueue = false; + @Input() systemQueue = false; - private modelValue: QueueInfo; + @Input() + expanded = false; - queueFormGroup: FormGroup; + @Output() + removeQueue = new EventEmitter(); + queueFormGroup: FormGroup; submitStrategies: string[] = []; processingStrategies: string[] = []; - + queueTitle = ''; hideBatchSize = false; + private modelValue: QueueInfo; private propagateChange = null; private propagateChangePending = false; + private valueChange$: Subscription = null; constructor(private dialog: MatDialog, private utils: UtilsService, @@ -114,10 +123,13 @@ export class QueueFormComponent implements ControlValueAccessor, OnInit, Validat description: [''] }) }); - this.queueFormGroup.valueChanges.subscribe(() => { + this.valueChange$ = this.queueFormGroup.valueChanges.subscribe(() => { this.updateModel(); }); - this.queueFormGroup.get('name').valueChanges.subscribe((value) => this.queueFormGroup.patchValue({topic: `tb_rule_engine.${value}`})); + this.queueFormGroup.get('name').valueChanges.subscribe((value) => { + this.queueFormGroup.patchValue({topic: `tb_rule_engine.${value}`}); + this.queueTitle = this.utils.customTranslation(value, value); + }); this.queueFormGroup.get('submitStrategy').get('type').valueChanges.subscribe(() => { this.submitStrategyTypeChanged(); }); @@ -128,6 +140,13 @@ export class QueueFormComponent implements ControlValueAccessor, OnInit, Validat } } + ngOnDestroy() { + if (this.valueChange$) { + this.valueChange$.unsubscribe(); + this.valueChange$ = null; + } + } + setDisabledState(isDisabled: boolean): void { this.disabled = isDisabled; if (this.disabled) { @@ -141,6 +160,10 @@ export class QueueFormComponent implements ControlValueAccessor, OnInit, Validat writeValue(value: QueueInfo): void { this.propagateChangePending = false; this.modelValue = value; + if (!this.modelValue.name) { + this.expanded = true; + } + this.queueTitle = this.utils.customTranslation(value.name, value.name); if (isDefinedAndNotNull(this.modelValue)) { this.queueFormGroup.patchValue(this.modelValue, {emitEvent: false}); } diff --git a/ui-ngx/src/app/modules/home/pages/tenant-profile/tenant-profiles-table-config.resolver.ts b/ui-ngx/src/app/modules/home/pages/tenant-profile/tenant-profiles-table-config.resolver.ts index 3b1773842d..7484d8b403 100644 --- a/ui-ngx/src/app/modules/home/pages/tenant-profile/tenant-profiles-table-config.resolver.ts +++ b/ui-ngx/src/app/modules/home/pages/tenant-profile/tenant-profiles-table-config.resolver.ts @@ -33,6 +33,8 @@ import { TenantProfileComponent } from '../../components/profile/tenant-profile. import { TenantProfileTabsComponent } from './tenant-profile-tabs.component'; import { DialogService } from '@core/services/dialog.service'; import { ImportExportService } from '@home/components/import-export/import-export.service'; +import { map } from 'rxjs/operators'; +import { guid } from '@core/utils'; @Injectable() export class TenantProfilesTableConfigResolver implements Resolve> { @@ -84,7 +86,12 @@ export class TenantProfilesTableConfigResolver implements Resolve this.translate.instant('tenant-profile.delete-tenant-profiles-text'); this.config.entitiesFetchFunction = pageLink => this.tenantProfileService.getTenantProfiles(pageLink); - this.config.loadEntity = id => this.tenantProfileService.getTenantProfile(id.id); + this.config.loadEntity = id => this.tenantProfileService.getTenantProfile(id.id).pipe( + map(tenantProfile => ({ + ...tenantProfile, + profileData: {...tenantProfile.profileData, queueConfiguration: this.addId(tenantProfile.profileData.queueConfiguration)}, + })) + ); this.config.saveEntity = tenantProfile => this.tenantProfileService.saveTenantProfile(tenantProfile); this.config.deleteEntity = id => this.tenantProfileService.deleteTenantProfile(id.id); this.config.onEntityAction = action => this.onTenantProfileAction(action); @@ -93,6 +100,15 @@ export class TenantProfilesTableConfigResolver implements Resolve { + value.id = guid(); + queuesWithId.push(value); + }); + return queuesWithId; + } + resolve(): EntityTableConfig { this.config.tableTitle = this.translate.instant('tenant-profile.tenant-profiles'); diff --git a/ui-ngx/src/app/shared/components/queue/queue-autocomplete.component.html b/ui-ngx/src/app/shared/components/queue/queue-autocomplete.component.html index 471d191b61..c8af29fbd4 100644 --- a/ui-ngx/src/app/shared/components/queue/queue-autocomplete.component.html +++ b/ui-ngx/src/app/shared/components/queue/queue-autocomplete.component.html @@ -15,7 +15,7 @@ limitations under the License. --> - + - + [displayWith]="displayQueueFn" + > + - {{getDescription(queue)}} + {{getDescription(queue)}}
diff --git a/ui-ngx/src/app/shared/components/queue/queue-autocomplete.component.scss b/ui-ngx/src/app/shared/components/queue/queue-autocomplete.component.scss new file mode 100644 index 0000000000..45a4ab9808 --- /dev/null +++ b/ui-ngx/src/app/shared/components/queue/queue-autocomplete.component.scss @@ -0,0 +1,29 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +::ng-deep { + .queue-option { + .mat-option-text { + display: inline; + } + .queue-option-description { + display: block; + overflow: hidden; + text-overflow: ellipsis; + white-space: nowrap; + } + } +} diff --git a/ui-ngx/src/app/shared/components/queue/queue-autocomplete.component.ts b/ui-ngx/src/app/shared/components/queue/queue-autocomplete.component.ts index 0caa487d66..897e041649 100644 --- a/ui-ngx/src/app/shared/components/queue/queue-autocomplete.component.ts +++ b/ui-ngx/src/app/shared/components/queue/queue-autocomplete.component.ts @@ -36,7 +36,7 @@ import { emptyPageData } from '@shared/models/page/page-data'; @Component({ selector: 'tb-queue-autocomplete', templateUrl: './queue-autocomplete.component.html', - styleUrls: [], + styleUrls: ['./queue-autocomplete.component.scss'], providers: [{ provide: NG_VALUE_ACCESSOR, useExisting: forwardRef(() => QueueAutocompleteComponent), @@ -207,7 +207,10 @@ export class QueueAutocompleteComponent implements ControlValueAccessor, OnInit getDescription(value) { return value.additionalInfo?.description ? value.additionalInfo.description : - `Submit Strategy: ${value.submitStrategy.type}, Processing Strategy: ${value.processingStrategy.type}`; + this.translate.instant( + 'queue.alt-description', + {submitStrategy: value.submitStrategy.type, processingStrategy: value.processingStrategy.type} + ); } clear() { diff --git a/ui-ngx/src/app/shared/models/queue.models.ts b/ui-ngx/src/app/shared/models/queue.models.ts index 86bd99a19d..6c67b4db1b 100644 --- a/ui-ngx/src/app/shared/models/queue.models.ts +++ b/ui-ngx/src/app/shared/models/queue.models.ts @@ -43,6 +43,7 @@ export enum QueueProcessingStrategyTypes { } export interface QueueInfo extends BaseData { + generatedId?: string; name: string; packProcessingTimeout: number; partitions: number; diff --git a/ui-ngx/src/app/shared/models/tenant.model.ts b/ui-ngx/src/app/shared/models/tenant.model.ts index 14c0d232df..33d1367c01 100644 --- a/ui-ngx/src/app/shared/models/tenant.model.ts +++ b/ui-ngx/src/app/shared/models/tenant.model.ts @@ -98,7 +98,7 @@ export function createTenantProfileConfiguration(type: TenantProfileType): Tenan export interface TenantProfileData { configuration: TenantProfileConfiguration; - queueConfiguration?: QueueInfo; + queueConfiguration?: Array; } export interface TenantProfile extends BaseData { 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 49b4f82c17..db52c04ac1 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -2794,6 +2794,7 @@ "select-name": "Select queue name", "name": "Name", "name-required": "Queue name is required!", + "name-unique": "Queue name is not unique!", "queue-required": "Queue is required!", "topic-required": "Queue topic is required!", "poll-interval-required": "Poll interval is required!", @@ -2840,7 +2841,8 @@ "delete": "Delete queue", "copyId": "Copy queue Id", "idCopiedMessage": "Queue Id has been copied to clipboard", - "description": "Description" + "description": "Description", + "alt-description": "Submit Strategy: {{submitStrategy}}, Processing Strategy: {{processingStrategy}}" }, "tenant": { "tenant": "Tenant",