Browse Source
Conflicts: ui-ngx/src/app/modules/home/components/widget/widget-components.module.tspull/8337/head
526 changed files with 12999 additions and 3269 deletions
File diff suppressed because one or more lines are too long
@ -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; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,194 @@ |
|||||
|
/** |
||||
|
* 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.notification.channels; |
||||
|
|
||||
|
import com.fasterxml.jackson.annotation.JsonInclude; |
||||
|
import com.fasterxml.jackson.annotation.JsonProperty; |
||||
|
import com.fasterxml.jackson.databind.node.ObjectNode; |
||||
|
import com.google.common.base.Strings; |
||||
|
import lombok.Data; |
||||
|
import lombok.RequiredArgsConstructor; |
||||
|
import lombok.Setter; |
||||
|
import org.apache.commons.codec.binary.Base64; |
||||
|
import org.apache.commons.lang3.StringUtils; |
||||
|
import org.springframework.boot.web.client.RestTemplateBuilder; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.springframework.web.client.RestTemplate; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod; |
||||
|
import org.thingsboard.server.common.data.notification.info.NotificationInfo; |
||||
|
import org.thingsboard.server.common.data.notification.targets.MicrosoftTeamsNotificationTargetConfig; |
||||
|
import org.thingsboard.server.common.data.notification.template.MicrosoftTeamsDeliveryMethodNotificationTemplate; |
||||
|
import org.thingsboard.server.common.data.notification.template.MicrosoftTeamsDeliveryMethodNotificationTemplate.Button.LinkType; |
||||
|
import org.thingsboard.server.service.notification.NotificationProcessingContext; |
||||
|
import org.thingsboard.server.service.security.system.SystemSecurityService; |
||||
|
|
||||
|
import java.time.Duration; |
||||
|
import java.time.temporal.ChronoUnit; |
||||
|
import java.util.List; |
||||
|
import java.util.Optional; |
||||
|
|
||||
|
@Component |
||||
|
@RequiredArgsConstructor |
||||
|
public class MicrosoftTeamsNotificationChannel implements NotificationChannel<MicrosoftTeamsNotificationTargetConfig, MicrosoftTeamsDeliveryMethodNotificationTemplate> { |
||||
|
|
||||
|
private final SystemSecurityService systemSecurityService; |
||||
|
|
||||
|
@Setter |
||||
|
private RestTemplate restTemplate = new RestTemplateBuilder() |
||||
|
.setConnectTimeout(Duration.of(15, ChronoUnit.SECONDS)) |
||||
|
.setReadTimeout(Duration.of(15, ChronoUnit.SECONDS)) |
||||
|
.build(); |
||||
|
|
||||
|
@Override |
||||
|
public void sendNotification(MicrosoftTeamsNotificationTargetConfig targetConfig, MicrosoftTeamsDeliveryMethodNotificationTemplate processedTemplate, NotificationProcessingContext ctx) throws Exception { |
||||
|
Message message = new Message(); |
||||
|
message.setThemeColor(Strings.emptyToNull(processedTemplate.getThemeColor())); |
||||
|
if (StringUtils.isEmpty(processedTemplate.getSubject())) { |
||||
|
message.setText(processedTemplate.getBody()); |
||||
|
} else { |
||||
|
message.setSummary(processedTemplate.getSubject()); |
||||
|
Message.Section section = new Message.Section(); |
||||
|
section.setActivityTitle(processedTemplate.getSubject()); |
||||
|
section.setActivitySubtitle(processedTemplate.getBody()); |
||||
|
message.setSections(List.of(section)); |
||||
|
} |
||||
|
var button = processedTemplate.getButton(); |
||||
|
if (button != null && button.isEnabled()) { |
||||
|
String uri; |
||||
|
if (button.getLinkType() == LinkType.DASHBOARD) { |
||||
|
String state = null; |
||||
|
if (button.isSetEntityIdInState() || StringUtils.isNotEmpty(button.getDashboardState())) { |
||||
|
ObjectNode stateObject = JacksonUtil.newObjectNode(); |
||||
|
if (button.isSetEntityIdInState()) { |
||||
|
stateObject.putObject("params") |
||||
|
.set("entityId", Optional.ofNullable(ctx.getRequest().getInfo()) |
||||
|
.map(NotificationInfo::getStateEntityId) |
||||
|
.map(JacksonUtil::valueToTree) |
||||
|
.orElse(null)); |
||||
|
} else { |
||||
|
stateObject.putObject("params"); |
||||
|
} |
||||
|
if (StringUtils.isNotEmpty(button.getDashboardState())) { |
||||
|
stateObject.put("id", button.getDashboardState()); |
||||
|
} |
||||
|
state = Base64.encodeBase64String(JacksonUtil.OBJECT_MAPPER.writeValueAsBytes(List.of(stateObject))); |
||||
|
} |
||||
|
String baseUrl = systemSecurityService.getBaseUrl(ctx.getTenantId(), null, null); |
||||
|
if (StringUtils.isEmpty(baseUrl)) { |
||||
|
throw new IllegalStateException("Failed to determine base url to construct dashboard link"); |
||||
|
} |
||||
|
uri = baseUrl + "/dashboards/" + button.getDashboardId(); |
||||
|
if (state != null) { |
||||
|
uri += "?state=" + state; |
||||
|
} |
||||
|
} else { |
||||
|
uri = button.getLink(); |
||||
|
} |
||||
|
if (StringUtils.isNotBlank(uri) && button.getText() != null) { |
||||
|
Message.ActionCard actionCard = new Message.ActionCard(); |
||||
|
actionCard.setType("OpenUri"); |
||||
|
actionCard.setName(button.getText()); |
||||
|
var target = new Message.ActionCard.Target("default", uri); |
||||
|
actionCard.setTargets(List.of(target)); |
||||
|
message.setPotentialAction(List.of(actionCard)); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
restTemplate.postForEntity(targetConfig.getWebhookUrl(), message, String.class); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void check(TenantId tenantId) throws Exception { |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public NotificationDeliveryMethod getDeliveryMethod() { |
||||
|
return NotificationDeliveryMethod.MICROSOFT_TEAMS; |
||||
|
} |
||||
|
|
||||
|
@Data |
||||
|
public static class Message { |
||||
|
@JsonProperty("@type") |
||||
|
private final String type = "MessageCard"; |
||||
|
@JsonProperty("@context") |
||||
|
private final String context = "http://schema.org/extensions"; |
||||
|
private String themeColor; |
||||
|
private String summary; |
||||
|
private String text; |
||||
|
private List<Section> sections; |
||||
|
private List<ActionCard> potentialAction; |
||||
|
|
||||
|
@Data |
||||
|
public static class Section { |
||||
|
private String activityTitle; |
||||
|
private String activitySubtitle; |
||||
|
private String activityImage; |
||||
|
private List<Fact> facts; |
||||
|
private boolean markdown; |
||||
|
|
||||
|
@Data |
||||
|
public static class Fact { |
||||
|
private final String name; |
||||
|
private final String value; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Data |
||||
|
@JsonInclude(JsonInclude.Include.NON_NULL) |
||||
|
public static class ActionCard { |
||||
|
@JsonProperty("@type") |
||||
|
private String type; // ActionCard, OpenUri
|
||||
|
private String name; |
||||
|
private List<Input> inputs; // for ActionCard
|
||||
|
private List<Action> actions; // for ActionCard
|
||||
|
private List<Target> targets; |
||||
|
|
||||
|
@Data |
||||
|
public static class Input { |
||||
|
@JsonProperty("@type") |
||||
|
private String type; // TextInput, DateInput, MultichoiceInput
|
||||
|
private String id; |
||||
|
private boolean isMultiple; |
||||
|
private String title; |
||||
|
private boolean isMultiSelect; |
||||
|
|
||||
|
@Data |
||||
|
public static class Choice { |
||||
|
private final String display; |
||||
|
private final String value; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Data |
||||
|
public static class Action { |
||||
|
@JsonProperty("@type") |
||||
|
private final String type; // HttpPOST
|
||||
|
private final String name; |
||||
|
private final String target; // url
|
||||
|
} |
||||
|
|
||||
|
@Data |
||||
|
public static class Target { |
||||
|
private final String os; |
||||
|
private final String uri; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
|
|
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue