Browse Source
# Conflicts: # common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.javapull/8757/head
1092 changed files with 36097 additions and 11711 deletions
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
@ -0,0 +1,78 @@ |
|||
/** |
|||
* 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.actors.shared; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.actors.ActorSystemContext; |
|||
import org.thingsboard.server.actors.TbActor; |
|||
import org.thingsboard.server.actors.TbActorId; |
|||
import org.thingsboard.server.actors.TbStringActorId; |
|||
import org.thingsboard.server.actors.service.ContextAwareActor; |
|||
import org.thingsboard.server.actors.service.ContextBasedCreator; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.msg.TbActorMsg; |
|||
import org.thingsboard.server.common.msg.aware.RuleChainAwareMsg; |
|||
import org.thingsboard.server.common.msg.queue.RuleEngineException; |
|||
|
|||
import java.util.UUID; |
|||
|
|||
@Slf4j |
|||
public class RuleChainErrorActor extends ContextAwareActor { |
|||
|
|||
private final TenantId tenantId; |
|||
private final RuleEngineException error; |
|||
|
|||
private RuleChainErrorActor(ActorSystemContext systemContext, TenantId tenantId, RuleEngineException error) { |
|||
super(systemContext); |
|||
this.tenantId = tenantId; |
|||
this.error = error; |
|||
} |
|||
|
|||
@Override |
|||
protected boolean doProcess(TbActorMsg msg) { |
|||
if (msg instanceof RuleChainAwareMsg) { |
|||
log.debug("[{}] Reply with {} for message {}", tenantId, error.getMessage(), msg); |
|||
var rcMsg = (RuleChainAwareMsg) msg; |
|||
rcMsg.getMsg().getCallback().onFailure(error); |
|||
return true; |
|||
} else { |
|||
return false; |
|||
} |
|||
} |
|||
|
|||
public static class ActorCreator extends ContextBasedCreator { |
|||
|
|||
private final TenantId tenantId; |
|||
private final RuleEngineException error; |
|||
|
|||
public ActorCreator(ActorSystemContext context, TenantId tenantId, RuleEngineException error) { |
|||
super(context); |
|||
this.tenantId = tenantId; |
|||
this.error = error; |
|||
} |
|||
|
|||
@Override |
|||
public TbActorId createActorId() { |
|||
return new TbStringActorId(UUID.randomUUID().toString()); |
|||
} |
|||
|
|||
@Override |
|||
public TbActor createActor() { |
|||
return new RuleChainErrorActor(context, tenantId, error); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,107 @@ |
|||
/** |
|||
* 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 io.swagger.annotations.ApiOperation; |
|||
import io.swagger.annotations.ApiParam; |
|||
import io.swagger.annotations.ApiResponse; |
|||
import io.swagger.annotations.ApiResponses; |
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.http.HttpHeaders; |
|||
import org.springframework.http.MediaType; |
|||
import org.springframework.http.ResponseEntity; |
|||
import org.springframework.security.access.prepost.PreAuthorize; |
|||
import org.springframework.web.bind.annotation.PathVariable; |
|||
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.RestController; |
|||
import org.thingsboard.server.common.data.Device; |
|||
import org.thingsboard.server.common.data.exception.ThingsboardException; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.dao.device.DeviceConnectivityService; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.security.permission.Operation; |
|||
import org.thingsboard.server.service.security.system.SystemSecurityService; |
|||
|
|||
import javax.servlet.http.HttpServletRequest; |
|||
import java.io.IOException; |
|||
import java.net.URISyntaxException; |
|||
|
|||
import static org.thingsboard.server.controller.ControllerConstants.DEVICE_ID; |
|||
import static org.thingsboard.server.controller.ControllerConstants.DEVICE_ID_PARAM_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.PROTOCOL; |
|||
import static org.thingsboard.server.controller.ControllerConstants.PROTOCOL_PARAM_DESCRIPTION; |
|||
import static org.thingsboard.server.controller.ControllerConstants.TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH; |
|||
import static org.thingsboard.server.dao.util.DeviceConnectivityUtil.PEM_CERT_FILE_NAME; |
|||
|
|||
@RestController |
|||
@TbCoreComponent |
|||
@RequestMapping("/api") |
|||
@RequiredArgsConstructor |
|||
@Slf4j |
|||
public class DeviceConnectivityController extends BaseController { |
|||
|
|||
private final DeviceConnectivityService deviceConnectivityService; |
|||
private final SystemSecurityService systemSecurityService; |
|||
|
|||
@ApiOperation(value = "Get commands to publish device telemetry (getDevicePublishTelemetryCommands)", |
|||
notes = "Fetch the list of commands to publish device telemetry based on device profile " + |
|||
"If the user has the authority of 'Tenant Administrator', the server checks that the device is owned by the same tenant. " + |
|||
"If the user has the authority of 'Customer User', the server checks that the device is assigned to the same customer. " + |
|||
TENANT_OR_CUSTOMER_AUTHORITY_PARAGRAPH) |
|||
@ApiResponses(value = { |
|||
@ApiResponse(code = 200, message = "OK", |
|||
examples = @io.swagger.annotations.Example( |
|||
value = { |
|||
@io.swagger.annotations.ExampleProperty( |
|||
mediaType = "application/json", |
|||
value = "{\"http\":\"curl -v -X POST http://localhost:8080/api/v1/0ySs4FTOn5WU15XLmal8/telemetry --header Content-Type:application/json --data {temperature:25}\"," + |
|||
"\"mqtt\":\"mosquitto_pub -d -q 1 -h localhost -t v1/devices/me/telemetry -i myClient1 -u myUsername1 -P myPassword -m {temperature:25}\"," + |
|||
"\"coap\":\"coap-client -m POST coap://localhost:5683/api/v1/0ySs4FTOn5WU15XLmal8/telemetry -t json -e {temperature:25}\"}")}))}) |
|||
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") |
|||
@RequestMapping(value = "/device-connectivity/{deviceId}", method = RequestMethod.GET) |
|||
@ResponseBody |
|||
public JsonNode getDevicePublishTelemetryCommands(@ApiParam(value = DEVICE_ID_PARAM_DESCRIPTION) |
|||
@PathVariable(DEVICE_ID) String strDeviceId, HttpServletRequest request) throws ThingsboardException, URISyntaxException { |
|||
checkParameter(DEVICE_ID, strDeviceId); |
|||
DeviceId deviceId = new DeviceId(toUUID(strDeviceId)); |
|||
Device device = checkDeviceId(deviceId, Operation.READ_CREDENTIALS); |
|||
|
|||
String baseUrl = systemSecurityService.getBaseUrl(getTenantId(), getCurrentUser().getCustomerId(), request); |
|||
return deviceConnectivityService.findDevicePublishTelemetryCommands(baseUrl, device); |
|||
} |
|||
|
|||
@ApiOperation(value = "Download server certificate using file path defined in device.connectivity properties (downloadServerCertificate)", notes = "Download server certificate.") |
|||
@RequestMapping(value = "/device-connectivity/{protocol}/certificate/download", method = RequestMethod.GET) |
|||
@ResponseBody |
|||
public ResponseEntity<org.springframework.core.io.Resource> downloadServerCertificate(@ApiParam(value = PROTOCOL_PARAM_DESCRIPTION) |
|||
@PathVariable(PROTOCOL) String protocol) throws ThingsboardException, IOException { |
|||
checkParameter(PROTOCOL, protocol); |
|||
var pemCert = |
|||
checkNotNull(deviceConnectivityService.getPemCertFile(protocol), protocol + " pem cert file is not found!"); |
|||
|
|||
return ResponseEntity.ok() |
|||
.header(HttpHeaders.CONTENT_DISPOSITION, "attachment;filename=" + PEM_CERT_FILE_NAME) |
|||
.header("x-filename", PEM_CERT_FILE_NAME) |
|||
.contentLength(pemCert.contentLength()) |
|||
.contentType(MediaType.APPLICATION_OCTET_STREAM) |
|||
.body(pemCert); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,158 @@ |
|||
/** |
|||
* 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.edge; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.stereotype.Component; |
|||
import org.springframework.transaction.event.TransactionalEventListener; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.cluster.TbClusterService; |
|||
import org.thingsboard.server.common.data.OtaPackageInfo; |
|||
import org.thingsboard.server.common.data.User; |
|||
import org.thingsboard.server.common.data.alarm.AlarmApiCallResult; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventType; |
|||
import org.thingsboard.server.common.data.relation.EntityRelation; |
|||
import org.thingsboard.server.common.data.relation.RelationTypeGroup; |
|||
import org.thingsboard.server.common.data.rule.RuleChain; |
|||
import org.thingsboard.server.common.data.rule.RuleChainType; |
|||
import org.thingsboard.server.common.data.security.Authority; |
|||
import org.thingsboard.server.dao.edge.EdgeSynchronizationManager; |
|||
import org.thingsboard.server.dao.eventsourcing.ActionEntityEvent; |
|||
import org.thingsboard.server.dao.eventsourcing.DeleteEntityEvent; |
|||
import org.thingsboard.server.dao.eventsourcing.RelationActionEvent; |
|||
import org.thingsboard.server.dao.eventsourcing.SaveEntityEvent; |
|||
|
|||
import javax.annotation.PostConstruct; |
|||
|
|||
import static org.thingsboard.server.service.entitiy.DefaultTbNotificationEntityService.edgeTypeByActionType; |
|||
|
|||
|
|||
/** |
|||
* This event listener does not support async event processing because relay on ThreadLocal |
|||
* Another possible approach is to implement a special annotation and a bunch of classes similar to TransactionalApplicationListener |
|||
* This class is the simplest approach to maintain edge synchronization within the single class. |
|||
* <p> |
|||
* For async event publishers, you have to decide whether publish event on creating async task in the same thread where dao method called |
|||
* @Autowired |
|||
* EdgeEventSynchronizationManager edgeSynchronizationManager |
|||
* ... |
|||
* //some async write action make future
|
|||
* if (!edgeSynchronizationManager.isSync()) { |
|||
* future.addCallback(eventPublisher.publishEvent(...)) |
|||
* } |
|||
* */ |
|||
@Component |
|||
@RequiredArgsConstructor |
|||
@Slf4j |
|||
public class EdgeEventSourcingListener { |
|||
|
|||
private final TbClusterService tbClusterService; |
|||
private final EdgeSynchronizationManager edgeSynchronizationManager; |
|||
|
|||
@PostConstruct |
|||
public void init() { |
|||
log.info("EdgeEventSourcingListener initiated"); |
|||
} |
|||
|
|||
@TransactionalEventListener(fallbackExecution = true) |
|||
public void handleEvent(SaveEntityEvent<?> event) { |
|||
if (edgeSynchronizationManager.isSync()) { |
|||
return; |
|||
} |
|||
try { |
|||
if (!isValidEdgeEventEntity(event.getEntity())) { |
|||
return; |
|||
} |
|||
log.trace("[{}] SaveEntityEvent called: {}", event.getTenantId(), event); |
|||
EdgeEventActionType action = Boolean.TRUE.equals(event.getAdded()) ? EdgeEventActionType.ADDED : EdgeEventActionType.UPDATED; |
|||
tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), null, event.getEntityId(), |
|||
null, null, action); |
|||
} catch (Exception e) { |
|||
log.error("[{}] failed to process SaveEntityEvent: {}", event.getTenantId(), event); |
|||
} |
|||
} |
|||
|
|||
@TransactionalEventListener(fallbackExecution = true) |
|||
public void handleEvent(DeleteEntityEvent<?> event) { |
|||
if (edgeSynchronizationManager.isSync()) { |
|||
return; |
|||
} |
|||
try { |
|||
log.trace("[{}] DeleteEntityEvent called: {}", event.getTenantId(), event); |
|||
tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), event.getEdgeId(), event.getEntityId(), |
|||
JacksonUtil.toString(event.getEntity()), null, EdgeEventActionType.DELETED); |
|||
} catch (Exception e) { |
|||
log.error("[{}] failed to process DeleteEntityEvent: {}", event.getTenantId(), event); |
|||
} |
|||
} |
|||
|
|||
@TransactionalEventListener(fallbackExecution = true) |
|||
public void handleEvent(ActionEntityEvent event) { |
|||
if (edgeSynchronizationManager.isSync()) { |
|||
return; |
|||
} |
|||
try { |
|||
log.trace("[{}] ActionEntityEvent called: {}", event.getTenantId(), event); |
|||
tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), event.getEdgeId(), event.getEntityId(), |
|||
event.getBody(), null, edgeTypeByActionType(event.getActionType())); |
|||
} catch (Exception e) { |
|||
log.error("[{}] failed to process ActionEntityEvent: {}", event.getTenantId(), event); |
|||
} |
|||
} |
|||
|
|||
@TransactionalEventListener(fallbackExecution = true) |
|||
public void handleEvent(RelationActionEvent event) { |
|||
if (edgeSynchronizationManager.isSync()) { |
|||
return; |
|||
} |
|||
try { |
|||
EntityRelation relation = event.getRelation(); |
|||
if (relation == null) { |
|||
log.trace("[{}] skipping RelationActionEvent event in case relation is null: {}", event.getTenantId(), event); |
|||
return; |
|||
} |
|||
if (!RelationTypeGroup.COMMON.equals(relation.getTypeGroup())) { |
|||
log.trace("[{}] skipping RelationActionEvent event in case NOT COMMON relation type group: {}", event.getTenantId(), event); |
|||
return; |
|||
} |
|||
log.trace("[{}] RelationActionEvent called: {}", event.getTenantId(), event); |
|||
tbClusterService.sendNotificationMsgToEdge(event.getTenantId(), null, null, |
|||
JacksonUtil.toString(relation), EdgeEventType.RELATION, edgeTypeByActionType(event.getActionType())); |
|||
} catch (Exception e) { |
|||
log.error("[{}] failed to process RelationActionEvent: {}", event.getTenantId(), event); |
|||
} |
|||
} |
|||
|
|||
private boolean isValidEdgeEventEntity(Object entity) { |
|||
if (entity instanceof OtaPackageInfo) { |
|||
OtaPackageInfo otaPackageInfo = (OtaPackageInfo) entity; |
|||
return otaPackageInfo.hasUrl() || otaPackageInfo.isHasData(); |
|||
} else if (entity instanceof RuleChain) { |
|||
RuleChain ruleChain = (RuleChain) entity; |
|||
return RuleChainType.EDGE.equals(ruleChain.getType()); |
|||
} else if (entity instanceof User) { |
|||
User user = (User) entity; |
|||
return !Authority.SYS_ADMIN.equals(user.getAuthority()); |
|||
} else if (entity instanceof AlarmApiCallResult) { |
|||
AlarmApiCallResult alarmApiCallResult = (AlarmApiCallResult) entity; |
|||
return alarmApiCallResult.isModified(); |
|||
} |
|||
// Default: If the entity doesn't match any of the conditions, consider it as valid.
|
|||
return true; |
|||
} |
|||
} |
|||
@ -0,0 +1,67 @@ |
|||
/** |
|||
* 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.edge.rpc.constructor; |
|||
|
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.Tenant; |
|||
import org.thingsboard.server.gen.edge.v1.TenantUpdateMsg; |
|||
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
|
|||
@Component |
|||
@TbCoreComponent |
|||
public class TenantMsgConstructor { |
|||
|
|||
public TenantUpdateMsg constructTenantUpdateMsg(UpdateMsgType msgType, Tenant tenant) { |
|||
TenantUpdateMsg.Builder builder = TenantUpdateMsg.newBuilder() |
|||
.setMsgType(msgType) |
|||
.setIdMSB(tenant.getId().getId().getMostSignificantBits()) |
|||
.setIdLSB(tenant.getId().getId().getLeastSignificantBits()) |
|||
.setTitle(tenant.getTitle()) |
|||
.setProfileIdMSB(tenant.getTenantProfileId().getId().getMostSignificantBits()) |
|||
.setProfileIdLSB(tenant.getTenantProfileId().getId().getLeastSignificantBits()) |
|||
.setRegion(tenant.getRegion()); |
|||
if (tenant.getCountry() != null) { |
|||
builder.setCountry(tenant.getCountry()); |
|||
} |
|||
if (tenant.getState() != null) { |
|||
builder.setState(tenant.getState()); |
|||
} |
|||
if (tenant.getCity() != null) { |
|||
builder.setCity(tenant.getCity()); |
|||
} |
|||
if (tenant.getAddress() != null) { |
|||
builder.setAddress(tenant.getAddress()); |
|||
} |
|||
if (tenant.getAddress2() != null) { |
|||
builder.setAddress2(tenant.getAddress2()); |
|||
} |
|||
if (tenant.getZip() != null) { |
|||
builder.setZip(tenant.getZip()); |
|||
} |
|||
if (tenant.getPhone() != null) { |
|||
builder.setPhone(tenant.getPhone()); |
|||
} |
|||
if (tenant.getEmail() != null) { |
|||
builder.setEmail(tenant.getEmail()); |
|||
} |
|||
if (tenant.getAdditionalInfo() != null) { |
|||
builder.setAdditionalInfo(JacksonUtil.toString(tenant.getAdditionalInfo())); |
|||
} |
|||
return builder.build(); |
|||
} |
|||
} |
|||
@ -0,0 +1,48 @@ |
|||
/** |
|||
* 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.edge.rpc.constructor; |
|||
|
|||
import com.google.protobuf.ByteString; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.server.common.data.TenantProfile; |
|||
import org.thingsboard.server.gen.edge.v1.TenantProfileUpdateMsg; |
|||
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
|||
import org.thingsboard.server.queue.util.DataDecodingEncodingService; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
|
|||
@Component |
|||
@TbCoreComponent |
|||
public class TenantProfileMsgConstructor { |
|||
|
|||
@Autowired |
|||
private DataDecodingEncodingService dataDecodingEncodingService; |
|||
|
|||
public TenantProfileUpdateMsg constructTenantProfileUpdateMsg(UpdateMsgType msgType, TenantProfile tenantProfile) { |
|||
TenantProfileUpdateMsg.Builder builder = TenantProfileUpdateMsg.newBuilder() |
|||
.setMsgType(msgType) |
|||
.setIdMSB(tenantProfile.getId().getId().getMostSignificantBits()) |
|||
.setIdLSB(tenantProfile.getId().getId().getLeastSignificantBits()) |
|||
.setName(tenantProfile.getName()) |
|||
.setDefault(tenantProfile.isDefault()) |
|||
.setIsolatedRuleChain(tenantProfile.isIsolatedTbRuleEngine()) |
|||
.setProfileDataBytes(ByteString.copyFrom(dataDecodingEncodingService.encode(tenantProfile.getProfileData()))); |
|||
if (tenantProfile.getDescription() != null) { |
|||
builder.setDescription(tenantProfile.getDescription()); |
|||
} |
|||
return builder.build(); |
|||
} |
|||
} |
|||
@ -0,0 +1,51 @@ |
|||
/** |
|||
* 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.edge.rpc.fetch; |
|||
|
|||
import lombok.AllArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.common.data.EdgeUtils; |
|||
import org.thingsboard.server.common.data.Tenant; |
|||
import org.thingsboard.server.common.data.edge.Edge; |
|||
import org.thingsboard.server.common.data.edge.EdgeEvent; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventType; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.page.PageData; |
|||
import org.thingsboard.server.common.data.page.PageLink; |
|||
import org.thingsboard.server.dao.tenant.TenantService; |
|||
|
|||
import java.util.List; |
|||
|
|||
@AllArgsConstructor |
|||
@Slf4j |
|||
public class TenantEdgeEventFetcher extends BasePageableEdgeEventFetcher<Tenant> { |
|||
|
|||
private final TenantService tenantService; |
|||
|
|||
@Override |
|||
PageData<Tenant> fetchPageData(TenantId tenantId, Edge edge, PageLink pageLink) { |
|||
Tenant tenant = tenantService.findTenantById(tenantId); |
|||
// returns PageData object to be in sync with other fetchers
|
|||
return new PageData<>(List.of(tenant), 1, 1, false); |
|||
} |
|||
|
|||
@Override |
|||
EdgeEvent constructEdgeEvent(TenantId tenantId, Edge edge, Tenant entity) { |
|||
return EdgeUtils.constructEdgeEvent(tenantId, edge.getId(), EdgeEventType.TENANT, |
|||
EdgeEventActionType.UPDATED, entity.getId(), null); |
|||
} |
|||
} |
|||
@ -0,0 +1,77 @@ |
|||
/** |
|||
* 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.edge.rpc.processor.asset; |
|||
|
|||
import com.datastax.oss.driver.api.core.uuid.Uuids; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.data.util.Pair; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.StringUtils; |
|||
import org.thingsboard.server.common.data.asset.Asset; |
|||
import org.thingsboard.server.common.data.id.AssetId; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.common.data.id.CustomerId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.gen.edge.v1.AssetUpdateMsg; |
|||
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
|||
|
|||
import java.util.UUID; |
|||
|
|||
@Slf4j |
|||
public abstract class BaseAssetProcessor extends BaseEdgeProcessor { |
|||
|
|||
protected Pair<Boolean, Boolean> saveOrUpdateAsset(TenantId tenantId, AssetId assetId, AssetUpdateMsg assetUpdateMsg, CustomerId customerId) { |
|||
boolean created = false; |
|||
boolean assetNameUpdated = false; |
|||
assetCreationLock.lock(); |
|||
try { |
|||
Asset asset = assetService.findAssetById(tenantId, assetId); |
|||
String assetName = assetUpdateMsg.getName(); |
|||
if (asset == null) { |
|||
created = true; |
|||
asset = new Asset(); |
|||
asset.setTenantId(tenantId); |
|||
asset.setCreatedTime(Uuids.unixTimestamp(assetId.getId())); |
|||
} |
|||
Asset assetByName = assetService.findAssetByTenantIdAndName(tenantId, assetName); |
|||
if (assetByName != null && !assetByName.getId().equals(assetId)) { |
|||
assetName = assetName + "_" + StringUtils.randomAlphanumeric(15); |
|||
log.warn("Asset with name {} already exists. Renaming asset name to {}", |
|||
assetUpdateMsg.getName(), assetName); |
|||
assetNameUpdated = true; |
|||
} |
|||
asset.setName(assetName); |
|||
asset.setType(assetUpdateMsg.getType()); |
|||
asset.setLabel(assetUpdateMsg.hasLabel() ? assetUpdateMsg.getLabel() : null); |
|||
asset.setAdditionalInfo(assetUpdateMsg.hasAdditionalInfo() |
|||
? JacksonUtil.toJsonNode(assetUpdateMsg.getAdditionalInfo()) : null); |
|||
|
|||
UUID assetProfileUUID = safeGetUUID(assetUpdateMsg.getAssetProfileIdMSB(), assetUpdateMsg.getAssetProfileIdLSB()); |
|||
asset.setAssetProfileId(assetProfileUUID != null ? new AssetProfileId(assetProfileUUID) : null); |
|||
|
|||
asset.setCustomerId(customerId); |
|||
|
|||
assetValidator.validate(asset, Asset::getTenantId); |
|||
if (created) { |
|||
asset.setId(assetId); |
|||
} |
|||
assetService.saveAsset(asset, false); |
|||
} finally { |
|||
assetCreationLock.unlock(); |
|||
} |
|||
return Pair.of(created, assetNameUpdated); |
|||
} |
|||
} |
|||
@ -0,0 +1,69 @@ |
|||
/** |
|||
* 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.edge.rpc.processor.asset; |
|||
|
|||
import com.datastax.oss.driver.api.core.uuid.Uuids; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.server.common.data.asset.AssetProfile; |
|||
import org.thingsboard.server.common.data.id.AssetProfileId; |
|||
import org.thingsboard.server.common.data.id.DashboardId; |
|||
import org.thingsboard.server.common.data.id.RuleChainId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.gen.edge.v1.AssetProfileUpdateMsg; |
|||
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
|||
|
|||
import java.nio.charset.StandardCharsets; |
|||
import java.util.UUID; |
|||
|
|||
@Slf4j |
|||
public class BaseAssetProfileProcessor extends BaseEdgeProcessor { |
|||
|
|||
protected boolean saveOrUpdateAssetProfile(TenantId tenantId, AssetProfileId assetProfileId, AssetProfileUpdateMsg assetProfileUpdateMsg) { |
|||
boolean created = false; |
|||
assetCreationLock.lock(); |
|||
try { |
|||
AssetProfile assetProfile = assetProfileService.findAssetProfileById(tenantId, assetProfileId); |
|||
String assetProfileName = assetProfileUpdateMsg.getName(); |
|||
if (assetProfile == null) { |
|||
created = true; |
|||
assetProfile = new AssetProfile(); |
|||
assetProfile.setTenantId(tenantId); |
|||
assetProfile.setCreatedTime(Uuids.unixTimestamp(assetProfileId.getId())); |
|||
} |
|||
assetProfile.setName(assetProfileName); |
|||
assetProfile.setDefault(assetProfileUpdateMsg.getDefault()); |
|||
assetProfile.setDefaultQueueName(assetProfileUpdateMsg.hasDefaultQueueName() ? assetProfileUpdateMsg.getDefaultQueueName() : null); |
|||
assetProfile.setDescription(assetProfileUpdateMsg.hasDescription() ? assetProfileUpdateMsg.getDescription() : null); |
|||
assetProfile.setImage(assetProfileUpdateMsg.hasImage() |
|||
? new String(assetProfileUpdateMsg.getImage().toByteArray(), StandardCharsets.UTF_8) : null); |
|||
|
|||
UUID defaultRuleChainUUID = safeGetUUID(assetProfileUpdateMsg.getDefaultRuleChainIdMSB(), assetProfileUpdateMsg.getDefaultRuleChainIdLSB()); |
|||
assetProfile.setDefaultRuleChainId(defaultRuleChainUUID != null ? new RuleChainId(defaultRuleChainUUID) : null); |
|||
|
|||
UUID defaultDashboardUUID = safeGetUUID(assetProfileUpdateMsg.getDefaultDashboardIdMSB(), assetProfileUpdateMsg.getDefaultDashboardIdLSB()); |
|||
assetProfile.setDefaultDashboardId(defaultDashboardUUID != null ? new DashboardId(defaultDashboardUUID) : null); |
|||
|
|||
assetProfileValidator.validate(assetProfile, AssetProfile::getTenantId); |
|||
if (created) { |
|||
assetProfile.setId(assetProfileId); |
|||
} |
|||
assetProfileService.saveAssetProfile(assetProfile, false); |
|||
} finally { |
|||
assetCreationLock.unlock(); |
|||
} |
|||
return created; |
|||
} |
|||
} |
|||
@ -0,0 +1,77 @@ |
|||
/** |
|||
* 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.edge.rpc.processor.dashboard; |
|||
|
|||
import com.datastax.oss.driver.api.core.uuid.Uuids; |
|||
import com.fasterxml.jackson.core.type.TypeReference; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.Dashboard; |
|||
import org.thingsboard.server.common.data.ShortCustomerInfo; |
|||
import org.thingsboard.server.common.data.id.CustomerId; |
|||
import org.thingsboard.server.common.data.id.DashboardId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.gen.edge.v1.DashboardUpdateMsg; |
|||
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
|||
|
|||
import java.util.Set; |
|||
|
|||
@Slf4j |
|||
public abstract class BaseDashboardProcessor extends BaseEdgeProcessor { |
|||
|
|||
protected boolean saveOrUpdateDashboard(TenantId tenantId, DashboardId dashboardId, DashboardUpdateMsg dashboardUpdateMsg, CustomerId customerId) { |
|||
boolean created = false; |
|||
Dashboard dashboard = dashboardService.findDashboardById(tenantId, dashboardId); |
|||
if (dashboard == null) { |
|||
created = true; |
|||
dashboard = new Dashboard(); |
|||
dashboard.setTenantId(tenantId); |
|||
dashboard.setCreatedTime(Uuids.unixTimestamp(dashboardId.getId())); |
|||
} |
|||
dashboard.setTitle(dashboardUpdateMsg.getTitle()); |
|||
dashboard.setConfiguration(JacksonUtil.toJsonNode(dashboardUpdateMsg.getConfiguration())); |
|||
Set<ShortCustomerInfo> assignedCustomers = null; |
|||
if (dashboardUpdateMsg.hasAssignedCustomers()) { |
|||
assignedCustomers = JacksonUtil.fromString(dashboardUpdateMsg.getAssignedCustomers(), new TypeReference<>() { |
|||
}); |
|||
dashboard.setAssignedCustomers(assignedCustomers); |
|||
} |
|||
|
|||
dashboardValidator.validate(dashboard, Dashboard::getTenantId); |
|||
if (created) { |
|||
dashboard.setId(dashboardId); |
|||
} |
|||
Dashboard savedDashboard = dashboardService.saveDashboard(dashboard, false); |
|||
if (assignedCustomers != null && !assignedCustomers.isEmpty()) { |
|||
for (ShortCustomerInfo assignedCustomer : assignedCustomers) { |
|||
if (assignedCustomer.getCustomerId().equals(customerId)) { |
|||
dashboardService.assignDashboardToCustomer(tenantId, dashboardId, assignedCustomer.getCustomerId()); |
|||
} |
|||
} |
|||
} else { |
|||
unassignCustomersFromDashboard(tenantId, savedDashboard); |
|||
} |
|||
return created; |
|||
} |
|||
|
|||
private void unassignCustomersFromDashboard(TenantId tenantId, Dashboard dashboard) { |
|||
if (dashboard.getAssignedCustomers() != null && !dashboard.getAssignedCustomers().isEmpty()) { |
|||
for (ShortCustomerInfo assignedCustomer : dashboard.getAssignedCustomers()) { |
|||
dashboardService.unassignDashboardFromCustomer(tenantId, dashboard.getId(), assignedCustomer.getCustomerId()); |
|||
} |
|||
} |
|||
} |
|||
} |
|||
@ -0,0 +1,102 @@ |
|||
/** |
|||
* 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.edge.rpc.processor.device; |
|||
|
|||
import com.datastax.oss.driver.api.core.uuid.Uuids; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.beans.factory.annotation.Autowired; |
|||
import org.thingsboard.server.common.data.DeviceProfile; |
|||
import org.thingsboard.server.common.data.DeviceProfileProvisionType; |
|||
import org.thingsboard.server.common.data.DeviceProfileType; |
|||
import org.thingsboard.server.common.data.DeviceTransportType; |
|||
import org.thingsboard.server.common.data.StringUtils; |
|||
import org.thingsboard.server.common.data.device.profile.DeviceProfileData; |
|||
import org.thingsboard.server.common.data.id.DashboardId; |
|||
import org.thingsboard.server.common.data.id.DeviceProfileId; |
|||
import org.thingsboard.server.common.data.id.OtaPackageId; |
|||
import org.thingsboard.server.common.data.id.RuleChainId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.gen.edge.v1.DeviceProfileUpdateMsg; |
|||
import org.thingsboard.server.queue.util.DataDecodingEncodingService; |
|||
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
|||
|
|||
import java.nio.charset.StandardCharsets; |
|||
import java.util.Optional; |
|||
import java.util.UUID; |
|||
|
|||
@Slf4j |
|||
public class BaseDeviceProfileProcessor extends BaseEdgeProcessor { |
|||
|
|||
@Autowired |
|||
private DataDecodingEncodingService dataDecodingEncodingService; |
|||
|
|||
protected boolean saveOrUpdateDeviceProfile(TenantId tenantId, DeviceProfileId deviceProfileId, DeviceProfileUpdateMsg deviceProfileUpdateMsg) { |
|||
boolean created = false; |
|||
deviceCreationLock.lock(); |
|||
try { |
|||
DeviceProfile deviceProfile = deviceProfileService.findDeviceProfileById(tenantId, deviceProfileId); |
|||
if (deviceProfile == null) { |
|||
created = true; |
|||
deviceProfile = new DeviceProfile(); |
|||
deviceProfile.setTenantId(tenantId); |
|||
deviceProfile.setCreatedTime(Uuids.unixTimestamp(deviceProfileId.getId())); |
|||
} |
|||
deviceProfile.setName(deviceProfileUpdateMsg.getName()); |
|||
deviceProfile.setDescription(deviceProfileUpdateMsg.hasDescription() ? deviceProfileUpdateMsg.getDescription() : null); |
|||
deviceProfile.setDefault(deviceProfileUpdateMsg.getDefault()); |
|||
deviceProfile.setType(DeviceProfileType.valueOf(deviceProfileUpdateMsg.getType())); |
|||
deviceProfile.setTransportType(deviceProfileUpdateMsg.hasTransportType() |
|||
? DeviceTransportType.valueOf(deviceProfileUpdateMsg.getTransportType()) : DeviceTransportType.DEFAULT); |
|||
deviceProfile.setImage(deviceProfileUpdateMsg.hasImage() |
|||
? new String(deviceProfileUpdateMsg.getImage().toByteArray(), StandardCharsets.UTF_8) : null); |
|||
deviceProfile.setProvisionType(deviceProfileUpdateMsg.hasProvisionType() |
|||
? DeviceProfileProvisionType.valueOf(deviceProfileUpdateMsg.getProvisionType()) : DeviceProfileProvisionType.DISABLED); |
|||
deviceProfile.setProvisionDeviceKey(deviceProfileUpdateMsg.hasProvisionDeviceKey() |
|||
? deviceProfileUpdateMsg.getProvisionDeviceKey() : null); |
|||
deviceProfile.setDefaultQueueName(deviceProfileUpdateMsg.getDefaultQueueName()); |
|||
|
|||
Optional<DeviceProfileData> profileDataOpt = |
|||
dataDecodingEncodingService.decode(deviceProfileUpdateMsg.getProfileDataBytes().toByteArray()); |
|||
deviceProfile.setProfileData(profileDataOpt.orElse(null)); |
|||
|
|||
UUID defaultRuleChainUUID = safeGetUUID(deviceProfileUpdateMsg.getDefaultRuleChainIdMSB(), deviceProfileUpdateMsg.getDefaultRuleChainIdLSB()); |
|||
deviceProfile.setDefaultRuleChainId(defaultRuleChainUUID != null ? new RuleChainId(defaultRuleChainUUID) : null); |
|||
|
|||
UUID defaultDashboardUUID = safeGetUUID(deviceProfileUpdateMsg.getDefaultDashboardIdMSB(), deviceProfileUpdateMsg.getDefaultDashboardIdLSB()); |
|||
deviceProfile.setDefaultDashboardId(defaultDashboardUUID != null ? new DashboardId(defaultDashboardUUID) : null); |
|||
|
|||
String defaultQueueName = StringUtils.isNotBlank(deviceProfileUpdateMsg.getDefaultQueueName()) |
|||
? deviceProfileUpdateMsg.getDefaultQueueName() : null; |
|||
deviceProfile.setDefaultQueueName(defaultQueueName); |
|||
|
|||
UUID firmwareUUID = safeGetUUID(deviceProfileUpdateMsg.getFirmwareIdMSB(), deviceProfileUpdateMsg.getFirmwareIdLSB()); |
|||
deviceProfile.setFirmwareId(firmwareUUID != null ? new OtaPackageId(firmwareUUID) : null); |
|||
|
|||
UUID softwareUUID = safeGetUUID(deviceProfileUpdateMsg.getSoftwareIdMSB(), deviceProfileUpdateMsg.getSoftwareIdLSB()); |
|||
deviceProfile.setSoftwareId(softwareUUID != null ? new OtaPackageId(softwareUUID) : null); |
|||
|
|||
|
|||
deviceProfileValidator.validate(deviceProfile, DeviceProfile::getTenantId); |
|||
if (created) { |
|||
deviceProfile.setId(deviceProfileId); |
|||
} |
|||
deviceProfileService.saveDeviceProfile(deviceProfile, false); |
|||
} finally { |
|||
deviceCreationLock.unlock(); |
|||
} |
|||
return created; |
|||
} |
|||
} |
|||
@ -0,0 +1,76 @@ |
|||
/** |
|||
* 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.edge.rpc.processor.entityview; |
|||
|
|||
import com.datastax.oss.driver.api.core.uuid.Uuids; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.data.util.Pair; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.EntityView; |
|||
import org.thingsboard.server.common.data.StringUtils; |
|||
import org.thingsboard.server.common.data.id.AssetId; |
|||
import org.thingsboard.server.common.data.id.CustomerId; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.EntityViewId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.gen.edge.v1.EdgeEntityType; |
|||
import org.thingsboard.server.gen.edge.v1.EntityViewUpdateMsg; |
|||
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
|||
|
|||
import java.util.UUID; |
|||
|
|||
@Slf4j |
|||
public abstract class BaseEntityViewProcessor extends BaseEdgeProcessor { |
|||
|
|||
protected Pair<Boolean, Boolean> saveOrUpdateEntityView(TenantId tenantId, EntityViewId entityViewId, EntityViewUpdateMsg entityViewUpdateMsg, CustomerId customerId) { |
|||
boolean created = false; |
|||
boolean entityViewNameUpdated = false; |
|||
EntityView entityView = entityViewService.findEntityViewById(tenantId, entityViewId); |
|||
String entityViewName = entityViewUpdateMsg.getName(); |
|||
if (entityView == null) { |
|||
created = true; |
|||
entityView = new EntityView(); |
|||
entityView.setTenantId(tenantId); |
|||
entityView.setCreatedTime(Uuids.unixTimestamp(entityViewId.getId())); |
|||
} |
|||
EntityView entityViewByName = entityViewService.findEntityViewByTenantIdAndName(tenantId, entityViewName); |
|||
if (entityViewByName != null && !entityViewByName.getId().equals(entityViewId)) { |
|||
entityViewName = entityViewName + "_" + StringUtils.randomAlphanumeric(15); |
|||
log.warn("Entity view with name {} already exists. Renaming entity view name to {}", |
|||
entityViewUpdateMsg.getName(), entityViewName); |
|||
entityViewNameUpdated = true; |
|||
} |
|||
entityView.setName(entityViewName); |
|||
entityView.setType(entityViewUpdateMsg.getType()); |
|||
entityView.setCustomerId(customerId); |
|||
entityView.setAdditionalInfo(entityViewUpdateMsg.hasAdditionalInfo() ? |
|||
JacksonUtil.toJsonNode(entityViewUpdateMsg.getAdditionalInfo()) : null); |
|||
|
|||
UUID entityIdUUID = safeGetUUID(entityViewUpdateMsg.getEntityIdMSB(), entityViewUpdateMsg.getEntityIdLSB()); |
|||
if (EdgeEntityType.DEVICE.equals(entityViewUpdateMsg.getEntityType())) { |
|||
entityView.setEntityId(entityIdUUID != null ? new DeviceId(entityIdUUID) : null); |
|||
} else if (EdgeEntityType.ASSET.equals(entityViewUpdateMsg.getEntityType())) { |
|||
entityView.setEntityId(entityIdUUID != null ? new AssetId(entityIdUUID) : null); |
|||
} |
|||
|
|||
entityViewValidator.validate(entityView, EntityView::getTenantId); |
|||
if (created) { |
|||
entityView.setId(entityViewId); |
|||
} |
|||
entityViewService.saveEntityView(entityView, false); |
|||
return Pair.of(created, entityViewNameUpdated); |
|||
} |
|||
} |
|||
@ -0,0 +1,57 @@ |
|||
/** |
|||
* 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.edge.rpc.processor.tenant; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.server.common.data.EdgeUtils; |
|||
import org.thingsboard.server.common.data.Tenant; |
|||
import org.thingsboard.server.common.data.TenantProfile; |
|||
import org.thingsboard.server.common.data.edge.EdgeEvent; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; |
|||
import org.thingsboard.server.gen.edge.v1.TenantProfileUpdateMsg; |
|||
import org.thingsboard.server.gen.edge.v1.TenantUpdateMsg; |
|||
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
|||
|
|||
@Component |
|||
@Slf4j |
|||
@TbCoreComponent |
|||
public class TenantEdgeProcessor extends BaseEdgeProcessor { |
|||
|
|||
public DownlinkMsg convertTenantEventToDownlink(EdgeEvent edgeEvent) { |
|||
TenantId tenantId = new TenantId(edgeEvent.getEntityId()); |
|||
DownlinkMsg downlinkMsg = null; |
|||
if (EdgeEventActionType.UPDATED.equals(edgeEvent.getAction())) { |
|||
Tenant tenant = tenantService.findTenantById(tenantId); |
|||
if (tenant != null) { |
|||
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); |
|||
TenantUpdateMsg tenantUpdateMsg = tenantMsgConstructor.constructTenantUpdateMsg(msgType, tenant); |
|||
TenantProfile tenantProfile = tenantProfileService.findTenantProfileById(tenantId, tenant.getTenantProfileId()); |
|||
TenantProfileUpdateMsg tenantProfileUpdateMsg = tenantProfileMsgConstructor.constructTenantProfileUpdateMsg(msgType, tenantProfile); |
|||
downlinkMsg = DownlinkMsg.newBuilder() |
|||
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|||
.addTenantUpdateMsg(tenantUpdateMsg) |
|||
.addTenantProfileUpdateMsg(tenantProfileUpdateMsg) |
|||
.build(); |
|||
} |
|||
} |
|||
return downlinkMsg; |
|||
} |
|||
} |
|||
@ -0,0 +1,53 @@ |
|||
/** |
|||
* 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.edge.rpc.processor.tenant; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.server.common.data.EdgeUtils; |
|||
import org.thingsboard.server.common.data.TenantProfile; |
|||
import org.thingsboard.server.common.data.edge.EdgeEvent; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|||
import org.thingsboard.server.common.data.id.TenantProfileId; |
|||
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; |
|||
import org.thingsboard.server.gen.edge.v1.TenantProfileUpdateMsg; |
|||
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
|||
|
|||
@Component |
|||
@Slf4j |
|||
@TbCoreComponent |
|||
public class TenantProfileEdgeProcessor extends BaseEdgeProcessor { |
|||
|
|||
public DownlinkMsg convertTenantProfileEventToDownlink(EdgeEvent edgeEvent) { |
|||
TenantProfileId tenantProfileId = new TenantProfileId(edgeEvent.getEntityId()); |
|||
DownlinkMsg downlinkMsg = null; |
|||
if (EdgeEventActionType.UPDATED.equals(edgeEvent.getAction())) { |
|||
TenantProfile tenantProfile = tenantProfileService.findTenantProfileById(edgeEvent.getTenantId(), tenantProfileId); |
|||
if (tenantProfile != null) { |
|||
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); |
|||
TenantProfileUpdateMsg tenantProfileUpdateMsg = |
|||
tenantProfileMsgConstructor.constructTenantProfileUpdateMsg(msgType, tenantProfile); |
|||
downlinkMsg = DownlinkMsg.newBuilder() |
|||
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|||
.addTenantProfileUpdateMsg(tenantProfileUpdateMsg) |
|||
.build(); |
|||
} |
|||
} |
|||
return downlinkMsg; |
|||
} |
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue