Browse Source

Merge branch 'master' into feature/image-resources

pull/9542/head
ViacheslavKlimov 3 years ago
parent
commit
5337943747
  1. 3
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java
  2. 21
      application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java
  3. 2
      application/src/main/java/org/thingsboard/server/config/CustomOAuth2AuthorizationRequestResolver.java
  4. 44
      application/src/main/java/org/thingsboard/server/config/TbRuleEngineSecurityConfiguration.java
  5. 2
      application/src/main/java/org/thingsboard/server/controller/DeviceConnectivityController.java
  6. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java
  7. 3
      application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java
  8. 15
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  9. 22
      application/src/main/java/org/thingsboard/server/service/queue/ProtoUtils.java
  10. 2
      application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/HttpCookieOAuth2AuthorizationRequestRepository.java
  11. 229
      application/src/test/java/org/thingsboard/server/controller/DeviceConnectivityControllerTest.java
  12. 29
      application/src/test/java/org/thingsboard/server/controller/DeviceControllerTest.java
  13. 30
      application/src/test/java/org/thingsboard/server/edge/RelationEdgeTest.java
  14. 4
      common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java
  15. 4
      common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java
  16. 8
      common/cluster-api/src/main/proto/queue.proto
  17. 4
      common/data/src/main/java/org/thingsboard/server/common/data/FstStatsService.java
  18. 2
      common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java
  19. 36
      common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceDeleteMsg.java
  20. 8
      common/queue/src/main/java/org/thingsboard/server/queue/util/ProtoWithFSTService.java
  21. 28
      common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbUtils.java
  22. 90
      common/script/script-api/src/test/java/org/thingsboard/script/api/tbel/TbUtilsTest.java
  23. 20
      common/stats/src/main/java/org/thingsboard/server/common/stats/FstStatsServiceImpl.java
  24. 20
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceConnectivityServiceImpl.java
  25. 73
      dao/src/main/java/org/thingsboard/server/dao/util/DeviceConnectivityUtil.java

3
application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java

