Browse Source

Merge branch 'develop/3.5' into improvements/notification-system

pull/8265/head
ViacheslavKlimov 4 years ago
parent
commit
14acf7922c
  1. 85
      application/src/main/java/org/thingsboard/server/controller/AdminController.java
  2. 2
      application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java
  3. 216
      application/src/main/java/org/thingsboard/server/service/system/DefaultSystemInfoService.java
  4. 25
      application/src/main/java/org/thingsboard/server/service/system/SystemInfoService.java
  5. 5
      application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java
  6. 52
      application/src/test/java/org/thingsboard/server/controller/BaseEntityQueryControllerTest.java
  7. 372
      application/src/test/java/org/thingsboard/server/controller/BaseHomePageApiTest.java
  8. 23
      application/src/test/java/org/thingsboard/server/controller/sql/HomePageApiSqlTest.java
  9. 11
      common/cluster-api/src/main/proto/queue.proto
  10. 27
      common/data/src/main/java/org/thingsboard/server/common/data/FeaturesInfo.java
  11. 29
      common/data/src/main/java/org/thingsboard/server/common/data/SystemInfo.java
  12. 43
      common/data/src/main/java/org/thingsboard/server/common/data/SystemInfoData.java
  13. 5
      common/data/src/main/java/org/thingsboard/server/common/data/UpdateMessage.java
  14. 38
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java
  15. 8
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java
  16. 11
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java
  17. 3
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java
  18. 1
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/PartitionService.java
  19. 3
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java
  20. 32
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java
  21. 4
      common/util/pom.xml
  22. 108
      common/util/src/main/java/org/thingsboard/common/util/SystemUtil.java
  23. 2
      dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java
  24. 2
      dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java
  25. 6
      dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java
  26. 6
      pom.xml
  27. 7
      rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java

85
application/src/main/java/org/thingsboard/server/controller/AdminController.java

@ -21,29 +21,41 @@ import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import io.swagger.annotations.ApiOperation;
import io.swagger.annotations.ApiParam;
import org.springframework.beans.factory.annotation.Autowired;
import lombok.RequiredArgsConstructor;
import org.springframework.context.annotation.Lazy;
import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestMethod;
import org.springframework.web.bind.annotation.ResponseBody;
import org.springframework.web.bind.annotation.ResponseStatus;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.context.request.async.DeferredResult;
import org.thingsboard.rule.engine.api.MailService;
import org.thingsboard.rule.engine.api.SmsService;
import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.FeaturesInfo;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.SystemInfo;
import org.thingsboard.server.common.data.UpdateMessage;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.security.model.JwtPair;
import org.thingsboard.server.common.data.security.model.JwtSettings;
import org.thingsboard.server.common.data.security.model.SecuritySettings;
import org.thingsboard.server.common.data.sms.config.TestSmsRequest;
import org.thingsboard.server.common.data.sync.vc.AutoCommitSettings;
import org.thingsboard.server.common.data.sync.vc.RepositorySettings;
import org.thingsboard.server.common.data.sync.vc.RepositorySettingsInfo;
import org.thingsboard.server.common.data.security.model.JwtSettings;
import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService;
import org.thingsboard.server.dao.settings.AdminSettingsService;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.common.data.security.model.JwtPair;
import org.thingsboard.server.service.security.auth.jwt.settings.JwtSettingsService;
import org.thingsboard.server.service.security.model.SecurityUser;
import org.thingsboard.server.service.security.model.token.JwtTokenFactory;
import org.thingsboard.server.service.security.permission.Operation;
@ -51,43 +63,30 @@ import org.thingsboard.server.service.security.permission.Resource;
import org.thingsboard.server.service.security.system.SystemSecurityService;
import org.thingsboard.server.service.sync.vc.EntitiesVersionControlService;
import org.thingsboard.server.service.sync.vc.autocommit.TbAutoCommitSettingsService;
import org.thingsboard.server.service.system.SystemInfoService;
import org.thingsboard.server.service.update.UpdateService;
import static org.thingsboard.server.controller.ControllerConstants.*;
import static org.thingsboard.server.controller.ControllerConstants.SYSTEM_AUTHORITY_PARAGRAPH;
import static org.thingsboard.server.controller.ControllerConstants.TENANT_AUTHORITY_PARAGRAPH;
@RestController
@TbCoreComponent
@RequestMapping("/api/admin")
@RequiredArgsConstructor
public class AdminController extends BaseController {
@Autowired
private MailService mailService;
@Autowired
private SmsService smsService;
@Autowired
private AdminSettingsService adminSettingsService;
@Autowired
private SystemSecurityService systemSecurityService;
private final MailService mailService;
private final SmsService smsService;
private final AdminSettingsService adminSettingsService;
private final SystemSecurityService systemSecurityService;
@Lazy
@Autowired
private JwtSettingsService jwtSettingsService;
private final JwtSettingsService jwtSettingsService;
@Lazy
@Autowired
private JwtTokenFactory tokenFactory;
@Autowired
private EntitiesVersionControlService versionControlService;
@Autowired
private TbAutoCommitSettingsService autoCommitSettingsService;
@Autowired
private UpdateService updateService;
private final JwtTokenFactory tokenFactory;
private final EntitiesVersionControlService versionControlService;
private final TbAutoCommitSettingsService autoCommitSettingsService;
private final UpdateService updateService;
private final SystemInfoService systemInfoService;
@ApiOperation(value = "Get the Administration Settings object using key (getAdminSettings)",
notes = "Get the Administration Settings object using specified string key. Referencing non-existing key will cause an error." + SYSTEM_AUTHORITY_PARAGRAPH)
@ -109,7 +108,6 @@ public class AdminController extends BaseController {
}
}
@ApiOperation(value = "Get the Administration Settings object using key (getAdminSettings)",
notes = "Creates or Updates the Administration Settings. Platform generates random Administration Settings Id during settings creation. " +
"The Administration Settings Id will be present in the response. Specify the Administration Settings Id when you would like to update the Administration Settings. " +
@ -318,7 +316,6 @@ public class AdminController extends BaseController {
}
}
@ApiOperation(value = "Check repository access (checkRepositoryAccess)",
notes = "Attempts to check repository access. " + TENANT_AUTHORITY_PARAGRAPH)
@PreAuthorize("hasAuthority('TENANT_ADMIN')")
@ -399,4 +396,24 @@ public class AdminController extends BaseController {
}
}
@ApiOperation(value = "Get system info (getSystemInfo)",
notes = "Get main information about system. "
+ SYSTEM_AUTHORITY_PARAGRAPH)
@PreAuthorize("hasAuthority('SYS_ADMIN')")
@RequestMapping(value = "/systemInfo", method = RequestMethod.GET)
@ResponseBody
public SystemInfo getSystemInfo() throws ThingsboardException {
return systemInfoService.getSystemInfo();
}
@ApiOperation(value = "Get features info (getFeaturesInfo)",
notes = "Get information about enabled/disabled features. "
+ SYSTEM_AUTHORITY_PARAGRAPH)
@PreAuthorize("hasAuthority('SYS_ADMIN')")
@RequestMapping(value = "/featuresInfo", method = RequestMethod.GET)
@ResponseBody
public FeaturesInfo getFeaturesInfo() {
return systemInfoService.getFeaturesInfo();
}
}

2
application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java

@ -56,7 +56,7 @@ public class EntityQueryController extends BaseController {
private static final int MAX_PAGE_SIZE = 100;
@ApiOperation(value = "Count Entities by Query", notes = ENTITY_COUNT_QUERY_DESCRIPTION)
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')")
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/entitiesQuery/count", method = RequestMethod.POST)
@ResponseBody
public long countEntitiesByQuery(

216
application/src/main/java/org/thingsboard/server/service/system/DefaultSystemInfoService.java

@ -0,0 +1,216 @@
/**
* Copyright © 2016-2023 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.
*/
package org.thingsboard.server.service.system;
import com.fasterxml.jackson.databind.JsonNode;
import com.google.common.util.concurrent.FutureCallback;
import com.google.protobuf.ProtocolStringList;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.ApiUsageState;
import org.thingsboard.server.common.data.FeaturesInfo;
import org.thingsboard.server.common.data.SystemInfo;
import org.thingsboard.server.common.data.SystemInfoData;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.DoubleDataEntry;
import org.thingsboard.server.common.data.kv.JsonDataEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.stats.TbApiUsageStateClient;
import org.thingsboard.server.dao.oauth2.OAuth2Service;
import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.dao.settings.AdminSettingsService;
import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo;
import org.thingsboard.server.queue.discovery.DiscoveryService;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.discovery.TbApplicationEventListener;
import org.thingsboard.server.queue.discovery.TbServiceInfoProvider;
import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService;
import javax.annotation.Nullable;
import javax.annotation.PreDestroy;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import static org.thingsboard.common.util.SystemUtil.getCpuUsage;
import static org.thingsboard.common.util.SystemUtil.getFreeDiscSpace;
import static org.thingsboard.common.util.SystemUtil.getFreeMemory;
import static org.thingsboard.common.util.SystemUtil.getMemoryUsage;
import static org.thingsboard.common.util.SystemUtil.getTotalCpuUsage;
import static org.thingsboard.common.util.SystemUtil.getTotalDiscSpace;
import static org.thingsboard.common.util.SystemUtil.getTotalMemory;
@TbCoreComponent
@Service
@RequiredArgsConstructor
@Slf4j
public class DefaultSystemInfoService extends TbApplicationEventListener<PartitionChangeEvent> implements SystemInfoService {
public static final FutureCallback<Integer> CALLBACK = new FutureCallback<>() {
@Override
public void onSuccess(@Nullable Integer result) {
}
@Override
public void onFailure(Throwable t) {
log.warn("Failed to persist system info", t);
}
};
private final TbServiceInfoProvider serviceInfoProvider;
private final PartitionService partitionService;
private final DiscoveryService discoveryService;
private final TelemetrySubscriptionService telemetryService;
private final TbApiUsageStateClient apiUsageStateClient;
private final AdminSettingsService adminSettingsService;
private final OAuth2Service oAuth2Service;
private volatile ScheduledExecutorService scheduler;
@Override
protected void onTbApplicationEvent(PartitionChangeEvent partitionChangeEvent) {
if (ServiceType.TB_CORE.equals(partitionChangeEvent.getServiceType())) {
boolean myPartition = partitionService.resolve(ServiceType.TB_CORE, TenantId.SYS_TENANT_ID, TenantId.SYS_TENANT_ID).isMyPartition();
synchronized (this) {
if (myPartition) {
if (scheduler == null) {
scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("tb-system-info-scheduler"));
scheduler.scheduleAtFixedRate(this::saveCurrentSystemInfo, 0, 1, TimeUnit.MINUTES);
}
} else {
destroy();
}
}
}
}
@Override
public SystemInfo getSystemInfo() {
SystemInfo systemInfo = new SystemInfo();
ServiceInfo serviceInfo = serviceInfoProvider.getServiceInfoWithCurrentSystemInfo();
if (discoveryService.isMonolith()) {
systemInfo.setMonolith(true);
systemInfo.setSystemData(Collections.singletonList(createSystemInfoData(serviceInfo)));
} else {
systemInfo.setSystemData(getSystemData(serviceInfo));
}
return systemInfo;
}
protected void saveCurrentSystemInfo() {
if (discoveryService.isMonolith()) {
saveCurrentMonolithSystemInfo();
} else {
saveCurrentClusterSystemInfo();
}
}
@Override
public FeaturesInfo getFeaturesInfo() {
FeaturesInfo featuresInfo = new FeaturesInfo();
featuresInfo.setEmailEnabled(isEmailEnabled());
featuresInfo.setSmsEnabled(adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, "sms") != null);
featuresInfo.setOauthEnabled(oAuth2Service.findOAuth2Info().isEnabled());
featuresInfo.setTwoFaEnabled(adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, "twoFaSettings") != null);
featuresInfo.setNotificationEnabled(adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, "notifications") != null);
return featuresInfo;
}
private boolean isEmailEnabled() {
AdminSettings mailSettings = adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, "mail");
if (mailSettings != null) {
JsonNode mailFrom = mailSettings.getJsonValue().get("mailFrom");
if (mailFrom != null) {
return DataValidator.doValidateEmail(mailFrom.asText());
}
}
return false;
}
private void saveCurrentClusterSystemInfo() {
long ts = System.currentTimeMillis();
List<SystemInfoData> clusterSystemData = getSystemData(serviceInfoProvider.getServiceInfoWithCurrentSystemInfo());
BasicTsKvEntry clusterDataKv = new BasicTsKvEntry(ts, new JsonDataEntry("clusterSystemData", JacksonUtil.toString(clusterSystemData)));
doSave(Collections.singletonList(clusterDataKv));
}
private void saveCurrentMonolithSystemInfo() {
long ts = System.currentTimeMillis();
List<TsKvEntry> tsList = new ArrayList<>();
getMemoryUsage().ifPresent(v -> tsList.add(new BasicTsKvEntry(ts, new LongDataEntry("memoryUsage", v))));
getTotalMemory().ifPresent(v -> tsList.add(new BasicTsKvEntry(ts, new LongDataEntry("totalMemory", v))));
getFreeMemory().ifPresent(v -> tsList.add(new BasicTsKvEntry(ts, new LongDataEntry("freeMemory", v))));
getCpuUsage().ifPresent(v -> tsList.add(new BasicTsKvEntry(ts, new DoubleDataEntry("cpuUsage", v))));
getTotalCpuUsage().ifPresent(v -> tsList.add(new BasicTsKvEntry(ts, new DoubleDataEntry("totalCpuUsage", v))));
getFreeDiscSpace().ifPresent(v -> tsList.add(new BasicTsKvEntry(ts, new LongDataEntry("freeDiscSpace", v))));
getTotalDiscSpace().ifPresent(v -> tsList.add(new BasicTsKvEntry(ts, new LongDataEntry("totalDiscSpace", v))));
doSave(tsList);
}
private void doSave(List<TsKvEntry> telemetry) {
ApiUsageState apiUsageState = apiUsageStateClient.getApiUsageState(TenantId.SYS_TENANT_ID);
telemetryService.saveAndNotifyInternal(TenantId.SYS_TENANT_ID, apiUsageState.getId(), telemetry, CALLBACK);
}
private List<SystemInfoData> getSystemData(ServiceInfo serviceInfo) {
List<SystemInfoData> clusterSystemData = new ArrayList<>();
clusterSystemData.add(createSystemInfoData(serviceInfo));
this.discoveryService.getOtherServers()
.stream()
.map(this::createSystemInfoData)
.forEach(clusterSystemData::add);
return clusterSystemData;
}
private SystemInfoData createSystemInfoData(ServiceInfo serviceInfo) {
ProtocolStringList serviceTypes = serviceInfo.getServiceTypesList();
SystemInfoData infoData = new SystemInfoData();
infoData.setServiceId(serviceInfo.getServiceId());
infoData.setServiceType(serviceTypes.size() > 1 ? "MONOLITH" : serviceTypes.get(0));
infoData.setMemoryUsage(serviceInfo.getSystemInfo().getMemoryUsage());
infoData.setTotalMemory(serviceInfo.getSystemInfo().getTotalMemory());
infoData.setFreeMemory(serviceInfo.getSystemInfo().getFreeMemory());
infoData.setCpuUsage(serviceInfo.getSystemInfo().getCpuUsage());
infoData.setTotalCpuUsage(serviceInfo.getSystemInfo().getTotalCpuUsage());
infoData.setFreeDiscSpace(serviceInfo.getSystemInfo().getFreeDiscSpace());
infoData.setTotalDiscSpace(serviceInfo.getSystemInfo().getTotalDiscSpace());
return infoData;
}
@PreDestroy
private void destroy() {
if (scheduler != null) {
scheduler.shutdownNow();
scheduler = null;
}
}
}

25
application/src/main/java/org/thingsboard/server/service/system/SystemInfoService.java

@ -0,0 +1,25 @@
/**
* Copyright © 2016-2023 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.
*/
package org.thingsboard.server.service.system;
import org.thingsboard.server.common.data.FeaturesInfo;
import org.thingsboard.server.common.data.SystemInfo;
public interface SystemInfoService {
SystemInfo getSystemInfo();
FeaturesInfo getFeaturesInfo();
}

5
application/src/main/java/org/thingsboard/server/service/update/DefaultUpdateService.java

@ -72,7 +72,7 @@ public class DefaultUpdateService implements UpdateService {
@PostConstruct
private void init() {
updateMessage = new UpdateMessage("", false);
updateMessage = new UpdateMessage("", false, "");
if (updatesEnabled) {
try {
platform = System.getProperty("platform", "unknown");
@ -131,7 +131,8 @@ public class DefaultUpdateService implements UpdateService {
UpdateMessage prevUpdateMessage = updateMessage;
updateMessage = new UpdateMessage(
response.get("message").asText(),
response.get("updateAvailable").asBoolean()
response.get("updateAvailable").asBoolean(),
version
);
if (updateMessage.isUpdateAvailable() && !updateMessage.equals(prevUpdateMessage)) {
notificationRuleProcessingService.process(NewPlatformVersionTrigger.builder()

52
application/src/test/java/org/thingsboard/server/controller/BaseEntityQueryControllerTest.java

@ -92,7 +92,7 @@ public abstract class BaseEntityQueryControllerTest extends AbstractControllerTe
}
@Test
public void testCountEntitiesByQuery() throws Exception {
public void testTenantCountEntitiesByQuery() throws Exception {
List<Device> devices = new ArrayList<>();
for (int i = 0; i < 97; i++) {
Device device = new Device();
@ -139,6 +139,56 @@ public abstract class BaseEntityQueryControllerTest extends AbstractControllerTe
Assert.assertEquals(97, count2.longValue());
}
@Test
public void testSysAdminCountEntitiesByQuery() throws Exception {
List<Device> devices = new ArrayList<>();
for (int i = 0; i < 97; i++) {
Device device = new Device();
device.setName("Device" + i);
device.setType("default");
device.setLabel("testLabel" + (int) (Math.random() * 1000));
devices.add(doPost("/api/device", device, Device.class));
Thread.sleep(1);
}
DeviceTypeFilter filter = new DeviceTypeFilter();
filter.setDeviceType("default");
filter.setDeviceNameFilter("");
loginSysAdmin();
EntityCountQuery countQuery = new EntityCountQuery(filter);
Long count = doPostWithResponse("/api/entitiesQuery/count", countQuery, Long.class);
Assert.assertEquals(97, count.longValue());
filter.setDeviceType("unknown");
count = doPostWithResponse("/api/entitiesQuery/count", countQuery, Long.class);
Assert.assertEquals(0, count.longValue());
filter.setDeviceType("default");
filter.setDeviceNameFilter("Device1");
count = doPostWithResponse("/api/entitiesQuery/count", countQuery, Long.class);
Assert.assertEquals(11, count.longValue());
EntityListFilter entityListFilter = new EntityListFilter();
entityListFilter.setEntityType(EntityType.DEVICE);
entityListFilter.setEntityList(devices.stream().map(Device::getId).map(DeviceId::toString).collect(Collectors.toList()));
countQuery = new EntityCountQuery(entityListFilter);
count = doPostWithResponse("/api/entitiesQuery/count", countQuery, Long.class);
Assert.assertEquals(97, count.longValue());
EntityTypeFilter filter2 = new EntityTypeFilter();
filter2.setEntityType(EntityType.DEVICE);
EntityCountQuery countQuery2 = new EntityCountQuery(filter2);
Long count2 = doPostWithResponse("/api/entitiesQuery/count", countQuery2, Long.class);
Assert.assertEquals(97, count2.longValue());
}
@Test
public void testSimpleFindEntityDataByQuery() throws Exception {
List<Device> devices = new ArrayList<>();

372
application/src/test/java/org/thingsboard/server/controller/BaseHomePageApiTest.java

@ -0,0 +1,372 @@
/**
* Copyright © 2016-2023 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.
*/
package org.thingsboard.server.controller;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.collect.Lists;
import lombok.extern.slf4j.Slf4j;
import org.junit.Assert;
import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.ApiUsageState;
import org.thingsboard.server.common.data.Customer;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.FeaturesInfo;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.oauth2.MapperType;
import org.thingsboard.server.common.data.oauth2.OAuth2CustomMapperConfig;
import org.thingsboard.server.common.data.oauth2.OAuth2DomainInfo;
import org.thingsboard.server.common.data.oauth2.OAuth2Info;
import org.thingsboard.server.common.data.oauth2.OAuth2MapperConfig;
import org.thingsboard.server.common.data.oauth2.OAuth2ParamsInfo;
import org.thingsboard.server.common.data.oauth2.OAuth2RegistrationInfo;
import org.thingsboard.server.common.data.oauth2.SchemeType;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.query.ApiUsageStateFilter;
import org.thingsboard.server.common.data.query.EntityCountQuery;
import org.thingsboard.server.common.data.query.EntityData;
import org.thingsboard.server.common.data.query.EntityTypeFilter;
import org.thingsboard.server.common.data.query.TsValue;
import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.common.stats.TbApiUsageStateClient;
import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityCountCmd;
import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityCountUpdate;
import org.thingsboard.server.service.ws.telemetry.cmd.v2.EntityDataUpdate;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@Slf4j
public abstract class BaseHomePageApiTest extends AbstractControllerTest {
@Autowired
private TbApiUsageStateClient apiUsageStateClient;
//For system administrator
@Test
public void testTenantsCountWsCmd() throws Exception {
loginSysAdmin();
List<Tenant> tenants = new ArrayList<>();
for (int i = 0; i < 100; i++) {
Tenant tenant = new Tenant();
tenant.setTitle("tenant" + i);
tenants.add(doPost("/api/tenant", tenant, Tenant.class));
}
EntityTypeFilter ef = new EntityTypeFilter();
ef.setEntityType(EntityType.TENANT);
EntityCountCmd cmd = new EntityCountCmd(1, new EntityCountQuery(ef, Collections.emptyList()));
getWsClient().send(cmd);
EntityCountUpdate update = getWsClient().parseCountReply(getWsClient().waitForReply());
Assert.assertEquals(1, update.getCmdId());
Assert.assertEquals(101, update.getCount());
for (Tenant tenant : tenants) {
doDelete("/api/tenant/" + tenant.getId().toString());
}
}
@Test
public void testTenantProfilesCountWsCmd() throws Exception {
loginSysAdmin();
List<TenantProfile> tenantProfiles = new ArrayList<>();
for (int i = 0; i < 100; i++) {
TenantProfile tenantProfile = new TenantProfile();
tenantProfile.setName("tenantProfile" + i);
tenantProfiles.add(doPost("/api/tenantProfile", tenantProfile, TenantProfile.class));
}
EntityTypeFilter ef = new EntityTypeFilter();
ef.setEntityType(EntityType.TENANT_PROFILE);
EntityCountCmd cmd = new EntityCountCmd(1, new EntityCountQuery(ef, Collections.emptyList()));
getWsClient().send(cmd);
EntityCountUpdate update = getWsClient().parseCountReply(getWsClient().waitForReply());
Assert.assertEquals(1, update.getCmdId());
Assert.assertEquals(101, update.getCount());
for (TenantProfile tenantProfile : tenantProfiles) {
doDelete("/api/tenantProfile/" + tenantProfile.getId().toString());
}
}
@Test
public void testUsersCountWsCmd() throws Exception {
loginSysAdmin();
List<User> users = new ArrayList<>();
for (int i = 0; i < 100; i++) {
User user = new User();
user.setEmail(i + "user@thingsboard.org");
user.setTenantId(tenantId);
user.setAuthority(Authority.TENANT_ADMIN);
users.add(doPost("/api/user", user, User.class));
}
EntityTypeFilter ef = new EntityTypeFilter();
ef.setEntityType(EntityType.USER);
EntityCountCmd cmd = new EntityCountCmd(1, new EntityCountQuery(ef, Collections.emptyList()));
getWsClient().send(cmd);
EntityCountUpdate update = getWsClient().parseCountReply(getWsClient().waitForReply());
Assert.assertEquals(1, update.getCmdId());
Assert.assertEquals(103, update.getCount());
for (User user : users) {
doDelete("/api/user/" + user.getId().toString());
}
}
@Test
public void testCustomersCountWsCmd() throws Exception {
loginTenantAdmin();
List<Customer> customers = new ArrayList<>();
for (int i = 0; i < 100; i++) {
Customer customer = new Customer();
customer.setTitle("customer" + i);
customers.add(doPost("/api/customer", customer, Customer.class));
}
loginSysAdmin();
EntityTypeFilter ef = new EntityTypeFilter();
ef.setEntityType(EntityType.CUSTOMER);
EntityCountCmd cmd = new EntityCountCmd(1, new EntityCountQuery(ef, Collections.emptyList()));
getWsClient().send(cmd);
EntityCountUpdate update = getWsClient().parseCountReply(getWsClient().waitForReply());
Assert.assertEquals(1, update.getCmdId());
Assert.assertEquals(101, update.getCount());
loginTenantAdmin();
for (Customer customer : customers) {
doDelete("/api/customer/" + customer.getId().toString());
}
}
@Test
public void testDevicesCountWsCmd() throws Exception {
loginTenantAdmin();
List<Device> devices = new ArrayList<>();
for (int i = 0; i < 100; i++) {
Device device = new Device();
device.setName("device" + i);
devices.add(doPost("/api/device", device, Device.class));
}
loginSysAdmin();
EntityTypeFilter ef = new EntityTypeFilter();
ef.setEntityType(EntityType.DEVICE);
EntityCountCmd cmd = new EntityCountCmd(1, new EntityCountQuery(ef, Collections.emptyList()));
getWsClient().send(cmd);
EntityCountUpdate update = getWsClient().parseCountReply(getWsClient().waitForReply());
Assert.assertEquals(1, update.getCmdId());
Assert.assertEquals(100, update.getCount());
loginTenantAdmin();
for (Device device : devices) {
doDelete("/api/device/" + device.getId().toString());
}
}
@Test
public void testAssetsCountWsCmd() throws Exception {
loginTenantAdmin();
List<Asset> assets = new ArrayList<>();
for (int i = 0; i < 100; i++) {
Asset asset = new Asset();
asset.setName("asset" + i);
assets.add(doPost("/api/asset", asset, Asset.class));
}
loginSysAdmin();
EntityTypeFilter ef = new EntityTypeFilter();
ef.setEntityType(EntityType.ASSET);
EntityCountCmd cmd = new EntityCountCmd(1, new EntityCountQuery(ef, Collections.emptyList()));
getWsClient().send(cmd);
EntityCountUpdate update = getWsClient().parseCountReply(getWsClient().waitForReply());
Assert.assertEquals(1, update.getCmdId());
Assert.assertEquals(100, update.getCount());
loginTenantAdmin();
for (Asset asset : assets) {
doDelete("/api/asset/" + asset.getId().toString());
}
}
@Test
public void testSystemInfoTimeSeriesWsCmd() throws Exception {
ApiUsageState apiUsageState = apiUsageStateClient.getApiUsageState(TenantId.SYS_TENANT_ID);
Assert.assertNotNull(apiUsageState);
loginSysAdmin();
long now = System.currentTimeMillis();
EntityDataUpdate update = getWsClient().sendEntityDataQuery(new ApiUsageStateFilter());
Assert.assertEquals(1, update.getCmdId());
PageData<EntityData> pageData = update.getData();
Assert.assertNotNull(pageData);
Assert.assertEquals(1, pageData.getData().size());
Assert.assertEquals(apiUsageState.getId(), pageData.getData().get(0).getEntityId());
update = getWsClient().subscribeTsUpdate(
List.of("memoryUsage", "totalMemory", "freeMemory", "cpuUsage", "totalCpuUsage", "freeDiscSpace", "totalDiscSpace"),
now, TimeUnit.HOURS.toMillis(1));
Assert.assertEquals(1, update.getCmdId());
List<EntityData> listData = update.getUpdate();
Assert.assertNotNull(listData);
Assert.assertEquals(1, listData.size());
Assert.assertEquals(apiUsageState.getId(), listData.get(0).getEntityId());
Assert.assertEquals(7, listData.get(0).getTimeseries().size());
for (TsValue[] tsv : listData.get(0).getTimeseries().values()) {
Assert.assertTrue(tsv.length > 1);
}
}
@Test
public void testGetFeaturesInfo() throws Exception {
loginSysAdmin();
FeaturesInfo featuresInfo = doGet("/api/admin/featuresInfo", FeaturesInfo.class);
Assert.assertNotNull(featuresInfo);
Assert.assertFalse(featuresInfo.isEmailEnabled());
Assert.assertFalse(featuresInfo.isSmsEnabled());
Assert.assertFalse(featuresInfo.isTwoFaEnabled());
Assert.assertFalse(featuresInfo.isNotificationEnabled());
Assert.assertFalse(featuresInfo.isOauthEnabled());
AdminSettings mailSettings = doGet("/api/admin/settings/mail", AdminSettings.class);
JsonNode jsonValue = mailSettings.getJsonValue();
((ObjectNode) jsonValue).put("mailFrom", "test@thingsboard.org");
mailSettings.setJsonValue(jsonValue);
doPost("/api/admin/settings", mailSettings).andExpect(status().isOk());
featuresInfo = doGet("/api/admin/featuresInfo", FeaturesInfo.class);
Assert.assertTrue(featuresInfo.isEmailEnabled());
Assert.assertFalse(featuresInfo.isSmsEnabled());
Assert.assertFalse(featuresInfo.isTwoFaEnabled());
Assert.assertFalse(featuresInfo.isNotificationEnabled());
Assert.assertFalse(featuresInfo.isOauthEnabled());
AdminSettings smsSettings = new AdminSettings();
smsSettings.setKey("sms");
smsSettings.setJsonValue(JacksonUtil.newObjectNode());
doPost("/api/admin/settings", smsSettings).andExpect(status().isOk());
featuresInfo = doGet("/api/admin/featuresInfo", FeaturesInfo.class);
Assert.assertTrue(featuresInfo.isEmailEnabled());
Assert.assertTrue(featuresInfo.isSmsEnabled());
Assert.assertFalse(featuresInfo.isTwoFaEnabled());
Assert.assertFalse(featuresInfo.isNotificationEnabled());
Assert.assertFalse(featuresInfo.isOauthEnabled());
AdminSettings twoFaSettingsSettings = new AdminSettings();
twoFaSettingsSettings.setKey("twoFaSettings");
twoFaSettingsSettings.setJsonValue(JacksonUtil.newObjectNode());
doPost("/api/admin/settings", twoFaSettingsSettings).andExpect(status().isOk());
featuresInfo = doGet("/api/admin/featuresInfo", FeaturesInfo.class);
Assert.assertTrue(featuresInfo.isEmailEnabled());
Assert.assertTrue(featuresInfo.isSmsEnabled());
Assert.assertTrue(featuresInfo.isTwoFaEnabled());
Assert.assertFalse(featuresInfo.isNotificationEnabled());
Assert.assertFalse(featuresInfo.isOauthEnabled());
AdminSettings notificationsSettings = new AdminSettings();
notificationsSettings.setKey("notifications");
notificationsSettings.setJsonValue(JacksonUtil.newObjectNode());
doPost("/api/admin/settings", notificationsSettings).andExpect(status().isOk());
featuresInfo = doGet("/api/admin/featuresInfo", FeaturesInfo.class);
Assert.assertTrue(featuresInfo.isEmailEnabled());
Assert.assertTrue(featuresInfo.isSmsEnabled());
Assert.assertTrue(featuresInfo.isTwoFaEnabled());
Assert.assertTrue(featuresInfo.isNotificationEnabled());
Assert.assertFalse(featuresInfo.isOauthEnabled());
OAuth2Info oAuth2Info = createDefaultOAuth2Info();
doPost("/api/oauth2/config", oAuth2Info).andExpect(status().isOk());
featuresInfo = doGet("/api/admin/featuresInfo", FeaturesInfo.class);
Assert.assertNotNull(featuresInfo);
Assert.assertTrue(featuresInfo.isEmailEnabled());
Assert.assertTrue(featuresInfo.isSmsEnabled());
Assert.assertTrue(featuresInfo.isTwoFaEnabled());
Assert.assertTrue(featuresInfo.isNotificationEnabled());
Assert.assertTrue(featuresInfo.isOauthEnabled());
}
private OAuth2Info createDefaultOAuth2Info() {
return new OAuth2Info(true, Lists.newArrayList(
OAuth2ParamsInfo.builder()
.domainInfos(Lists.newArrayList(
OAuth2DomainInfo.builder().name("domain").scheme(SchemeType.MIXED).build()
))
.mobileInfos(Collections.emptyList())
.clientRegistrations(Lists.newArrayList(
validRegistrationInfo()
))
.build()
));
}
private OAuth2RegistrationInfo validRegistrationInfo() {
return OAuth2RegistrationInfo.builder()
.clientId(UUID.randomUUID().toString())
.clientSecret(UUID.randomUUID().toString())
.authorizationUri(UUID.randomUUID().toString())
.accessTokenUri(UUID.randomUUID().toString())
.scope(Arrays.asList(UUID.randomUUID().toString(), UUID.randomUUID().toString()))
.platforms(Collections.emptyList())
.userInfoUri(UUID.randomUUID().toString())
.userNameAttributeName(UUID.randomUUID().toString())
.jwkSetUri(UUID.randomUUID().toString())
.clientAuthenticationMethod(UUID.randomUUID().toString())
.loginButtonLabel(UUID.randomUUID().toString())
.mapperConfig(
OAuth2MapperConfig.builder()
.type(MapperType.CUSTOM)
.custom(
OAuth2CustomMapperConfig.builder()
.url(UUID.randomUUID().toString())
.build()
)
.build()
)
.build();
}
}

23
application/src/test/java/org/thingsboard/server/controller/sql/HomePageApiSqlTest.java

@ -0,0 +1,23 @@
/**
* Copyright © 2016-2023 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.
*/
package org.thingsboard.server.controller.sql;
import org.thingsboard.server.controller.BaseHomePageApiTest;
import org.thingsboard.server.dao.service.DaoSqlTest;
@DaoSqlTest
public class HomePageApiSqlTest extends BaseHomePageApiTest {
}

11
common/cluster-api/src/main/proto/queue.proto

@ -27,6 +27,17 @@ message ServiceInfo {
string serviceId = 1;
repeated string serviceTypes = 2;
repeated string transports = 6;
SystemInfoProto systemInfo = 10;
}
message SystemInfoProto {
double cpuUsage = 1;
double totalCpuUsage = 2;
int64 memoryUsage = 3;
int64 totalMemory = 4;
int64 freeMemory = 5;
int64 freeDiscSpace = 6;
int64 totalDiscSpace = 7;
}
/**

27
common/data/src/main/java/org/thingsboard/server/common/data/FeaturesInfo.java

@ -0,0 +1,27 @@
/**
* Copyright © 2016-2023 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.
*/
package org.thingsboard.server.common.data;
import lombok.Data;
@Data
public class FeaturesInfo {
boolean isEmailEnabled;
boolean isSmsEnabled;
boolean isNotificationEnabled;
boolean isOauthEnabled;
boolean isTwoFaEnabled;
}

29
common/data/src/main/java/org/thingsboard/server/common/data/SystemInfo.java

@ -0,0 +1,29 @@
/**
* Copyright © 2016-2023 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.
*/
package org.thingsboard.server.common.data;
import io.swagger.annotations.ApiModelProperty;
import lombok.Data;
import java.util.List;
@Data
public class SystemInfo {
@ApiModelProperty(position = 1, value = "Is monolith.")
private boolean isMonolith;
@ApiModelProperty(position = 2, value = "System data.")
private List<SystemInfoData> systemData;
}

43
common/data/src/main/java/org/thingsboard/server/common/data/SystemInfoData.java

@ -0,0 +1,43 @@
/**
* Copyright © 2016-2023 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.
*/
package org.thingsboard.server.common.data;
import io.swagger.annotations.ApiModelProperty;
import lombok.Data;
import java.util.Map;
@Data
public class SystemInfoData {
@ApiModelProperty(position = 1, value = "Service Id.")
private String serviceId;
@ApiModelProperty(position = 2, value = "Service type.")
private String serviceType;
@ApiModelProperty(position = 3, value = "CPU usage.")
private Double cpuUsage;
@ApiModelProperty(position = 4, value = "Total CPU usage.")
private Double totalCpuUsage;
@ApiModelProperty(position = 5, value = "Memory usage in bytes.")
private Long memoryUsage;
@ApiModelProperty(position = 6, value = "Total memory in bytes.")
private Long totalMemory;
@ApiModelProperty(position = 6, value = "Free memory in bytes.")
private Long freeMemory;
@ApiModelProperty(position = 7, value = "Free disc space in bytes.")
private Long freeDiscSpace;
@ApiModelProperty(position = 7, value = "Total disc space in bytes.")
private Long totalDiscSpace;
}

5
common/data/src/main/java/org/thingsboard/server/common/data/UpdateMessage.java

@ -25,7 +25,8 @@ public class UpdateMessage {
@ApiModelProperty(position = 1, value = "The message about new platform update available.")
private final String message;
@ApiModelProperty(position = 1, value = "'True' if new platform update is available.")
@ApiModelProperty(position = 2, value = "'True' if new platform update is available.")
private final boolean isUpdateAvailable;
@ApiModelProperty(position = 3, value = "Current ThingsBoard version.")
private final String currentVersion;
}

38
common/queue/src/main/java/org/thingsboard/server/queue/discovery/DefaultTbServiceInfoProvider.java

@ -24,6 +24,7 @@ import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.TbTransportService;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo;
import org.thingsboard.server.queue.util.AfterContextReady;
@ -36,6 +37,14 @@ import java.util.Collections;
import java.util.List;
import java.util.stream.Collectors;
import static org.thingsboard.common.util.SystemUtil.getCpuUsage;
import static org.thingsboard.common.util.SystemUtil.getFreeDiscSpace;
import static org.thingsboard.common.util.SystemUtil.getFreeMemory;
import static org.thingsboard.common.util.SystemUtil.getMemoryUsage;
import static org.thingsboard.common.util.SystemUtil.getTotalCpuUsage;
import static org.thingsboard.common.util.SystemUtil.getTotalDiscSpace;
import static org.thingsboard.common.util.SystemUtil.getTotalMemory;
@Component
@Slf4j
public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider {
@ -69,11 +78,8 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider {
} else {
serviceTypes = Collections.singletonList(ServiceType.of(serviceType));
}
ServiceInfo.Builder builder = ServiceInfo.newBuilder()
.setServiceId(serviceId)
.addAllServiceTypes(serviceTypes.stream().map(ServiceType::name).collect(Collectors.toList()));
serviceInfo = builder.build();
serviceInfo = getServiceInfoWithCurrentSystemInfo();
}
@AfterContextReady
@ -99,4 +105,28 @@ public class DefaultTbServiceInfoProvider implements TbServiceInfoProvider {
return serviceTypes.contains(serviceType);
}
@Override
public ServiceInfo getServiceInfoWithCurrentSystemInfo() {
ServiceInfo.Builder builder = ServiceInfo.newBuilder()
.setServiceId(serviceId)
.addAllServiceTypes(serviceTypes.stream().map(ServiceType::name).collect(Collectors.toList()))
.setSystemInfo(getCurrentSystemInfoProto());
return builder.build();
}
private TransportProtos.SystemInfoProto getCurrentSystemInfoProto() {
TransportProtos.SystemInfoProto.Builder builder = TransportProtos.SystemInfoProto.newBuilder();
getMemoryUsage().ifPresent(builder::setMemoryUsage);
getTotalMemory().ifPresent(builder::setTotalMemory);
getFreeMemory().ifPresent(builder::setFreeMemory);
getCpuUsage().ifPresent(builder::setCpuUsage);
getTotalCpuUsage().ifPresent(builder::setTotalCpuUsage);
getFreeDiscSpace().ifPresent(builder::setFreeDiscSpace);
getTotalDiscSpace().ifPresent(builder::setTotalDiscSpace);
return builder.build();
}
}

8
common/queue/src/main/java/org/thingsboard/server/queue/discovery/DiscoveryService.java

@ -15,6 +15,14 @@
*/
package org.thingsboard.server.queue.discovery;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.List;
public interface DiscoveryService {
List<TransportProtos.ServiceInfo> getOtherServers();
boolean isMonolith();
}

11
common/queue/src/main/java/org/thingsboard/server/queue/discovery/DummyDiscoveryService.java

@ -22,9 +22,11 @@ import org.springframework.context.annotation.DependsOn;
import org.springframework.context.event.EventListener;
import org.springframework.core.annotation.Order;
import org.springframework.stereotype.Service;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.util.AfterStartUp;
import java.util.Collections;
import java.util.List;
@Service
@ConditionalOnProperty(prefix = "zk", value = "enabled", havingValue = "false", matchIfMissing = true)
@ -46,4 +48,13 @@ public class DummyDiscoveryService implements DiscoveryService {
partitionService.recalculatePartitions(serviceInfoProvider.getServiceInfo(), Collections.emptyList());
}
@Override
public List<TransportProtos.ServiceInfo> getOtherServers() {
return Collections.emptyList();
}
@Override
public boolean isMonolith() {
return true;
}
}

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

@ -22,7 +22,6 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.QueueId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
@ -217,7 +216,7 @@ public class HashPartitionService implements PartitionService {
}
queueServicesMap.values().forEach(list -> list.sort(Comparator.comparing(ServiceInfo::getServiceId)));
final ConcurrentMap<QueueKey, List<Integer>> newPartitions = new ConcurrentHashMap<>();
final ConcurrentMap<QueueKey, List<Integer>> newPartitions = new ConcurrentHashMap<>();
partitionSizesMap.forEach((queueKey, size) -> {
for (int i = 0; i < size; i++) {
ServiceInfo serviceInfo = resolveByPartitionIdx(queueServicesMap.get(queueKey), queueKey, i);

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

@ -16,7 +16,6 @@
package org.thingsboard.server.queue.discovery;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.QueueId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;

3
common/queue/src/main/java/org/thingsboard/server/queue/discovery/TbServiceInfoProvider.java

@ -16,6 +16,7 @@
package org.thingsboard.server.queue.discovery;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ServiceInfo;
public interface TbServiceInfoProvider {
@ -28,4 +29,6 @@ public interface TbServiceInfoProvider {
boolean isService(ServiceType serviceType);
ServiceInfo getServiceInfoWithCurrentSystemInfo();
}

32
common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java

@ -16,6 +16,7 @@
package org.thingsboard.server.queue.discovery;
import com.google.protobuf.InvalidProtocolBufferException;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.CuratorFrameworkFactory;
@ -33,21 +34,19 @@ import org.apache.zookeeper.KeeperException;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.event.EventListener;
import org.springframework.core.annotation.Order;
import org.springframework.stereotype.Service;
import org.springframework.util.Assert;
import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.queue.discovery.event.ServiceListChangedEvent;
import org.thingsboard.server.queue.util.AfterStartUp;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.util.List;
import java.util.NoSuchElementException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import static org.apache.curator.framework.recipes.cache.PathChildrenCacheEvent.Type.CHILD_REMOVED;
@ -71,7 +70,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
private final TbServiceInfoProvider serviceInfoProvider;
private final PartitionService partitionService;
private ExecutorService reconnectExecutorService;
private ScheduledExecutorService zkExecutorService;
private CuratorFramework client;
private PathChildrenCache cache;
private String nodePath;
@ -93,7 +92,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
Assert.notNull(zkConnectionTimeout, missingProperty("zk.connection_timeout_ms"));
Assert.notNull(zkSessionTimeout, missingProperty("zk.session_timeout_ms"));
reconnectExecutorService = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("zk-discovery"));
zkExecutorService = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("zk-discovery"));
log.info("Initializing discovery service using ZK connect string: {}", zkUrl);
@ -101,7 +100,8 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
initZkClient();
}
private List<TransportProtos.ServiceInfo> getOtherServers() {
@Override
public List<TransportProtos.ServiceInfo> getOtherServers() {
return cache.getCurrentData().stream()
.filter(cd -> !cd.getPath().equals(nodePath))
.map(cd -> {
@ -115,6 +115,11 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
.collect(Collectors.toList());
}
@Override
public boolean isMonolith() {
return false;
}
@AfterStartUp(order = AfterStartUp.DISCOVERY_SERVICE)
public void onApplicationEvent(ApplicationReadyEvent event) {
if (stopped) {
@ -128,15 +133,17 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
return;
}
log.info("Going to publish current server...");
publishCurrentServer();
zkExecutorService.scheduleAtFixedRate(this::publishCurrentServer, 0, 1, TimeUnit.MINUTES);
log.info("Going to recalculate partitions...");
recalculatePartitions();
}
@SneakyThrows
public synchronized void publishCurrentServer() {
TransportProtos.ServiceInfo self = serviceInfoProvider.getServiceInfo();
if (currentServerExists()) {
log.info("[{}] ZK node for current instance already exists, NOT created new one: {}", self.getServiceId(), nodePath);
log.trace("[{}] Updating ZK node for current instance: {}", self.getServiceId(), nodePath);
client.setData().forPath(nodePath, serviceInfoProvider.getServiceInfoWithCurrentSystemInfo().toByteArray());
} else {
try {
log.info("[{}] Creating ZK node for current instance", self.getServiceId());
@ -175,7 +182,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
return (client, newState) -> {
log.info("[{}] ZK state changed: {}", self.getServiceId(), newState);
if (newState == ConnectionState.LOST) {
reconnectExecutorService.submit(this::reconnect);
zkExecutorService.submit(this::reconnect);
}
};
}
@ -240,7 +247,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
@PreDestroy
public void destroy() {
destroyZkClient();
reconnectExecutorService.shutdownNow();
zkExecutorService.shutdownNow();
log.info("Stopped discovery service");
}
@ -280,10 +287,9 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
log.error("Failed to decode server instance for node {}", data.getPath(), e);
throw e;
}
log.info("Processing [{}] event for [{}]", pathChildrenCacheEvent.getType(), instance.getServiceId());
log.debug("Processing [{}] event for [{}]", pathChildrenCacheEvent.getType(), instance.getServiceId());
switch (pathChildrenCacheEvent.getType()) {
case CHILD_ADDED:
case CHILD_UPDATED:
case CHILD_REMOVED:
recalculatePartitions();
break;

4
common/util/pom.xml

@ -88,6 +88,10 @@
<groupId>org.thingsboard.common</groupId>
<artifactId>data</artifactId>
</dependency>
<dependency>
<groupId>com.github.dblock</groupId>
<artifactId>oshi-core</artifactId>
</dependency>
</dependencies>
<build>

108
common/util/src/main/java/org/thingsboard/common/util/SystemUtil.java

@ -0,0 +1,108 @@
/**
* Copyright © 2016-2023 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.
*/
package org.thingsboard.common.util;
import lombok.extern.slf4j.Slf4j;
import oshi.SystemInfo;
import oshi.hardware.HardwareAbstractionLayer;
import java.lang.management.ManagementFactory;
import java.lang.management.MemoryMXBean;
import java.nio.file.FileStore;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.Optional;
@Slf4j
public class SystemUtil {
private static final HardwareAbstractionLayer HARDWARE;
static {
SystemInfo si = new SystemInfo();
HARDWARE = si.getHardware();
}
public static Optional<Long> getMemoryUsage() {
try {
MemoryMXBean memoryMXBean = ManagementFactory.getMemoryMXBean();
return Optional.of(memoryMXBean.getHeapMemoryUsage().getUsed());
} catch (Exception e) {
log.debug("Failed to get memory usage!!!", e);
}
return Optional.empty();
}
public static Optional<Long> getTotalMemory() {
try {
return Optional.of(HARDWARE.getMemory().getTotal());
} catch (Exception e) {
log.debug("Failed to get total memory!!!", e);
}
return Optional.empty();
}
public static Optional<Long> getFreeMemory() {
try {
return Optional.of(HARDWARE.getMemory().getAvailable());
} catch (Exception e) {
log.debug("Failed to get free memory!!!", e);
}
return Optional.empty();
}
public static Optional<Double> getCpuUsage() {
try {
return Optional.of(prepare(HARDWARE.getProcessor().getSystemLoadAverage()));
} catch (Exception e) {
log.debug("Failed to get cpu usage!!!", e);
}
return Optional.empty();
}
public static Optional<Double> getTotalCpuUsage() {
try {
return Optional.of(prepare(HARDWARE.getProcessor().getSystemCpuLoad() * 100));
} catch (Exception e) {
log.debug("Failed to get total cpu usage!!!", e);
}
return Optional.empty();
}
public static Optional<Long> getFreeDiscSpace() {
try {
FileStore store = Files.getFileStore(Paths.get("/"));
return Optional.of(store.getUsableSpace());
} catch (Exception e) {
log.debug("Failed to get free disc space!!!", e);
}
return Optional.empty();
}
public static Optional<Long> getTotalDiscSpace() {
try {
FileStore store = Files.getFileStore(Paths.get("/"));
return Optional.of(store.getTotalSpace());
} catch (Exception e) {
log.debug("Failed to get total disc space!!!", e);
}
return Optional.empty();
}
private static Double prepare(Double d) {
return (int) (d * 100) / 100.0;
}
}

2
dao/src/main/java/org/thingsboard/server/dao/entity/BaseEntityService.java

@ -25,10 +25,10 @@ import org.thingsboard.server.common.data.HasLabel;
import org.thingsboard.server.common.data.HasName;
import org.thingsboard.server.common.data.HasTitle;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.id.NameLabelAndCustomerDetails;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.HasId;
import org.thingsboard.server.common.data.id.NameLabelAndCustomerDetails;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.query.EntityCountQuery;

2
dao/src/main/java/org/thingsboard/server/dao/service/DataValidator.java

@ -92,7 +92,7 @@ public abstract class DataValidator<D extends BaseData<?>> {
}
}
private static boolean doValidateEmail(String email) {
public static boolean doValidateEmail(String email) {
if (email == null) {
return false;
}

6
dao/src/main/java/org/thingsboard/server/dao/sql/query/DefaultEntityQueryRepository.java

@ -241,6 +241,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository {
entityTableMap.put(EntityType.RULE_CHAIN, "rule_chain");
entityTableMap.put(EntityType.DEVICE_PROFILE, "device_profile");
entityTableMap.put(EntityType.ASSET_PROFILE, "asset_profile");
entityTableMap.put(EntityType.TENANT_PROFILE, "tenant_profile");
}
public static EntityType[] RELATION_QUERY_ENTITY_TYPES = new EntityType[]{
@ -311,12 +312,13 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository {
@Override
public long countEntitiesByQuery(TenantId tenantId, CustomerId customerId, EntityCountQuery query) {
EntityType entityType = resolveEntityType(query.getEntityFilter());
QueryContext ctx = new QueryContext(new QuerySecurityContext(tenantId, customerId, entityType));
QueryContext ctx = new QueryContext(new QuerySecurityContext(tenantId, customerId, entityType, TenantId.SYS_TENANT_ID.equals(tenantId)));
if (query.getKeyFilters() == null || query.getKeyFilters().isEmpty()) {
ctx.append("select count(e.id) from ");
ctx.append(addEntityTableQuery(ctx, query.getEntityFilter()));
ctx.append(" e where ");
ctx.append(buildEntityWhere(ctx, query.getEntityFilter(), Collections.emptyList()));
return transactionTemplate.execute(status -> {
long startTs = System.currentTimeMillis();
try {
@ -514,7 +516,7 @@ public class DefaultEntityQueryRepository implements EntityQueryRepository {
}
private String buildPermissionQuery(QueryContext ctx, EntityFilter entityFilter) {
if(ctx.isIgnorePermissionCheck()){
if (ctx.isIgnorePermissionCheck()) {
return "1=1";
}
switch (entityFilter.getType()) {

6
pom.xml

@ -150,6 +150,7 @@
<allure-testng.version>2.21.0</allure-testng.version>
<allure-maven.version>2.12.0</allure-maven.version>
<slack-api.version>1.12.1</slack-api.version>
<oshi.version>3.4.0</oshi.version>
</properties>
<modules>
@ -1994,6 +1995,11 @@
<artifactId>exp4j</artifactId>
<version>${exp4j.version}</version>
</dependency>
<dependency>
<groupId>com.github.dblock</groupId>
<artifactId>oshi-core</artifactId>
<version>${oshi.version}</version>
</dependency>
</dependencies>
</dependencyManagement>

7
rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java

@ -32,7 +32,6 @@ import org.springframework.http.ResponseEntity;
import org.springframework.http.client.support.HttpRequestWrapper;
import org.springframework.util.LinkedMultiValueMap;
import org.springframework.util.MultiValueMap;
import org.thingsboard.server.common.data.StringUtils;
import org.springframework.web.client.HttpClientErrorException;
import org.springframework.web.client.RestTemplate;
import org.thingsboard.common.util.ThingsBoardExecutors;
@ -56,6 +55,8 @@ import org.thingsboard.server.common.data.EventInfo;
import org.thingsboard.server.common.data.OtaPackage;
import org.thingsboard.server.common.data.OtaPackageInfo;
import org.thingsboard.server.common.data.SaveDeviceWithCredentialsRequest;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.SystemInfo;
import org.thingsboard.server.common.data.TbResource;
import org.thingsboard.server.common.data.TbResourceInfo;
import org.thingsboard.server.common.data.Tenant;
@ -394,6 +395,10 @@ public class RestClient implements Closeable {
}
}
public SystemInfo getSystemInfo() {
return restTemplate.getForEntity(baseURL + "/api/admin/systemInfo", SystemInfo.class).getBody();
}
public Optional<Alarm> getAlarmById(AlarmId alarmId) {
try {
ResponseEntity<Alarm> alarm = restTemplate.getForEntity(baseURL + "/api/alarm/{alarmId}", Alarm.class, alarmId.getId());

Loading…
Cancel
Save