@ -64,6 +64,9 @@ public class DeviceActor extends ContextAwareActor {
case DEVICE_ATTRIBUTES_UPDATE_TO_DEVICE_ACTOR_MSG: case DEVICE_ATTRIBUTES_UPDATE_TO_DEVICE_ACTOR_MSG:
processor.processAttributesUpdate((DeviceAttributesEventNotificationMsg) msg); processor.processAttributesUpdate((DeviceAttributesEventNotificationMsg) msg);
break; break;
case DEVICE_DELETE_TO_DEVICE_ACTOR_MSG:
ctx.stop(ctx.getSelf());
break;
case DEVICE_CREDENTIALS_UPDATE_TO_DEVICE_ACTOR_MSG: case DEVICE_CREDENTIALS_UPDATE_TO_DEVICE_ACTOR_MSG:
processor.processCredentialsUpdate(msg); processor.processCredentialsUpdate(msg);
break; break;

21
application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java

@ -53,10 +53,13 @@ import org.thingsboard.server.common.msg.queue.PartitionChangeMsg;
import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg;
import org.thingsboard.server.common.msg.queue.RuleEngineException; import org.thingsboard.server.common.msg.queue.RuleEngineException;
import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.rule.engine.DeviceDeleteMsg;
import org.thingsboard.server.service.edge.rpc.EdgeRpcService; import org.thingsboard.server.service.edge.rpc.EdgeRpcService;
import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper; import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper;
import java.util.HashSet;
import java.util.List; import java.util.List;
import java.util.Set;
@Slf4j @Slf4j
public class TenantActor extends RuleChainManagerActor { public class TenantActor extends RuleChainManagerActor {
@ -65,8 +68,11 @@ public class TenantActor extends RuleChainManagerActor {
private boolean isCore; private boolean isCore;
private ApiUsageState apiUsageState; private ApiUsageState apiUsageState;
private Set<DeviceId> deletedDevices;
private TenantActor(ActorSystemContext systemContext, TenantId tenantId) { private TenantActor(ActorSystemContext systemContext, TenantId tenantId) {
super(systemContext, tenantId); super(systemContext, tenantId);
this.deletedDevices = new HashSet<>();
} }
boolean cantFindTenant = false; boolean cantFindTenant = false;
@ -221,6 +227,10 @@ public class TenantActor extends RuleChainManagerActor {
if (!isCore) { if (!isCore) {
log.warn("RECEIVED INVALID MESSAGE: {}", msg); log.warn("RECEIVED INVALID MESSAGE: {}", msg);
} }
if (deletedDevices.contains(msg.getDeviceId())) {
log.debug("RECEIVED MESSAGE FOR DELETED DEVICE: {}", msg);
return;
}
TbActorRef deviceActor = getOrCreateDeviceActor(msg.getDeviceId()); TbActorRef deviceActor = getOrCreateDeviceActor(msg.getDeviceId());
if (priority) { if (priority) {
deviceActor.tellWithHighPriority(msg); deviceActor.tellWithHighPriority(msg);
@ -240,7 +250,8 @@ public class TenantActor extends RuleChainManagerActor {
log.info("[{}] Received API state update. Going to ENABLE Rule Engine execution.", tenantId); log.info("[{}] Received API state update. Going to ENABLE Rule Engine execution.", tenantId);
initRuleChains(); initRuleChains();
} }
} else if (msg.getEntityId().getEntityType() == EntityType.EDGE) { }
if (msg.getEntityId().getEntityType() == EntityType.EDGE) {
EdgeId edgeId = new EdgeId(msg.getEntityId().getId()); EdgeId edgeId = new EdgeId(msg.getEntityId().getId());
EdgeRpcService edgeRpcService = systemContext.getEdgeRpcService(); EdgeRpcService edgeRpcService = systemContext.getEdgeRpcService();
if (msg.getEvent() == ComponentLifecycleEvent.DELETED) { if (msg.getEvent() == ComponentLifecycleEvent.DELETED) {
@ -249,7 +260,13 @@ public class TenantActor extends RuleChainManagerActor {
Edge edge = systemContext.getEdgeService().findEdgeById(tenantId, edgeId); Edge edge = systemContext.getEdgeService().findEdgeById(tenantId, edgeId);
edgeRpcService.updateEdge(tenantId, edge); edgeRpcService.updateEdge(tenantId, edge);
} }
} else if (isRuleEngine) { }
if (msg.getEntityId().getEntityType() == EntityType.DEVICE && ComponentLifecycleEvent.DELETED == msg.getEvent()) {
DeviceId deviceId = (DeviceId) msg.getEntityId();
onToDeviceActorMsg(new DeviceDeleteMsg(tenantId, deviceId), true);
deletedDevices.add(deviceId);
}
if (isRuleEngine) {
TbActorRef target = getEntityActorRef(msg.getEntityId()); TbActorRef target = getEntityActorRef(msg.getEntityId());
if (target != null) { if (target != null) {
if (msg.getEntityId().getEntityType() == EntityType.RULE_CHAIN) { if (msg.getEntityId().getEntityType() == EntityType.RULE_CHAIN) {

2
application/src/main/java/org/thingsboard/server/config/CustomOAuth2AuthorizationRequestResolver.java

@ -38,6 +38,7 @@ import org.springframework.web.util.UriComponentsBuilder;
import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.dao.oauth2.OAuth2Configuration; import org.thingsboard.server.dao.oauth2.OAuth2Configuration;
import org.thingsboard.server.dao.oauth2.OAuth2Service; import org.thingsboard.server.dao.oauth2.OAuth2Service;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.security.auth.oauth2.TbOAuth2ParameterNames; import org.thingsboard.server.service.security.auth.oauth2.TbOAuth2ParameterNames;
import org.thingsboard.server.service.security.model.token.OAuth2AppTokenFactory; import org.thingsboard.server.service.security.model.token.OAuth2AppTokenFactory;
import org.thingsboard.server.utils.MiscUtils; import org.thingsboard.server.utils.MiscUtils;
@ -51,6 +52,7 @@ import java.util.HashMap;
import java.util.Map; import java.util.Map;
import java.util.UUID; import java.util.UUID;
@TbCoreComponent
@Service @Service
@Slf4j @Slf4j
public class CustomOAuth2AuthorizationRequestResolver implements OAuth2AuthorizationRequestResolver { public class CustomOAuth2AuthorizationRequestResolver implements OAuth2AuthorizationRequestResolver {

44
application/src/main/java/org/thingsboard/server/config/TbRuleEngineSecurityConfiguration.java

@ -0,0 +1,44 @@
/**
* 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.config;
import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression;
import org.springframework.boot.autoconfigure.security.SecurityProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.core.annotation.Order;
import org.springframework.security.config.annotation.method.configuration.EnableGlobalMethodSecurity;
import org.springframework.security.config.annotation.web.builders.HttpSecurity;
import org.springframework.security.config.annotation.web.configuration.EnableWebSecurity;
import org.springframework.security.web.SecurityFilterChain;
@Configuration
@EnableWebSecurity
@EnableGlobalMethodSecurity(prePostEnabled = true)
@Order(SecurityProperties.BASIC_AUTH_ORDER)
@ConditionalOnExpression("'${service.type:null}'=='tb-rule-engine'")
public class TbRuleEngineSecurityConfiguration {
@Bean
SecurityFilterChain filterChain(HttpSecurity http) throws Exception {
http.headers().cacheControl().and().frameOptions().disable()
.and().cors().and().csrf().disable()
.authorizeRequests()
.antMatchers("/actuator/prometheus").permitAll()
.anyRequest().authenticated();
return http.build();
}
}

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

@ -106,7 +106,7 @@ public class DeviceConnectivityController extends BaseController {
@RequestMapping(value = "/device-connectivity/gateway-launch/{deviceId}", method = RequestMethod.GET) @RequestMapping(value = "/device-connectivity/gateway-launch/{deviceId}", method = RequestMethod.GET)
@ResponseBody @ResponseBody
public JsonNode getGatewayLaunchCommands(@ApiParam(value = DEVICE_ID_PARAM_DESCRIPTION) public JsonNode getGatewayLaunchCommands(@ApiParam(value = DEVICE_ID_PARAM_DESCRIPTION)
@PathVariable(DEVICE_ID) String strDeviceId, HttpServletRequest request) throws ThingsboardException, URISyntaxException { @PathVariable(DEVICE_ID) String strDeviceId, HttpServletRequest request) throws ThingsboardException, URISyntaxException {
checkParameter(DEVICE_ID, strDeviceId); checkParameter(DEVICE_ID, strDeviceId);
DeviceId deviceId = new DeviceId(toUUID(strDeviceId)); DeviceId deviceId = new DeviceId(toUUID(strDeviceId));
Device device = checkDeviceId(deviceId, Operation.READ_CREDENTIALS); Device device = checkDeviceId(deviceId, Operation.READ_CREDENTIALS);

2
application/src/main/java/org/thingsboard/server/service/edge/rpc/sync/DefaultEdgeRequestsService.java

@ -294,7 +294,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService {
private ListenableFuture<List<EntityRelation>> findRelationByQuery(TenantId tenantId, Edge edge, private ListenableFuture<List<EntityRelation>> findRelationByQuery(TenantId tenantId, Edge edge,
EntityId entityId, EntitySearchDirection direction) { EntityId entityId, EntitySearchDirection direction) {
EntityRelationsQuery query = new EntityRelationsQuery(); EntityRelationsQuery query = new EntityRelationsQuery();
query.setParameters(new RelationsSearchParameters(entityId, direction, -1, false)); query.setParameters(new RelationsSearchParameters(entityId, direction, 1, false));
return relationService.findByQuery(tenantId, query); return relationService.findByQuery(tenantId, query);
} }

3
application/src/main/java/org/thingsboard/server/service/entitiy/DefaultTbNotificationEntityService.java

@ -112,7 +112,7 @@ public class DefaultTbNotificationEntityService implements TbNotificationEntityS
public void notifyDeleteDevice(TenantId tenantId, DeviceId deviceId, CustomerId customerId, Device device, public void notifyDeleteDevice(TenantId tenantId, DeviceId deviceId, CustomerId customerId, Device device,
User user, Object... additionalInfo) { User user, Object... additionalInfo) {
gatewayNotificationsService.onDeviceDeleted(device); gatewayNotificationsService.onDeviceDeleted(device);
tbClusterService.onDeviceDeleted(device, null); tbClusterService.onDeviceDeleted(tenantId, device, null);
logEntityAction(tenantId, deviceId, device, customerId, ActionType.DELETED, user, additionalInfo); logEntityAction(tenantId, deviceId, device, customerId, ActionType.DELETED, user, additionalInfo);
} }
@ -126,6 +126,7 @@ public class DefaultTbNotificationEntityService implements TbNotificationEntityS
@Override @Override
public void notifyAssignDeviceToTenant(TenantId tenantId, TenantId newTenantId, DeviceId deviceId, CustomerId customerId, public void notifyAssignDeviceToTenant(TenantId tenantId, TenantId newTenantId, DeviceId deviceId, CustomerId customerId,
Device device, Tenant tenant, User user, Object... additionalInfo) { Device device, Tenant tenant, User user, Object... additionalInfo) {
tbClusterService.onDeviceAssignedToTenant(tenantId, device);
logEntityAction(tenantId, deviceId, device, customerId, ActionType.ASSIGNED_TO_TENANT, user, additionalInfo); logEntityAction(tenantId, deviceId, device, customerId, ActionType.ASSIGNED_TO_TENANT, user, additionalInfo);
pushAssignedFromNotification(tenant, newTenantId, device); pushAssignedFromNotification(tenant, newTenantId, device);
} }

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

@ -308,10 +308,17 @@ public class DefaultTbClusterService implements TbClusterService {
} }
@Override @Override
public void onDeviceDeleted(Device device, TbQueueCallback callback) { public void onDeviceDeleted(TenantId tenantId, Device device, TbQueueCallback callback) {
broadcastEntityDeleteToTransport(device.getTenantId(), device.getId(), device.getName(), callback); DeviceId deviceId = device.getId();
sendDeviceStateServiceEvent(device.getTenantId(), device.getId(), false, false, true); broadcastEntityDeleteToTransport(tenantId, deviceId, device.getName(), callback);
broadcastEntityStateChangeEvent(device.getTenantId(), device.getId(), ComponentLifecycleEvent.DELETED); sendDeviceStateServiceEvent(tenantId, deviceId, false, false, true);
broadcastEntityStateChangeEvent(tenantId, deviceId, ComponentLifecycleEvent.DELETED);
}
@Override
public void onDeviceAssignedToTenant(TenantId oldTenantId, Device device) {
onDeviceDeleted(oldTenantId, device, null);
sendDeviceStateServiceEvent(device.getTenantId(), device.getId(), true, false, false);
} }
@Override @Override

22
application/src/main/java/org/thingsboard/server/service/queue/ProtoUtils.java

@ -46,6 +46,7 @@ import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest;
import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequestActorMsg; import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequestActorMsg;
import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceAttributesEventNotificationMsg;
import org.thingsboard.server.common.msg.rule.engine.DeviceCredentialsUpdateNotificationMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceCredentialsUpdateNotificationMsg;
import org.thingsboard.server.common.msg.rule.engine.DeviceDeleteMsg;
import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceEdgeUpdateMsg;
import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.common.msg.rule.engine.DeviceNameOrTypeUpdateMsg;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
@ -384,6 +385,21 @@ public class ProtoUtils {
); );
} }
private static TransportProtos.DeviceDeleteMsgProto toProto(DeviceDeleteMsg msg) {
return TransportProtos.DeviceDeleteMsgProto.newBuilder()
.setTenantIdMSB(msg.getTenantId().getId().getMostSignificantBits())
.setTenantIdLSB(msg.getTenantId().getId().getLeastSignificantBits())
.setDeviceIdMSB(msg.getDeviceId().getId().getMostSignificantBits())
.setDeviceIdLSB(msg.getDeviceId().getId().getLeastSignificantBits())
.build();
}
private static DeviceDeleteMsg fromProto(TransportProtos.DeviceDeleteMsgProto proto) {
return new DeviceDeleteMsg(
TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())),
new DeviceId(new UUID(proto.getDeviceIdMSB(), proto.getDeviceIdLSB())));
}
public static TransportProtos.ToDeviceActorNotificationMsgProto toProto(ToDeviceActorNotificationMsg msg) { public static TransportProtos.ToDeviceActorNotificationMsgProto toProto(ToDeviceActorNotificationMsg msg) {
if (msg instanceof DeviceEdgeUpdateMsg) { if (msg instanceof DeviceEdgeUpdateMsg) {
DeviceEdgeUpdateMsg updateMsg = (DeviceEdgeUpdateMsg) msg; DeviceEdgeUpdateMsg updateMsg = (DeviceEdgeUpdateMsg) msg;
@ -413,6 +429,10 @@ public class ProtoUtils {
RemoveRpcActorMsg updateMsg = (RemoveRpcActorMsg) msg; RemoveRpcActorMsg updateMsg = (RemoveRpcActorMsg) msg;
TransportProtos.RemoveRpcActorMsgProto proto = toProto(updateMsg); TransportProtos.RemoveRpcActorMsgProto proto = toProto(updateMsg);
return TransportProtos.ToDeviceActorNotificationMsgProto.newBuilder().setRemoveRpcActorMsg(proto).build(); return TransportProtos.ToDeviceActorNotificationMsgProto.newBuilder().setRemoveRpcActorMsg(proto).build();
} else if (msg instanceof DeviceDeleteMsg) {
DeviceDeleteMsg updateMsg = (DeviceDeleteMsg) msg;
TransportProtos.DeviceDeleteMsgProto proto = toProto(updateMsg);
return TransportProtos.ToDeviceActorNotificationMsgProto.newBuilder().setDeviceDeleteMsg(proto).build();
} }
return null; return null;
} }
@ -432,6 +452,8 @@ public class ProtoUtils {
return fromProto(proto.getFromDeviceRpcResponseMsg()); return fromProto(proto.getFromDeviceRpcResponseMsg());
} else if (proto.hasRemoveRpcActorMsg()) { } else if (proto.hasRemoveRpcActorMsg()) {
return fromProto(proto.getRemoveRpcActorMsg()); return fromProto(proto.getRemoveRpcActorMsg());
} else if (proto.hasDeviceDeleteMsg()) {
return fromProto(proto.getDeviceDeleteMsg());
} }
return null; return null;
} }

2
application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/HttpCookieOAuth2AuthorizationRequestRepository.java

@ -18,11 +18,13 @@ package org.thingsboard.server.service.security.auth.oauth2;
import org.springframework.security.oauth2.client.web.AuthorizationRequestRepository; import org.springframework.security.oauth2.client.web.AuthorizationRequestRepository;
import org.springframework.security.oauth2.core.endpoint.OAuth2AuthorizationRequest; import org.springframework.security.oauth2.core.endpoint.OAuth2AuthorizationRequest;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.queue.util.TbCoreComponent;
import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletRequest;
import javax.servlet.http.HttpServletResponse; import javax.servlet.http.HttpServletResponse;
@Component @Component
@TbCoreComponent
public class HttpCookieOAuth2AuthorizationRequestRepository implements AuthorizationRequestRepository<OAuth2AuthorizationRequest> { public class HttpCookieOAuth2AuthorizationRequestRepository implements AuthorizationRequestRepository<OAuth2AuthorizationRequest> {
public static final String OAUTH2_AUTHORIZATION_REQUEST_COOKIE_NAME = "oauth2_auth_request"; public static final String OAUTH2_AUTHORIZATION_REQUEST_COOKIE_NAME = "oauth2_auth_request";
public static final String PREV_URI_PARAMETER = "prevUri"; public static final String PREV_URI_PARAMETER = "prevUri";

229
application/src/test/java/org/thingsboard/server/controller/DeviceConnectivityControllerTest.java

@ -18,8 +18,6 @@ package org.thingsboard.server.controller;
import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode; import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.ListeningExecutorService;
import com.google.common.util.concurrent.MoreExecutors;
import org.junit.After; import org.junit.After;
import org.junit.Assert; import org.junit.Assert;
import org.junit.Before; import org.junit.Before;
@ -32,7 +30,6 @@ import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.ContextConfiguration;
import org.springframework.test.context.TestPropertySource; import org.springframework.test.context.TestPropertySource;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.common.data.AdminSettings; import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfile;
@ -46,7 +43,6 @@ import org.thingsboard.server.common.data.device.profile.DefaultDeviceProfileCon
import org.thingsboard.server.common.data.device.profile.DeviceProfileData; import org.thingsboard.server.common.data.device.profile.DeviceProfileData;
import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.MqttDeviceProfileTransportConfiguration;
import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.data.security.DeviceCredentialsType;
@ -58,14 +54,16 @@ import java.nio.file.Path;
import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThat;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.CA_ROOT_CERT_PEM;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.COAP; import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.COAP;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.COAPS; import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.COAPS;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.DOCKER; import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.DOCKER;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.HTTP; import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.HTTP;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.HTTPS; import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.HTTPS;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.LINUX;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.MQTT; import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.MQTT;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.MQTTS; import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.MQTTS;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.CA_ROOT_CERT_PEM; import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.WINDOWS;
@TestPropertySource(properties = { @TestPropertySource(properties = {
"device.connectivity.mqtts.pem_cert_file=/tmp/" + CA_ROOT_CERT_PEM "device.connectivity.mqtts.pem_cert_file=/tmp/" + CA_ROOT_CERT_PEM
@ -294,6 +292,140 @@ public class DeviceConnectivityControllerTest extends AbstractControllerTest {
DEVICE_TELEMETRY_TOPIC, credentials.getCredentialsId())); DEVICE_TELEMETRY_TOPIC, credentials.getCredentialsId()));
} }
@Test
public void testFetchGatewayLaunchCommands() throws Exception {
Device device = new Device();
device.setName("My device");
device.setType("default");
ObjectNode additionalInfo = JacksonUtil.newObjectNode();
additionalInfo.put("gateway", true);
device.setAdditionalInfo(additionalInfo);
Device savedDevice = doPost("/api/device", device, Device.class);
DeviceCredentials credentials =
doGet("/api/device/" + savedDevice.getId().getId() + "/credentials", DeviceCredentials.class);
JsonNode commands =
doGetTyped("/api/device-connectivity/gateway-launch/" + savedDevice.getId().getId(), new TypeReference<>() {
});
JsonNode dockerMqttCommands = commands.get(MQTT);
assertThat(dockerMqttCommands.get(LINUX).asText()).isEqualTo(String.format("docker run -it -v ~/.tb-gateway/logs:/thingsboard_gateway/logs -v ~/.tb-gateway/extensions:/thingsboard_gateway/extensions -v ~/.tb-gateway/config:/thingsboard_gateway/config --network=host -p 5000:5000 --name tbGatewayLocalhost -e host=localhost -e port=1883 -e accessToken=%s --restart always thingsboard/tb-gateway", credentials.getCredentialsId()));
assertThat(dockerMqttCommands.get(WINDOWS).asText()).isEqualTo("docker run -it -v %HOMEPATH%/tb-gateway/logs:/thingsboard_gateway/logs -v %HOMEPATH%/tb-gateway/extensions:/thingsboard_gateway/extensions -v %HOMEPATH%/tb-gateway/config:/thingsboard_gateway/config --network=host -p 5000:5000 --name tbGatewayLocalhost -e host=localhost -e port=1883 -e accessToken=" + credentials.getCredentialsId() + " --restart always thingsboard/tb-gateway");
JsonNode dockerMqttsCommands = commands.get(MQTTS);
assertThat(dockerMqttsCommands.get(LINUX).asText()).isEqualTo(String.format("docker run -it -v ~/.tb-gateway/logs:/thingsboard_gateway/logs -v ~/.tb-gateway/extensions:/thingsboard_gateway/extensions -v ~/.tb-gateway/config:/thingsboard_gateway/config --network=host -p 5000:5000 --name tbGatewayLocalhost -e host=localhost -e port=8883 -e accessToken=%s --restart always thingsboard/tb-gateway", credentials.getCredentialsId()));
assertThat(dockerMqttsCommands.get(WINDOWS).asText()).isEqualTo("docker run -it -v %HOMEPATH%/tb-gateway/logs:/thingsboard_gateway/logs -v %HOMEPATH%/tb-gateway/extensions:/thingsboard_gateway/extensions -v %HOMEPATH%/tb-gateway/config:/thingsboard_gateway/config --network=host -p 5000:5000 --name tbGatewayLocalhost -e host=localhost -e port=8883 -e accessToken=" + credentials.getCredentialsId() + " --restart always thingsboard/tb-gateway");
}
@Test
public void testFetchPublishTelemetryCommandsForDeviceWithIpV6LocalhostAddress() throws Exception {
loginSysAdmin();
setConnectivityHost("::1");
loginTenantAdmin();
Device device = new Device();
device.setName("My device");
device.setType("default");
Device savedDevice = doPost("/api/device", device, Device.class);
DeviceCredentials credentials =
doGet("/api/device/" + savedDevice.getId().getId() + "/credentials", DeviceCredentials.class);
JsonNode commands =
doGetTyped("/api/device-connectivity/" + savedDevice.getId().getId(), new TypeReference<>() {
});
assertThat(commands).hasSize(3);
JsonNode httpCommands = commands.get(HTTP);
assertThat(httpCommands.get(HTTP).asText()).isEqualTo(String.format("curl -v -X POST http://[::1]:8080/api/v1/%s/telemetry " +
"--header Content-Type:application/json --data \"{temperature:25}\"",
credentials.getCredentialsId()));
assertThat(httpCommands.get(HTTPS).asText()).isEqualTo(String.format("curl -v -X POST https://[::1]/api/v1/%s/telemetry " +
"--header Content-Type:application/json --data \"{temperature:25}\"",
credentials.getCredentialsId()));
JsonNode mqttCommands = commands.get(MQTT);
assertThat(mqttCommands.get(MQTT).asText()).isEqualTo(String.format("mosquitto_pub -d -q 1 -h ::1 -p 1883 -t v1/devices/me/telemetry " +
"-u \"%s\" -m \"{temperature:25}\"", credentials.getCredentialsId()));
assertThat(mqttCommands.get(MQTTS).get(0).asText()).isEqualTo("curl -f -S -o ca-root.pem http://localhost:80/api/device-connectivity/mqtts/certificate/download");
assertThat(mqttCommands.get(MQTTS).get(1).asText()).isEqualTo(String.format("mosquitto_pub -d -q 1 --cafile ca-root.pem -h ::1 -p 8883 " +
"-t v1/devices/me/telemetry -u \"%s\" -m \"{temperature:25}\"", credentials.getCredentialsId()));
JsonNode dockerMqttCommands = commands.get(MQTT).get(DOCKER);
assertThat(dockerMqttCommands.get(MQTT).asText()).isEqualTo(String.format("docker run --rm -it --network=host thingsboard/mosquitto-clients mosquitto_pub -d -q 1 -h ::1" +
" -p 1883 -t v1/devices/me/telemetry -u \"%s\" -m \"{temperature:25}\"", credentials.getCredentialsId()));
assertThat(dockerMqttCommands.get(MQTTS).asText()).isEqualTo(String.format("docker run --rm -it --network=host thingsboard/mosquitto-clients " +
"/bin/sh -c \"curl -f -S -o ca-root.pem http://localhost:80/api/device-connectivity/mqtts/certificate/download && " +
"mosquitto_pub -d -q 1 --cafile ca-root.pem -h ::1 -p 8883 -t v1/devices/me/telemetry -u \"%s\" -m \"{temperature:25}\"\"",
credentials.getCredentialsId()));
JsonNode linuxCoapCommands = commands.get(COAP);
assertThat(linuxCoapCommands.get(COAP).asText()).isEqualTo(String.format("coap-client -v 6 -m POST coap://[::1]:5683/api/v1/%s/telemetry " +
"-t json -e \"{temperature:25}\"", credentials.getCredentialsId()));
assertThat(linuxCoapCommands.get(COAPS).asText()).isEqualTo(String.format("coap-client-openssl -v 6 -m POST coaps://[::1]:5684/api/v1/%s/telemetry" +
" -t json -e \"{temperature:25}\"", credentials.getCredentialsId()));
JsonNode dockerCoapCommands = commands.get(COAP).get(DOCKER);
assertThat(dockerCoapCommands.get(COAP).asText()).isEqualTo(String.format("docker run --rm -it --network=host" +
" thingsboard/coap-clients coap-client -v 6 -m POST coap://[::1]:5683/api/v1/%s/telemetry -t json -e \"{temperature:25}\"", credentials.getCredentialsId()));
assertThat(dockerCoapCommands.get(COAPS).asText()).isEqualTo(String.format("docker run --rm -it --network=host" +
" thingsboard/coap-clients coap-client-openssl -v 6 -m POST coaps://[::1]:5684/api/v1/%s/telemetry -t json -e \"{temperature:25}\"", credentials.getCredentialsId()));
}
@Test
public void testFetchPublishTelemetryCommandsForDeviceWithIpV6Address() throws Exception {
loginSysAdmin();
setConnectivityHost("1:1:1:1:1:1:1:1");
loginTenantAdmin();
Device device = new Device();
device.setName("My device");
device.setType("default");
Device savedDevice = doPost("/api/device", device, Device.class);
DeviceCredentials credentials =
doGet("/api/device/" + savedDevice.getId().getId() + "/credentials", DeviceCredentials.class);
JsonNode commands =
doGetTyped("/api/device-connectivity/" + savedDevice.getId().getId(), new TypeReference<>() {
});
assertThat(commands).hasSize(3);
JsonNode httpCommands = commands.get(HTTP);
assertThat(httpCommands.get(HTTP).asText()).isEqualTo(String.format("curl -v -X POST http://[1:1:1:1:1:1:1:1]:8080/api/v1/%s/telemetry " +
"--header Content-Type:application/json --data \"{temperature:25}\"",
credentials.getCredentialsId()));
assertThat(httpCommands.get(HTTPS).asText()).isEqualTo(String.format("curl -v -X POST https://[1:1:1:1:1:1:1:1]/api/v1/%s/telemetry " +
"--header Content-Type:application/json --data \"{temperature:25}\"",
credentials.getCredentialsId()));
JsonNode mqttCommands = commands.get(MQTT);
assertThat(mqttCommands.get(MQTT).asText()).isEqualTo(String.format("mosquitto_pub -d -q 1 -h 1:1:1:1:1:1:1:1 -p 1883 -t v1/devices/me/telemetry " +
"-u \"%s\" -m \"{temperature:25}\"", credentials.getCredentialsId()));
assertThat(mqttCommands.get(MQTTS).get(0).asText()).isEqualTo("curl -f -S -o ca-root.pem http://localhost:80/api/device-connectivity/mqtts/certificate/download");
assertThat(mqttCommands.get(MQTTS).get(1).asText()).isEqualTo(String.format("mosquitto_pub -d -q 1 --cafile ca-root.pem -h 1:1:1:1:1:1:1:1 -p 8883 " +
"-t v1/devices/me/telemetry -u \"%s\" -m \"{temperature:25}\"", credentials.getCredentialsId()));
JsonNode dockerMqttCommands = commands.get(MQTT).get(DOCKER);
assertThat(dockerMqttCommands.get(MQTT).asText()).isEqualTo(String.format("docker run --rm -it thingsboard/mosquitto-clients mosquitto_pub -d -q 1 -h 1:1:1:1:1:1:1:1" +
" -p 1883 -t v1/devices/me/telemetry -u \"%s\" -m \"{temperature:25}\"", credentials.getCredentialsId()));
assertThat(dockerMqttCommands.get(MQTTS).asText()).isEqualTo(String.format("docker run --rm -it thingsboard/mosquitto-clients " +
"/bin/sh -c \"curl -f -S -o ca-root.pem http://localhost:80/api/device-connectivity/mqtts/certificate/download && " +
"mosquitto_pub -d -q 1 --cafile ca-root.pem -h 1:1:1:1:1:1:1:1 -p 8883 -t v1/devices/me/telemetry -u \"%s\" -m \"{temperature:25}\"\"",
credentials.getCredentialsId()));
JsonNode linuxCoapCommands = commands.get(COAP);
assertThat(linuxCoapCommands.get(COAP).asText()).isEqualTo(String.format("coap-client -v 6 -m POST coap://[1:1:1:1:1:1:1:1]:5683/api/v1/%s/telemetry " +
"-t json -e \"{temperature:25}\"", credentials.getCredentialsId()));
assertThat(linuxCoapCommands.get(COAPS).asText()).isEqualTo(String.format("coap-client-openssl -v 6 -m POST coaps://[1:1:1:1:1:1:1:1]:5684/api/v1/%s/telemetry" +
" -t json -e \"{temperature:25}\"", credentials.getCredentialsId()));
JsonNode dockerCoapCommands = commands.get(COAP).get(DOCKER);
assertThat(dockerCoapCommands.get(COAP).asText()).isEqualTo(String.format("docker run --rm -it" +
" thingsboard/coap-clients coap-client -v 6 -m POST coap://[1:1:1:1:1:1:1:1]:5683/api/v1/%s/telemetry -t json -e \"{temperature:25}\"", credentials.getCredentialsId()));
assertThat(dockerCoapCommands.get(COAPS).asText()).isEqualTo(String.format("docker run --rm -it" +
" thingsboard/coap-clients coap-client-openssl -v 6 -m POST coaps://[1:1:1:1:1:1:1:1]:5684/api/v1/%s/telemetry -t json -e \"{temperature:25}\"", credentials.getCredentialsId()));
}
@Test @Test
public void testFetchPublishTelemetryCommandsForDeviceWithMqttBasicCreds() throws Exception { public void testFetchPublishTelemetryCommandsForDeviceWithMqttBasicCreds() throws Exception {
Device device = new Device(); Device device = new Device();
@ -518,47 +650,7 @@ public class DeviceConnectivityControllerTest extends AbstractControllerTest {
public void testFetchPublishTelemetryCommandsForDefaultDeviceIfHostIsNotLocalhost() throws Exception { public void testFetchPublishTelemetryCommandsForDefaultDeviceIfHostIsNotLocalhost() throws Exception {
loginSysAdmin(); loginSysAdmin();
ObjectNode config = JacksonUtil.newObjectNode(); setConnectivityHost("test.domain");
ObjectNode http = JacksonUtil.newObjectNode();
http.put("enabled", true);
http.put("host", "test.domain");
http.put("port", 8080);
config.set("http", http);
ObjectNode https = JacksonUtil.newObjectNode();
https.put("enabled", true);
https.put("host", "test.domain");
https.put("port", 443);
config.set("https", https);
ObjectNode mqtt = JacksonUtil.newObjectNode();
mqtt.put("enabled", true);
mqtt.put("host", "test.domain");
mqtt.put("port", 1883);
config.set("mqtt", mqtt);
ObjectNode mqtts = JacksonUtil.newObjectNode();
mqtts.put("enabled", true);
mqtts.put("host", "test.domain");
mqtts.put("port", 8883);
config.set("mqtts", mqtts);
ObjectNode coap = JacksonUtil.newObjectNode();
coap.put("enabled", true);
coap.put("host", "test.domain");
coap.put("port", 5683);
config.set("coap", coap);
ObjectNode coaps = JacksonUtil.newObjectNode();
coaps.put("enabled", true);
coaps.put("host", "test.domain");
coaps.put("port", 5684);
config.set("coaps", coaps);
AdminSettings adminSettings = doGet("/api/admin/settings/connectivity", AdminSettings.class);
adminSettings.setJsonValue(config);
doPost("/api/admin/settings", adminSettings).andExpect(status().isOk());
login("tenant2@thingsboard.org", "testPassword1"); login("tenant2@thingsboard.org", "testPassword1");
@ -612,4 +704,49 @@ public class DeviceConnectivityControllerTest extends AbstractControllerTest {
assertThat(dockerCoapCommands.get(COAPS).asText()).isEqualTo(String.format("docker run --rm -it " + assertThat(dockerCoapCommands.get(COAPS).asText()).isEqualTo(String.format("docker run --rm -it " +
"thingsboard/coap-clients coap-client-openssl -v 6 -m POST coaps://test.domain:5684/api/v1/%s/telemetry -t json -e \"{temperature:25}\"", credentials.getCredentialsId())); "thingsboard/coap-clients coap-client-openssl -v 6 -m POST coaps://test.domain:5684/api/v1/%s/telemetry -t json -e \"{temperature:25}\"", credentials.getCredentialsId()));
} }
private void setConnectivityHost(String host) throws Exception {
ObjectNode config = JacksonUtil.newObjectNode();
ObjectNode http = JacksonUtil.newObjectNode();
http.put("enabled", true);
http.put("host", host);
http.put("port", 8080);
config.set("http", http);
ObjectNode https = JacksonUtil.newObjectNode();
https.put("enabled", true);
https.put("host", host);
https.put("port", 443);
config.set("https", https);
ObjectNode mqtt = JacksonUtil.newObjectNode();
mqtt.put("enabled", true);
mqtt.put("host", host);
mqtt.put("port", 1883);
config.set("mqtt", mqtt);
ObjectNode mqtts = JacksonUtil.newObjectNode();
mqtts.put("enabled", true);
mqtts.put("host", host);
mqtts.put("port", 8883);
config.set("mqtts", mqtts);
ObjectNode coap = JacksonUtil.newObjectNode();
coap.put("enabled", true);
coap.put("host", host);
coap.put("port", 5683);
config.set("coap", coap);
ObjectNode coaps = JacksonUtil.newObjectNode();
coaps.put("enabled", true);
coaps.put("host", host);
coaps.put("port", 5684);
config.set("coaps", coaps);
AdminSettings adminSettings = doGet("/api/admin/settings/connectivity", AdminSettings.class);
adminSettings.setJsonValue(config);
doPost("/api/admin/settings", adminSettings).andExpect(status().isOk());
}
} }

29
application/src/test/java/org/thingsboard/server/controller/DeviceControllerTest.java

@ -51,7 +51,6 @@ import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.alarm.Alarm; import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.alarm.AlarmInfo; import org.thingsboard.server.common.data.alarm.AlarmInfo;
import org.thingsboard.server.common.data.alarm.AlarmSeverity; import org.thingsboard.server.common.data.alarm.AlarmSeverity;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.CustomerId;
@ -74,6 +73,7 @@ import org.thingsboard.server.dao.exception.DeviceCredentialsValidationException
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.service.gateway_device.GatewayNotificationsService; import org.thingsboard.server.service.gateway_device.GatewayNotificationsService;
import org.thingsboard.server.service.state.DeviceStateService;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.List; import java.util.List;
@ -82,6 +82,8 @@ import java.util.concurrent.TimeUnit;
import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThat;
import static org.hamcrest.Matchers.containsString; import static org.hamcrest.Matchers.containsString;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.Mockito.never; import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times; import static org.mockito.Mockito.times;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@ -106,6 +108,9 @@ public class DeviceControllerTest extends AbstractControllerTest {
@SpyBean @SpyBean
private GatewayNotificationsService gatewayNotificationsService; private GatewayNotificationsService gatewayNotificationsService;
@SpyBean
private DeviceStateService deviceStateService;
@Autowired @Autowired
private DeviceDao deviceDao; private DeviceDao deviceDao;
@ -1340,6 +1345,28 @@ public class DeviceControllerTest extends AbstractControllerTest {
ActionType.ASSIGNED_TO_TENANT, savedDifferentTenant.getId().getId().toString(), savedDifferentTenant.getTitle()); ActionType.ASSIGNED_TO_TENANT, savedDifferentTenant.getId().getId().toString(), savedDifferentTenant.getTitle());
testNotificationUpdateGatewayNever(); testNotificationUpdateGatewayNever();
Mockito.verify(deviceStateService, times(1)).onQueueMsg(
argThat(proto ->
proto.getTenantIdMSB() == savedTenant.getUuidId().getMostSignificantBits() &&
proto.getTenantIdLSB() == savedTenant.getUuidId().getLeastSignificantBits() &&
proto.getDeviceIdMSB() == savedDevice.getUuidId().getMostSignificantBits() &&
proto.getDeviceIdLSB() == savedDevice.getUuidId().getLeastSignificantBits() &&
proto.getDeleted()
),
any()
);
Mockito.verify(deviceStateService, times(1)).onQueueMsg(
argThat(proto ->
proto.getTenantIdMSB() == savedDifferentTenant.getUuidId().getMostSignificantBits() &&
proto.getTenantIdLSB() == savedDifferentTenant.getUuidId().getLeastSignificantBits() &&
proto.getDeviceIdMSB() == savedDevice.getUuidId().getMostSignificantBits() &&
proto.getDeviceIdLSB() == savedDevice.getUuidId().getLeastSignificantBits() &&
proto.getAdded()
),
any()
);
login("tenant9@thingsboard.org", "testPassword1"); login("tenant9@thingsboard.org", "testPassword1");
Device foundDevice1 = doGet("/api/device/" + assignedDevice.getId().getId(), Device.class); Device foundDevice1 = doGet("/api/device/" + assignedDevice.getId().getId(), Device.class);

30
application/src/test/java/org/thingsboard/server/edge/RelationEdgeTest.java

@ -129,14 +129,24 @@ public class RelationEdgeTest extends AbstractEdgeTest {
Device device = findDeviceByName("Edge Device 1"); Device device = findDeviceByName("Edge Device 1");
Asset asset = findAssetByName("Edge Asset 1"); Asset asset = findAssetByName("Edge Asset 1");
EntityRelation relation = new EntityRelation(); EntityRelation deviceToAssetRelation = new EntityRelation();
relation.setType("test"); deviceToAssetRelation.setType("test");
relation.setFrom(device.getId()); deviceToAssetRelation.setFrom(device.getId());
relation.setTo(asset.getId()); deviceToAssetRelation.setTo(asset.getId());
relation.setTypeGroup(RelationTypeGroup.COMMON); deviceToAssetRelation.setTypeGroup(RelationTypeGroup.COMMON);
edgeImitator.expectMessageAmount(1); edgeImitator.expectMessageAmount(1);
doPost("/api/relation", relation); doPost("/api/relation", deviceToAssetRelation);
Assert.assertTrue(edgeImitator.waitForMessages());
EntityRelation assetToTenantRelation = new EntityRelation();
assetToTenantRelation.setType("test");
assetToTenantRelation.setFrom(asset.getId());
assetToTenantRelation.setTo(tenantId);
assetToTenantRelation.setTypeGroup(RelationTypeGroup.COMMON);
edgeImitator.expectMessageAmount(1);
doPost("/api/relation", assetToTenantRelation);
Assert.assertTrue(edgeImitator.waitForMessages()); Assert.assertTrue(edgeImitator.waitForMessages());
UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder(); UplinkMsg.Builder uplinkMsgBuilder = UplinkMsg.newBuilder();
@ -159,16 +169,16 @@ public class RelationEdgeTest extends AbstractEdgeTest {
Assert.assertTrue(latestMessage instanceof RelationUpdateMsg); Assert.assertTrue(latestMessage instanceof RelationUpdateMsg);
RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage; RelationUpdateMsg relationUpdateMsg = (RelationUpdateMsg) latestMessage;
Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType()); Assert.assertEquals(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, relationUpdateMsg.getMsgType());
Assert.assertEquals(relation.getType(), relationUpdateMsg.getType()); Assert.assertEquals(deviceToAssetRelation.getType(), relationUpdateMsg.getType());
UUID fromUUID = new UUID(relationUpdateMsg.getFromIdMSB(), relationUpdateMsg.getFromIdLSB()); UUID fromUUID = new UUID(relationUpdateMsg.getFromIdMSB(), relationUpdateMsg.getFromIdLSB());
EntityId fromEntityId = EntityIdFactory.getByTypeAndUuid(relationUpdateMsg.getFromEntityType(), fromUUID); EntityId fromEntityId = EntityIdFactory.getByTypeAndUuid(relationUpdateMsg.getFromEntityType(), fromUUID);
Assert.assertEquals(relation.getFrom(), fromEntityId); Assert.assertEquals(deviceToAssetRelation.getFrom(), fromEntityId);
UUID toUUID = new UUID(relationUpdateMsg.getToIdMSB(), relationUpdateMsg.getToIdLSB()); UUID toUUID = new UUID(relationUpdateMsg.getToIdMSB(), relationUpdateMsg.getToIdLSB());
EntityId toEntityId = EntityIdFactory.getByTypeAndUuid(relationUpdateMsg.getToEntityType(), toUUID); EntityId toEntityId = EntityIdFactory.getByTypeAndUuid(relationUpdateMsg.getToEntityType(), toUUID);
Assert.assertEquals(relation.getTo(), toEntityId); Assert.assertEquals(deviceToAssetRelation.getTo(), toEntityId);
Assert.assertEquals(relation.getTypeGroup().name(), relationUpdateMsg.getTypeGroup()); Assert.assertEquals(deviceToAssetRelation.getTypeGroup().name(), relationUpdateMsg.getTypeGroup());
} }
} }

4
common/cache/src/main/java/org/thingsboard/server/cache/RedisTbTransactionalCache.java

@ -84,8 +84,10 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
} else if (Arrays.equals(rawValue, BINARY_NULL_VALUE)) { } else if (Arrays.equals(rawValue, BINARY_NULL_VALUE)) {
return SimpleTbCacheValueWrapper.empty(); return SimpleTbCacheValueWrapper.empty();
} else { } else {
long startTime = System.nanoTime();
V value = valueSerializer.deserialize(key, rawValue); V value = valueSerializer.deserialize(key, rawValue);
if (value != null) { if (value != null) {
fstStatsService.recordDecodeTime(value.getClass(), startTime);
fstStatsService.incrementDecode(value.getClass()); fstStatsService.incrementDecode(value.getClass());
} }
return SimpleTbCacheValueWrapper.wrap(value); return SimpleTbCacheValueWrapper.wrap(value);
@ -198,7 +200,9 @@ public abstract class RedisTbTransactionalCache<K extends Serializable, V extend
return BINARY_NULL_VALUE; return BINARY_NULL_VALUE;
} else { } else {
try { try {
long startTime = System.nanoTime();
var bytes = valueSerializer.serialize(value); var bytes = valueSerializer.serialize(value);
fstStatsService.recordEncodeTime(value.getClass(), startTime);
fstStatsService.incrementEncode(value.getClass()); fstStatsService.incrementEncode(value.getClass());
return bytes; return bytes;
} catch (Exception e) { } catch (Exception e) {

4
common/cluster-api/src/main/java/org/thingsboard/server/cluster/TbClusterService.java

@ -83,7 +83,9 @@ public interface TbClusterService extends TbQueueClusterService {
void onDeviceUpdated(Device device, Device old); void onDeviceUpdated(Device device, Device old);
void onDeviceDeleted(Device device, TbQueueCallback callback); void onDeviceDeleted(TenantId tenantId, Device device, TbQueueCallback callback);
void onDeviceAssignedToTenant(TenantId oldTenantId, Device device);
void onResourceChange(TbResource resource, TbQueueCallback callback); void onResourceChange(TbResource resource, TbQueueCallback callback);

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

@ -1016,6 +1016,13 @@ message RemoveRpcActorMsgProto {
int64 deviceIdLSB = 6; int64 deviceIdLSB = 6;
} }
message DeviceDeleteMsgProto {
int64 tenantIdMSB = 1;
int64 tenantIdLSB = 2;
int64 deviceIdMSB = 3;
int64 deviceIdLSB = 4;
}
message ToDeviceActorNotificationMsgProto { message ToDeviceActorNotificationMsgProto {
DeviceEdgeUpdateMsgProto deviceEdgeUpdateMsg = 1; DeviceEdgeUpdateMsgProto deviceEdgeUpdateMsg = 1;
DeviceNameOrTypeUpdateMsgProto deviceNameOrTypeMsg = 2; DeviceNameOrTypeUpdateMsgProto deviceNameOrTypeMsg = 2;
@ -1024,6 +1031,7 @@ message ToDeviceActorNotificationMsgProto {
ToDeviceRpcRequestActorMsgProto toDeviceRpcRequestMsg = 5; ToDeviceRpcRequestActorMsgProto toDeviceRpcRequestMsg = 5;
FromDeviceRpcResponseActorMsgProto fromDeviceRpcResponseMsg = 6; FromDeviceRpcResponseActorMsgProto fromDeviceRpcResponseMsg = 6;
RemoveRpcActorMsgProto removeRpcActorMsg = 7; RemoveRpcActorMsgProto removeRpcActorMsg = 7;
DeviceDeleteMsgProto deviceDeleteMsg = 8;
} }
/** /**

4
common/data/src/main/java/org/thingsboard/server/common/data/FstStatsService.java

@ -21,4 +21,8 @@ public interface FstStatsService {
void incrementDecode(Class<?> clazz); void incrementDecode(Class<?> clazz);
void recordEncodeTime(Class<?> clazz, long startTime);
void recordDecodeTime(Class<?> clazz, long startTime);
} }

2
common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java

@ -93,6 +93,8 @@ public enum MsgType {
DEVICE_NAME_OR_TYPE_UPDATE_TO_DEVICE_ACTOR_MSG, DEVICE_NAME_OR_TYPE_UPDATE_TO_DEVICE_ACTOR_MSG,
DEVICE_DELETE_TO_DEVICE_ACTOR_MSG,
DEVICE_EDGE_UPDATE_TO_DEVICE_ACTOR_MSG, DEVICE_EDGE_UPDATE_TO_DEVICE_ACTOR_MSG,
DEVICE_RPC_REQUEST_TO_DEVICE_ACTOR_MSG, DEVICE_RPC_REQUEST_TO_DEVICE_ACTOR_MSG,

36
common/message/src/main/java/org/thingsboard/server/common/msg/rule/engine/DeviceDeleteMsg.java

@ -0,0 +1,36 @@
/**
* 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.msg.rule.engine;
import lombok.Data;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.MsgType;
import org.thingsboard.server.common.msg.ToDeviceActorNotificationMsg;
@Data
public class DeviceDeleteMsg implements ToDeviceActorNotificationMsg {
private static final long serialVersionUID = 4679029228395462172L;
private final TenantId tenantId;
private final DeviceId deviceId;
@Override
public MsgType getMsgType() {
return MsgType.DEVICE_DELETE_TO_DEVICE_ACTOR_MSG;
}
}

8
common/queue/src/main/java/org/thingsboard/server/queue/util/ProtoWithFSTService.java

@ -36,8 +36,12 @@ public class ProtoWithFSTService implements DataDecodingEncodingService {
@Override @Override
public <T> Optional<T> decode(byte[] byteArray) { public <T> Optional<T> decode(byte[] byteArray) {
try { try {
long startTime = System.nanoTime();
Optional<T> optional = Optional.ofNullable(FSTUtils.decode(byteArray)); Optional<T> optional = Optional.ofNullable(FSTUtils.decode(byteArray));
optional.ifPresent(obj -> fstStatsService.incrementDecode(obj.getClass())); optional.ifPresent(obj -> {
fstStatsService.recordDecodeTime(obj.getClass(), startTime);
fstStatsService.incrementDecode(obj.getClass());
});
return optional; return optional;
} catch (IllegalArgumentException e) { } catch (IllegalArgumentException e) {
log.error("Error during deserialization message, [{}]", e.getMessage()); log.error("Error during deserialization message, [{}]", e.getMessage());
@ -48,7 +52,9 @@ public class ProtoWithFSTService implements DataDecodingEncodingService {
@Override @Override
public <T> byte[] encode(T msq) { public <T> byte[] encode(T msq) {
long startTime = System.nanoTime();
var bytes = FSTUtils.encode(msq); var bytes = FSTUtils.encode(msq);
fstStatsService.recordEncodeTime(msq.getClass(), startTime);
fstStatsService.incrementEncode(msq.getClass()); fstStatsService.incrementEncode(msq.getClass());
return bytes; return bytes;
} }

28
common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbUtils.java

@ -182,12 +182,12 @@ public class TbUtils {
return TbJson.parse(ctx, jsonStr); return TbJson.parse(ctx, jsonStr);
} }
public static String bytesToString(List<Byte> bytesList) { public static String bytesToString(List<?> bytesList) {
byte[] bytes = bytesFromList(bytesList); byte[] bytes = bytesFromList(bytesList);
return new String(bytes); return new String(bytes);
} }
public static String bytesToString(List<Byte> bytesList, String charsetName) throws UnsupportedEncodingException { public static String bytesToString(List<?> bytesList, String charsetName) throws UnsupportedEncodingException {
byte[] bytes = bytesFromList(bytesList); byte[] bytes = bytesFromList(bytesList);
return new String(bytes, charsetName); return new String(bytes, charsetName);
} }
@ -210,10 +210,20 @@ public class TbUtils {
} }
} }
private static byte[] bytesFromList(List<Byte> bytesList) { private static byte[] bytesFromList(List<?> bytesList) {
byte[] bytes = new byte[bytesList.size()]; byte[] bytes = new byte[bytesList.size()];
for (int i = 0; i < bytesList.size(); i++) { for (int i = 0; i < bytesList.size(); i++) {
bytes[i] = bytesList.get(i); Object objectVal = bytesList.get(i);
if (objectVal instanceof Integer) {
bytes[i] = isValidIntegerToByte((Integer) objectVal);
} else if (objectVal instanceof String) {
bytes[i] = isValidIntegerToByte(parseInt((String) objectVal));
} else if (objectVal instanceof Byte) {
bytes[i] = (byte) objectVal;
} else {
throw new NumberFormatException("The value '" + objectVal + "' could not be correctly converted to a byte. " +
"Must be a HexDecimal/String/Integer/Byte format !");
}
} }
return bytes; return bytes;
} }
@ -643,7 +653,7 @@ public class TbUtils {
} }
} }
public static boolean isValidRadix(String value, int radix) { private static boolean isValidRadix(String value, int radix) {
for (int i = 0; i < value.length(); i++) { for (int i = 0; i < value.length(); i++) {
if (i == 0 && value.charAt(i) == '-') { if (i == 0 && value.charAt(i) == '-') {
if (value.length() == 1) if (value.length() == 1)
@ -657,4 +667,12 @@ public class TbUtils {
return true; return true;
} }
private static byte isValidIntegerToByte (Integer val) {
if (val > 255 || val.intValue() < -128) {
throw new NumberFormatException("The value '" + val + "' could not be correctly converted to a byte. " +
"Integer to byte conversion requires the use of only 8 bits (with a range of min/max = -128/255)!");
} else {
return val.byteValue();
}
}
} }

90
common/script/script-api/src/test/java/org/thingsboard/script/api/tbel/TbUtilsTest.java

@ -31,6 +31,7 @@ import java.io.IOException;
import java.math.BigInteger; import java.math.BigInteger;
import java.nio.ByteBuffer; import java.nio.ByteBuffer;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Arrays;
import java.util.Calendar; import java.util.Calendar;
import java.util.List; import java.util.List;
import java.util.Random; import java.util.Random;
@ -385,6 +386,95 @@ public class TbUtilsTest {
Assert.assertThrows(IllegalAccessException.class, () -> TbUtils.stringToBytes(ctx, ((ExecutionHashMap) finalInputJson).get("hello"), "UTF-8")); Assert.assertThrows(IllegalAccessException.class, () -> TbUtils.stringToBytes(ctx, ((ExecutionHashMap) finalInputJson).get("hello"), "UTF-8"));
} }
@Test
public void bytesFromList() {
byte[] arrayBytes = {(byte)0x00, (byte)0x08, (byte)0x10, (byte)0x1C, (byte)0xFF, (byte)0xFC, (byte)0xAD, (byte)0x88, (byte)0x75, (byte)0x74, (byte)0x8A, (byte)0x82};
Object[] arrayMix = { "0x00", 8, "16", "0x1C", 255, (byte)0xFC, 173, 136, 117, 116, -118, "-126"};
String expected = new String(arrayBytes);
ArrayList<Byte> listBytes = new ArrayList<>(arrayBytes.length);
for (Byte element : arrayBytes) {
listBytes.add(element);
}
Assert.assertEquals(expected, TbUtils.bytesToString(listBytes));
ArrayList<Object> listMix = new ArrayList<>(arrayMix.length);
for (Object element : arrayMix) {
listMix.add(element);
}
Assert.assertEquals(expected, TbUtils.bytesToString(listMix));
}
@Test
public void bytesFromList_Error() {
List<String> listHex = new ArrayList<>();
listHex.add("0xFG");
try {
TbUtils.bytesToString(listHex);
Assert.fail("Should throw NumberFormatException");
} catch (NumberFormatException e) {
Assert.assertTrue(e.getMessage().contains("Failed radix: [16] for value: \"FG\"!"));
}
listHex.add(0, "1F");
try {
TbUtils.bytesToString(listHex);
Assert.fail("Should throw NumberFormatException");
} catch (NumberFormatException e) {
Assert.assertTrue(e.getMessage().contains("Failed radix: [10] for value: \"1F\"!"));
}
List<String> listIntString = new ArrayList<>();
listIntString.add("-129");
try {
TbUtils.bytesToString(listIntString);
Assert.fail("Should throw NumberFormatException");
} catch (NumberFormatException e) {
Assert.assertTrue(e.getMessage().contains("The value '-129' could not be correctly converted to a byte. " +
"Integer to byte conversion requires the use of only 8 bits (with a range of min/max = -128/255)!"));
}
listIntString.add(0, "256");
try {
TbUtils.bytesToString(listIntString);
Assert.fail("Should throw NumberFormatException");
} catch (NumberFormatException e) {
Assert.assertTrue(e.getMessage().contains("The value '256' could not be correctly converted to a byte. " +
"Integer to byte conversion requires the use of only 8 bits (with a range of min/max = -128/255)!"));
}
ArrayList<Integer> listIntBytes = new ArrayList<>();
listIntBytes.add(-129);
try {
TbUtils.bytesToString(listIntBytes);
Assert.fail("Should throw NumberFormatException");
} catch (NumberFormatException e) {
Assert.assertTrue(e.getMessage().contains("The value '-129' could not be correctly converted to a byte. " +
"Integer to byte conversion requires the use of only 8 bits (with a range of min/max = -128/255)!"));
}
listIntBytes.add(0, 256);
try {
TbUtils.bytesToString(listIntBytes);
Assert.fail("Should throw NumberFormatException");
} catch (NumberFormatException e) {
Assert.assertTrue(e.getMessage().contains("The value '256' could not be correctly converted to a byte. " +
"Integer to byte conversion requires the use of only 8 bits (with a range of min/max = -128/255)!"));
}
ArrayList<Object> listObjects = new ArrayList<>();
ArrayList<String> listStringObjects = new ArrayList<>();
listStringObjects.add("0xFD");
listObjects.add(listStringObjects);
try {
TbUtils.bytesToString(listObjects);
Assert.fail("Should throw NumberFormatException");
} catch (NumberFormatException e) {
Assert.assertTrue(e.getMessage().contains("The value '[0xFD]' could not be correctly converted to a byte. " +
"Must be a HexDecimal/String/Integer/Byte format !"));
}
}
private static List<Byte> toList(byte[] data) { private static List<Byte> toList(byte[] data) {
List<Byte> result = new ArrayList<>(data.length); List<Byte> result = new ArrayList<>(data.length);
for (Byte b : data) { for (Byte b : data) {

20
common/stats/src/main/java/org/thingsboard/server/common/stats/FstStatsServiceImpl.java

@ -15,28 +15,44 @@
*/ */
package org.thingsboard.server.common.stats; package org.thingsboard.server.common.stats;
import io.micrometer.core.instrument.Timer;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.FstStatsService; import org.thingsboard.server.common.data.FstStatsService;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
@Service @Service
public class FstStatsServiceImpl implements FstStatsService { public class FstStatsServiceImpl implements FstStatsService {
private final ConcurrentHashMap<String, StatsCounter> encodeCounters = new ConcurrentHashMap<>(); private final ConcurrentHashMap<String, StatsCounter> encodeCounters = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, StatsCounter> decodeCounters = new ConcurrentHashMap<>(); private final ConcurrentHashMap<String, StatsCounter> decodeCounters = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Timer> encodeTimers = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Timer> decodeTimer = new ConcurrentHashMap<>();
@Autowired @Autowired
private StatsFactory statsFactory; private StatsFactory statsFactory;
@Override @Override
public void incrementEncode(Class<?> clazz) { public void incrementEncode(Class<?> clazz) {
encodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fstEncode", key)).increment(); encodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fst_encode", key)).increment();
} }
@Override @Override
public void incrementDecode(Class<?> clazz) { public void incrementDecode(Class<?> clazz) {
decodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fstDecode", key)).increment(); decodeCounters.computeIfAbsent(clazz.getSimpleName(), key -> statsFactory.createStatsCounter("fst_decode", key)).increment();
}
@Override
public void recordEncodeTime(Class<?> clazz, long startTime) {
encodeTimers.computeIfAbsent(clazz.getSimpleName(),
key -> statsFactory.createTimer("fst_encode_time", "statsName", key)).record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS);
}
@Override
public void recordDecodeTime(Class<?> clazz, long startTime) {
decodeTimer.computeIfAbsent(clazz.getSimpleName(),
key -> statsFactory.createTimer("fst_decode_time", "statsName", key)).record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS);
} }
} }

20
dao/src/main/java/org/thingsboard/server/dao/device/DeviceConnectivityServiceImpl.java

@ -43,7 +43,6 @@ import org.thingsboard.server.dao.util.DeviceConnectivityUtil;
import java.io.InputStream; import java.io.InputStream;
import java.io.InputStreamReader; import java.io.InputStreamReader;
import java.net.URI;
import java.net.URISyntaxException; import java.net.URISyntaxException;
import java.nio.charset.StandardCharsets; import java.nio.charset.StandardCharsets;
import java.util.ArrayList; import java.util.ArrayList;
@ -64,6 +63,7 @@ import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.LINUX;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.MQTT; import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.MQTT;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.MQTTS; import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.MQTTS;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.WINDOWS; import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.WINDOWS;
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.getHost;
@Service("DeviceConnectivityDaoService") @Service("DeviceConnectivityDaoService")
@Slf4j @Slf4j
@ -230,7 +230,7 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService
deviceCredentials.getCredentialsType() != DeviceCredentialsType.ACCESS_TOKEN) { deviceCredentials.getCredentialsType() != DeviceCredentialsType.ACCESS_TOKEN) {
return null; return null;
} }
String hostName = getHost(baseUrl, properties); String hostName = getHost(baseUrl, properties, protocol);
String propertiesPort = properties.getPort(); String propertiesPort = properties.getPort();
String port = (propertiesPort.isEmpty() || HTTP_DEFAULT_PORT.equals(propertiesPort) || HTTPS_DEFAULT_PORT.equals(propertiesPort)) String port = (propertiesPort.isEmpty() || HTTP_DEFAULT_PORT.equals(propertiesPort) || HTTPS_DEFAULT_PORT.equals(propertiesPort))
? "" : ":" + propertiesPort; ? "" : ":" + propertiesPort;
@ -278,14 +278,14 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService
private String getMqttPublishCommand(String baseUrl, String deviceTelemetryTopic, DeviceCredentials deviceCredentials) throws URISyntaxException { private String getMqttPublishCommand(String baseUrl, String deviceTelemetryTopic, DeviceCredentials deviceCredentials) throws URISyntaxException {
DeviceConnectivityInfo properties = getConnectivity(MQTT); DeviceConnectivityInfo properties = getConnectivity(MQTT);
String mqttHost = getHost(baseUrl, properties); String mqttHost = getHost(baseUrl, properties, MQTT);
String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort(); String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort();
return DeviceConnectivityUtil.getMqttPublishCommand(MQTT, mqttHost, mqttPort, deviceTelemetryTopic, deviceCredentials); return DeviceConnectivityUtil.getMqttPublishCommand(MQTT, mqttHost, mqttPort, deviceTelemetryTopic, deviceCredentials);
} }
private List<String> getMqttsPublishCommand(String baseUrl, String deviceTelemetryTopic, DeviceCredentials deviceCredentials) throws URISyntaxException { private List<String> getMqttsPublishCommand(String baseUrl, String deviceTelemetryTopic, DeviceCredentials deviceCredentials) throws URISyntaxException {
DeviceConnectivityInfo properties = getConnectivity(MQTTS); DeviceConnectivityInfo properties = getConnectivity(MQTTS);
String mqttHost = getHost(baseUrl, properties); String mqttHost = getHost(baseUrl, properties, MQTTS);
String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort(); String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort();
String pubCommand = DeviceConnectivityUtil.getMqttPublishCommand(MQTTS, mqttHost, mqttPort, deviceTelemetryTopic, deviceCredentials); String pubCommand = DeviceConnectivityUtil.getMqttPublishCommand(MQTTS, mqttHost, mqttPort, deviceTelemetryTopic, deviceCredentials);
@ -301,7 +301,7 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService
private JsonNode getGatewayDockerCommands(String baseUrl, DeviceCredentials deviceCredentials, String mqttType) throws URISyntaxException { private JsonNode getGatewayDockerCommands(String baseUrl, DeviceCredentials deviceCredentials, String mqttType) throws URISyntaxException {
ObjectNode dockerLaunchCommands = JacksonUtil.newObjectNode(); ObjectNode dockerLaunchCommands = JacksonUtil.newObjectNode();
DeviceConnectivityInfo properties = getConnectivity(mqttType); DeviceConnectivityInfo properties = getConnectivity(mqttType);
String mqttHost = getHost(baseUrl, properties); String mqttHost = getHost(baseUrl, properties, mqttType);
String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort(); String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort();
Optional.ofNullable(DeviceConnectivityUtil.getGatewayLaunchCommand(LINUX, mqttHost, mqttPort, deviceCredentials)) Optional.ofNullable(DeviceConnectivityUtil.getGatewayLaunchCommand(LINUX, mqttHost, mqttPort, deviceCredentials))
.ifPresent(v -> dockerLaunchCommands.put(LINUX, v)); .ifPresent(v -> dockerLaunchCommands.put(LINUX, v));
@ -312,7 +312,7 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService
private String getDockerMqttPublishCommand(String protocol, String baseUrl, String deviceTelemetryTopic, DeviceCredentials deviceCredentials) throws URISyntaxException { private String getDockerMqttPublishCommand(String protocol, String baseUrl, String deviceTelemetryTopic, DeviceCredentials deviceCredentials) throws URISyntaxException {
DeviceConnectivityInfo properties = getConnectivity(protocol); DeviceConnectivityInfo properties = getConnectivity(protocol);
String mqttHost = getHost(baseUrl, properties); String mqttHost = getHost(baseUrl, properties, protocol);
String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort(); String mqttPort = properties.getPort().isEmpty() ? null : properties.getPort();
return DeviceConnectivityUtil.getDockerMqttPublishCommand(protocol, baseUrl, mqttHost, mqttPort, deviceTelemetryTopic, deviceCredentials); return DeviceConnectivityUtil.getDockerMqttPublishCommand(protocol, baseUrl, mqttHost, mqttPort, deviceTelemetryTopic, deviceCredentials);
} }
@ -352,20 +352,16 @@ public class DeviceConnectivityServiceImpl implements DeviceConnectivityService
private String getCoapPublishCommand(String protocol, String baseUrl, DeviceCredentials deviceCredentials) throws URISyntaxException { private String getCoapPublishCommand(String protocol, String baseUrl, DeviceCredentials deviceCredentials) throws URISyntaxException {
DeviceConnectivityInfo properties = getConnectivity(protocol); DeviceConnectivityInfo properties = getConnectivity(protocol);
String hostName = getHost(baseUrl, properties); String hostName = getHost(baseUrl, properties, protocol);
String port = properties.getPort().isEmpty() ? "" : ":" + properties.getPort(); String port = properties.getPort().isEmpty() ? "" : ":" + properties.getPort();
return DeviceConnectivityUtil.getCoapPublishCommand(protocol, hostName, port, deviceCredentials); return DeviceConnectivityUtil.getCoapPublishCommand(protocol, hostName, port, deviceCredentials);
} }
private String getDockerCoapPublishCommand(String protocol, String baseUrl, DeviceCredentials deviceCredentials) throws URISyntaxException { private String getDockerCoapPublishCommand(String protocol, String baseUrl, DeviceCredentials deviceCredentials) throws URISyntaxException {
DeviceConnectivityInfo properties = getConnectivity(protocol); DeviceConnectivityInfo properties = getConnectivity(protocol);
String host = getHost(baseUrl, properties); String host = getHost(baseUrl, properties, protocol);
String port = properties.getPort().isEmpty() ? "" : ":" + properties.getPort(); String port = properties.getPort().isEmpty() ? "" : ":" + properties.getPort();
return DeviceConnectivityUtil.getDockerCoapPublishCommand(protocol, host, port, deviceCredentials); return DeviceConnectivityUtil.getDockerCoapPublishCommand(protocol, host, port, deviceCredentials);
} }
private String getHost(String baseUrl, DeviceConnectivityInfo properties) throws URISyntaxException {
return properties.getHost().isEmpty() ? new URI(baseUrl).getHost() : properties.getHost();
}
} }

73
dao/src/main/java/org/thingsboard/server/dao/util/DeviceConnectivityUtil.java

@ -19,9 +19,14 @@ import org.apache.commons.lang3.StringUtils;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials; import org.thingsboard.server.common.data.device.credentials.BasicMqttCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.dao.device.DeviceConnectivityInfo;
import java.util.Arrays; import java.net.Inet6Address;
import java.util.List; import java.net.InetAddress;
import java.net.URI;
import java.net.URISyntaxException;
import java.net.UnknownHostException;
import java.util.regex.Pattern;
public class DeviceConnectivityUtil { public class DeviceConnectivityUtil {
@ -39,9 +44,11 @@ public class DeviceConnectivityUtil {
public static final String JSON_EXAMPLE_PAYLOAD = "\"{temperature:25}\""; public static final String JSON_EXAMPLE_PAYLOAD = "\"{temperature:25}\"";
public static final String DOCKER_RUN = "docker run --rm -it "; public static final String DOCKER_RUN = "docker run --rm -it ";
public static final String GATEWAY_DOCKER_RUN = "docker run -it "; public static final String GATEWAY_DOCKER_RUN = "docker run -it ";
public static final String NETWORK_HOST_PARAM = "--network=host ";
public static final String MQTT_IMAGE = "thingsboard/mosquitto-clients "; public static final String MQTT_IMAGE = "thingsboard/mosquitto-clients ";
public static final String COAP_IMAGE = "thingsboard/coap-clients "; public static final String COAP_IMAGE = "thingsboard/coap-clients ";
public static final List<String> LOCAL_HOSTS = Arrays.asList("localhost", "127.0.0.1"); private final static Pattern VALID_URL_PATTERN = Pattern.compile("^(https?)://[-a-zA-Z0-9+&@#/%?=~_|!:,.;]*[-a-zA-Z0-9+&@#/%=~_|]");
public static String getHttpPublishCommand(String protocol, String host, String port, DeviceCredentials deviceCredentials) { public static String getHttpPublishCommand(String protocol, String host, String port, DeviceCredentials deviceCredentials) {
return String.format("curl -v -X POST %s://%s%s/api/v1/%s/telemetry --header Content-Type:application/json --data " + JSON_EXAMPLE_PAYLOAD, return String.format("curl -v -X POST %s://%s%s/api/v1/%s/telemetry --header Content-Type:application/json --data " + JSON_EXAMPLE_PAYLOAD,
@ -71,7 +78,7 @@ public class DeviceConnectivityUtil {
command.append(" -u \"").append(credentials.getUserName()).append("\""); command.append(" -u \"").append(credentials.getUserName()).append("\"");
} }
if (credentials.getPassword() != null) { if (credentials.getPassword() != null) {
command.append(" -P \"").append(credentials.getPassword()).append("\"");; command.append(" -P \"").append(credentials.getPassword()).append("\"");
} }
} else { } else {
return null; return null;
@ -90,17 +97,19 @@ public class DeviceConnectivityUtil {
gatewayVolumePathPrefix = "%HOMEPATH%/tb-gateway"; gatewayVolumePathPrefix = "%HOMEPATH%/tb-gateway";
} }
String gatewayContainerName = "tbGateway" + StringUtils.capitalize(host.replace(".", "")); String gatewayContainerName = "tbGateway" + StringUtils.capitalize(host.replaceAll("[^A-Za-z0-9]", ""));
StringBuilder command = new StringBuilder(GATEWAY_DOCKER_RUN); StringBuilder command = new StringBuilder(GATEWAY_DOCKER_RUN);
command.append("-v {gatewayVolumePathPrefix}/logs:/thingsboard_gateway/logs".replace("{gatewayVolumePathPrefix}", gatewayVolumePathPrefix)); command.append("-v {gatewayVolumePathPrefix}/logs:/thingsboard_gateway/logs ".replace("{gatewayVolumePathPrefix}", gatewayVolumePathPrefix));
command.append(" -v {gatewayVolumePathPrefix}/extensions:/thingsboard_gateway/extensions".replace("{gatewayVolumePathPrefix}", gatewayVolumePathPrefix)); command.append("-v {gatewayVolumePathPrefix}/extensions:/thingsboard_gateway/extensions ".replace("{gatewayVolumePathPrefix}", gatewayVolumePathPrefix));
command.append(" -v {gatewayVolumePathPrefix}/config:/thingsboard_gateway/config".replace("{gatewayVolumePathPrefix}", gatewayVolumePathPrefix)); command.append("-v {gatewayVolumePathPrefix}/config:/thingsboard_gateway/config ".replace("{gatewayVolumePathPrefix}", gatewayVolumePathPrefix));
command.append(" --name ").append(gatewayContainerName); command.append(isLocalhost(host) ? NETWORK_HOST_PARAM : "");
command.append(" -e host=").append(host); command.append("-p 5000:5000 ");
command.append(" -e port=").append(port); command.append("--name ").append(gatewayContainerName).append(" ");
command.append("-e host=").append(host).append(" ");
switch(deviceCredentials.getCredentialsType()) { command.append("-e port=").append(port);
switch (deviceCredentials.getCredentialsType()) {
case ACCESS_TOKEN: case ACCESS_TOKEN:
command.append(" -e accessToken=").append(deviceCredentials.getCredentialsId()); command.append(" -e accessToken=").append(deviceCredentials.getCredentialsId());
break; break;
@ -139,7 +148,7 @@ public class DeviceConnectivityUtil {
} }
StringBuilder mqttDockerCommand = new StringBuilder(); StringBuilder mqttDockerCommand = new StringBuilder();
mqttDockerCommand.append(DOCKER_RUN).append(LOCAL_HOSTS.contains(host) ? "--network=host ":"").append(MQTT_IMAGE); mqttDockerCommand.append(DOCKER_RUN).append(isLocalhost(host) ? NETWORK_HOST_PARAM : "").append(MQTT_IMAGE);
if (MQTTS.equals(protocol)) { if (MQTTS.equals(protocol)) {
mqttDockerCommand.append("/bin/sh -c \"") mqttDockerCommand.append("/bin/sh -c \"")
@ -171,6 +180,40 @@ public class DeviceConnectivityUtil {
public static String getDockerCoapPublishCommand(String protocol, String host, String port, DeviceCredentials deviceCredentials) { public static String getDockerCoapPublishCommand(String protocol, String host, String port, DeviceCredentials deviceCredentials) {
String coapCommand = getCoapPublishCommand(protocol, host, port, deviceCredentials); String coapCommand = getCoapPublishCommand(protocol, host, port, deviceCredentials);
return coapCommand != null ? String.format("%s%s%s", DOCKER_RUN + (LOCAL_HOSTS.contains(host) ? "--network=host ":""), COAP_IMAGE, coapCommand) : null; return coapCommand != null ? String.format("%s%s%s", DOCKER_RUN + (isLocalhost(host) ? NETWORK_HOST_PARAM : ""), COAP_IMAGE, coapCommand) : null;
}
public static String getHost(String baseUrl, DeviceConnectivityInfo properties, String protocol) throws URISyntaxException {
String initialHost = properties.getHost().isEmpty() ? baseUrl : properties.getHost();
InetAddress inetAddress;
String host = null;
if (VALID_URL_PATTERN.matcher(initialHost).matches()) {
host = new URI(initialHost).getHost();
}
if (host == null) {
host = initialHost;
}
try {
host = host.replaceAll("^https?://", "");
inetAddress = InetAddress.getByName(host);
} catch (UnknownHostException e) {
return host;
}
if (inetAddress instanceof Inet6Address) {
host = host.replaceAll("[\\[\\]]", "");
if (!MQTT.equals(protocol) && !MQTTS.equals(protocol)) {
host = "[" + host + "]";
}
}
return host;
}
private static boolean isLocalhost(String host) {
try {
InetAddress inetAddress = InetAddress.getByName(host);
return inetAddress.isLoopbackAddress();
} catch (UnknownHostException e) {
return false;
}
} }
} }

Loading…
Cancel
Save