Browse Source
[3.5] Edge processors refactoring - move common part of edge / server to base classespull/8034/head
committed by
GitHub
32 changed files with 756 additions and 623 deletions
@ -1,195 +0,0 @@ |
|||||
/** |
|
||||
* Copyright © 2016-2022 The Thingsboard Authors |
|
||||
* |
|
||||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|
||||
* you may not use this file except in compliance with the License. |
|
||||
* You may obtain a copy of the License at |
|
||||
* |
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
|
||||
* |
|
||||
* Unless required by applicable law or agreed to in writing, software |
|
||||
* distributed under the License is distributed on an "AS IS" BASIS, |
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|
||||
* See the License for the specific language governing permissions and |
|
||||
* limitations under the License. |
|
||||
*/ |
|
||||
package org.thingsboard.server.service.edge.rpc.processor; |
|
||||
|
|
||||
import com.fasterxml.jackson.core.JsonProcessingException; |
|
||||
import com.google.common.util.concurrent.Futures; |
|
||||
import com.google.common.util.concurrent.ListenableFuture; |
|
||||
import lombok.extern.slf4j.Slf4j; |
|
||||
import org.springframework.stereotype.Component; |
|
||||
import org.thingsboard.common.util.JacksonUtil; |
|
||||
import org.thingsboard.server.common.data.EdgeUtils; |
|
||||
import org.thingsboard.server.common.data.EntityType; |
|
||||
import org.thingsboard.server.common.data.alarm.Alarm; |
|
||||
import org.thingsboard.server.common.data.alarm.AlarmSeverity; |
|
||||
import org.thingsboard.server.common.data.alarm.AlarmStatus; |
|
||||
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.AlarmId; |
|
||||
import org.thingsboard.server.common.data.id.EdgeId; |
|
||||
import org.thingsboard.server.common.data.id.EntityId; |
|
||||
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.gen.edge.v1.AlarmUpdateMsg; |
|
||||
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; |
|
||||
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
|
||||
import org.thingsboard.server.gen.transport.TransportProtos; |
|
||||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|
||||
|
|
||||
import java.util.ArrayList; |
|
||||
import java.util.List; |
|
||||
import java.util.UUID; |
|
||||
|
|
||||
@Component |
|
||||
@Slf4j |
|
||||
@TbCoreComponent |
|
||||
public class AlarmEdgeProcessor extends BaseEdgeProcessor { |
|
||||
|
|
||||
public ListenableFuture<Void> processAlarmFromEdge(TenantId tenantId, AlarmUpdateMsg alarmUpdateMsg) { |
|
||||
log.trace("[{}] processAlarmFromEdge [{}]", tenantId, alarmUpdateMsg); |
|
||||
EntityId originatorId = getAlarmOriginator(tenantId, alarmUpdateMsg.getOriginatorName(), |
|
||||
EntityType.valueOf(alarmUpdateMsg.getOriginatorType())); |
|
||||
if (originatorId == null) { |
|
||||
log.warn("Originator not found for the alarm msg {}", alarmUpdateMsg); |
|
||||
return Futures.immediateFuture(null); |
|
||||
} |
|
||||
try { |
|
||||
Alarm existentAlarm = alarmService.findLatestByOriginatorAndType(tenantId, originatorId, alarmUpdateMsg.getType()).get(); |
|
||||
switch (alarmUpdateMsg.getMsgType()) { |
|
||||
case ENTITY_CREATED_RPC_MESSAGE: |
|
||||
case ENTITY_UPDATED_RPC_MESSAGE: |
|
||||
if (existentAlarm == null || existentAlarm.getStatus().isCleared()) { |
|
||||
existentAlarm = new Alarm(); |
|
||||
existentAlarm.setTenantId(tenantId); |
|
||||
existentAlarm.setType(alarmUpdateMsg.getName()); |
|
||||
existentAlarm.setOriginator(originatorId); |
|
||||
existentAlarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity())); |
|
||||
existentAlarm.setStartTs(alarmUpdateMsg.getStartTs()); |
|
||||
existentAlarm.setClearTs(alarmUpdateMsg.getClearTs()); |
|
||||
existentAlarm.setPropagate(alarmUpdateMsg.getPropagate()); |
|
||||
} |
|
||||
existentAlarm.setStatus(AlarmStatus.valueOf(alarmUpdateMsg.getStatus())); |
|
||||
existentAlarm.setAckTs(alarmUpdateMsg.getAckTs()); |
|
||||
existentAlarm.setEndTs(alarmUpdateMsg.getEndTs()); |
|
||||
existentAlarm.setDetails(JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); |
|
||||
alarmService.createOrUpdateAlarm(existentAlarm); |
|
||||
break; |
|
||||
case ALARM_ACK_RPC_MESSAGE: |
|
||||
if (existentAlarm != null) { |
|
||||
alarmService.ackAlarm(tenantId, existentAlarm.getId(), alarmUpdateMsg.getAckTs()); |
|
||||
} |
|
||||
break; |
|
||||
case ALARM_CLEAR_RPC_MESSAGE: |
|
||||
if (existentAlarm != null) { |
|
||||
alarmService.clearAlarm(tenantId, existentAlarm.getId(), |
|
||||
JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails()), alarmUpdateMsg.getAckTs()); |
|
||||
} |
|
||||
break; |
|
||||
case ENTITY_DELETED_RPC_MESSAGE: |
|
||||
if (existentAlarm != null) { |
|
||||
alarmService.deleteAlarm(tenantId, existentAlarm.getId()); |
|
||||
} |
|
||||
break; |
|
||||
} |
|
||||
return Futures.immediateFuture(null); |
|
||||
} catch (Exception e) { |
|
||||
log.error("Failed to process alarm update msg [{}]", alarmUpdateMsg, e); |
|
||||
return Futures.immediateFailedFuture(new RuntimeException("Failed to process alarm update msg", e)); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
private EntityId getAlarmOriginator(TenantId tenantId, String entityName, EntityType entityType) { |
|
||||
switch (entityType) { |
|
||||
case DEVICE: |
|
||||
return deviceService.findDeviceByTenantIdAndName(tenantId, entityName).getId(); |
|
||||
case ASSET: |
|
||||
return assetService.findAssetByTenantIdAndName(tenantId, entityName).getId(); |
|
||||
case ENTITY_VIEW: |
|
||||
return entityViewService.findEntityViewByTenantIdAndName(tenantId, entityName).getId(); |
|
||||
default: |
|
||||
return null; |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
public DownlinkMsg convertAlarmEventToDownlink(EdgeEvent edgeEvent) { |
|
||||
AlarmId alarmId = new AlarmId(edgeEvent.getEntityId()); |
|
||||
DownlinkMsg downlinkMsg = null; |
|
||||
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); |
|
||||
switch (edgeEvent.getAction()) { |
|
||||
case ADDED: |
|
||||
case UPDATED: |
|
||||
case ALARM_ACK: |
|
||||
case ALARM_CLEAR: |
|
||||
try { |
|
||||
Alarm alarm = alarmService.findAlarmByIdAsync(edgeEvent.getTenantId(), alarmId).get(); |
|
||||
if (alarm != null) { |
|
||||
downlinkMsg = DownlinkMsg.newBuilder() |
|
||||
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|
||||
.addAlarmUpdateMsg(alarmMsgConstructor.constructAlarmUpdatedMsg(edgeEvent.getTenantId(), msgType, alarm)) |
|
||||
.build(); |
|
||||
} |
|
||||
} catch (Exception e) { |
|
||||
log.error("Can't process alarm msg [{}] [{}]", edgeEvent, msgType, e); |
|
||||
} |
|
||||
break; |
|
||||
case DELETED: |
|
||||
Alarm alarm = JacksonUtil.OBJECT_MAPPER.convertValue(edgeEvent.getBody(), Alarm.class); |
|
||||
AlarmUpdateMsg alarmUpdateMsg = |
|
||||
alarmMsgConstructor.constructAlarmUpdatedMsg(edgeEvent.getTenantId(), msgType, alarm); |
|
||||
downlinkMsg = DownlinkMsg.newBuilder() |
|
||||
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|
||||
.addAlarmUpdateMsg(alarmUpdateMsg) |
|
||||
.build(); |
|
||||
break; |
|
||||
} |
|
||||
return downlinkMsg; |
|
||||
} |
|
||||
|
|
||||
public ListenableFuture<Void> processAlarmNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) throws JsonProcessingException { |
|
||||
EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()); |
|
||||
AlarmId alarmId = new AlarmId(new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); |
|
||||
switch (actionType) { |
|
||||
case DELETED: |
|
||||
EdgeId edgeId = new EdgeId(new UUID(edgeNotificationMsg.getEdgeIdMSB(), edgeNotificationMsg.getEdgeIdLSB())); |
|
||||
Alarm deletedAlarm = JacksonUtil.OBJECT_MAPPER.readValue(edgeNotificationMsg.getBody(), Alarm.class); |
|
||||
return saveEdgeEvent(tenantId, edgeId, EdgeEventType.ALARM, actionType, alarmId, JacksonUtil.OBJECT_MAPPER.valueToTree(deletedAlarm)); |
|
||||
default: |
|
||||
ListenableFuture<Alarm> alarmFuture = alarmService.findAlarmByIdAsync(tenantId, alarmId); |
|
||||
return Futures.transformAsync(alarmFuture, alarm -> { |
|
||||
if (alarm == null) { |
|
||||
return Futures.immediateFuture(null); |
|
||||
} |
|
||||
EdgeEventType type = EdgeUtils.getEdgeEventTypeByEntityType(alarm.getOriginator().getEntityType()); |
|
||||
if (type == null) { |
|
||||
return Futures.immediateFuture(null); |
|
||||
} |
|
||||
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); |
|
||||
PageData<EdgeId> pageData; |
|
||||
List<ListenableFuture<Void>> futures = new ArrayList<>(); |
|
||||
do { |
|
||||
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator(), pageLink); |
|
||||
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { |
|
||||
for (EdgeId relatedEdgeId : pageData.getData()) { |
|
||||
futures.add(saveEdgeEvent(tenantId, |
|
||||
relatedEdgeId, |
|
||||
EdgeEventType.ALARM, |
|
||||
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()), |
|
||||
alarmId, |
|
||||
null)); |
|
||||
} |
|
||||
if (pageData.hasNext()) { |
|
||||
pageLink = pageLink.nextPageLink(); |
|
||||
} |
|
||||
} |
|
||||
} while (pageData != null && pageData.hasNext()); |
|
||||
return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); |
|
||||
}, dbCallbackExecutorService); |
|
||||
} |
|
||||
} |
|
||||
|
|
||||
} |
|
||||
@ -0,0 +1,102 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.edge.rpc.processor.alarm; |
||||
|
|
||||
|
import com.fasterxml.jackson.core.JsonProcessingException; |
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.EdgeUtils; |
||||
|
import org.thingsboard.server.common.data.alarm.Alarm; |
||||
|
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.AlarmId; |
||||
|
import org.thingsboard.server.common.data.id.EdgeId; |
||||
|
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.gen.edge.v1.AlarmUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
|
||||
|
import java.util.ArrayList; |
||||
|
import java.util.List; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
@Component |
||||
|
@Slf4j |
||||
|
@TbCoreComponent |
||||
|
public class AlarmEdgeProcessor extends BaseAlarmProcessor { |
||||
|
|
||||
|
public DownlinkMsg convertAlarmEventToDownlink(EdgeEvent edgeEvent) { |
||||
|
AlarmUpdateMsg alarmUpdateMsg = |
||||
|
convertAlarmEventToAlarmMsg(edgeEvent.getTenantId(), edgeEvent.getEntityId(), edgeEvent.getAction(), edgeEvent.getBody()); |
||||
|
if (alarmUpdateMsg != null) { |
||||
|
return DownlinkMsg.newBuilder() |
||||
|
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
||||
|
.addAlarmUpdateMsg(alarmUpdateMsg) |
||||
|
.build(); |
||||
|
} |
||||
|
return null; |
||||
|
} |
||||
|
|
||||
|
public ListenableFuture<Void> processAlarmNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) throws JsonProcessingException { |
||||
|
EdgeEventActionType actionType = EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()); |
||||
|
AlarmId alarmId = new AlarmId(new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB())); |
||||
|
switch (actionType) { |
||||
|
case DELETED: |
||||
|
EdgeId edgeId = new EdgeId(new UUID(edgeNotificationMsg.getEdgeIdMSB(), edgeNotificationMsg.getEdgeIdLSB())); |
||||
|
Alarm deletedAlarm = JacksonUtil.OBJECT_MAPPER.readValue(edgeNotificationMsg.getBody(), Alarm.class); |
||||
|
return saveEdgeEvent(tenantId, edgeId, EdgeEventType.ALARM, actionType, alarmId, JacksonUtil.OBJECT_MAPPER.valueToTree(deletedAlarm)); |
||||
|
default: |
||||
|
ListenableFuture<Alarm> alarmFuture = alarmService.findAlarmByIdAsync(tenantId, alarmId); |
||||
|
return Futures.transformAsync(alarmFuture, alarm -> { |
||||
|
if (alarm == null) { |
||||
|
return Futures.immediateFuture(null); |
||||
|
} |
||||
|
EdgeEventType type = EdgeUtils.getEdgeEventTypeByEntityType(alarm.getOriginator().getEntityType()); |
||||
|
if (type == null) { |
||||
|
return Futures.immediateFuture(null); |
||||
|
} |
||||
|
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); |
||||
|
PageData<EdgeId> pageData; |
||||
|
List<ListenableFuture<Void>> futures = new ArrayList<>(); |
||||
|
do { |
||||
|
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator(), pageLink); |
||||
|
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { |
||||
|
for (EdgeId relatedEdgeId : pageData.getData()) { |
||||
|
futures.add(saveEdgeEvent(tenantId, |
||||
|
relatedEdgeId, |
||||
|
EdgeEventType.ALARM, |
||||
|
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()), |
||||
|
alarmId, |
||||
|
null)); |
||||
|
} |
||||
|
if (pageData.hasNext()) { |
||||
|
pageLink = pageLink.nextPageLink(); |
||||
|
} |
||||
|
} |
||||
|
} while (pageData != null && pageData.hasNext()); |
||||
|
return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); |
||||
|
}, dbCallbackExecutorService); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,127 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.edge.rpc.processor.alarm; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.EntityType; |
||||
|
import org.thingsboard.server.common.data.alarm.Alarm; |
||||
|
import org.thingsboard.server.common.data.alarm.AlarmSeverity; |
||||
|
import org.thingsboard.server.common.data.alarm.AlarmStatus; |
||||
|
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
||||
|
import org.thingsboard.server.common.data.id.AlarmId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.gen.edge.v1.AlarmUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
@Slf4j |
||||
|
public abstract class BaseAlarmProcessor extends BaseEdgeProcessor { |
||||
|
|
||||
|
public ListenableFuture<Void> processAlarmMsg(TenantId tenantId, AlarmUpdateMsg alarmUpdateMsg) { |
||||
|
log.trace("[{}] processAlarmMsg [{}]", tenantId, alarmUpdateMsg); |
||||
|
EntityId originatorId = getAlarmOriginator(tenantId, alarmUpdateMsg.getOriginatorName(), |
||||
|
EntityType.valueOf(alarmUpdateMsg.getOriginatorType())); |
||||
|
if (originatorId == null) { |
||||
|
log.warn("Originator not found for the alarm msg {}", alarmUpdateMsg); |
||||
|
return Futures.immediateFuture(null); |
||||
|
} |
||||
|
try { |
||||
|
Alarm existentAlarm = alarmService.findLatestByOriginatorAndType(tenantId, originatorId, alarmUpdateMsg.getType()).get(); |
||||
|
switch (alarmUpdateMsg.getMsgType()) { |
||||
|
case ENTITY_CREATED_RPC_MESSAGE: |
||||
|
case ENTITY_UPDATED_RPC_MESSAGE: |
||||
|
if (existentAlarm == null || existentAlarm.getStatus().isCleared()) { |
||||
|
existentAlarm = new Alarm(); |
||||
|
existentAlarm.setTenantId(tenantId); |
||||
|
existentAlarm.setType(alarmUpdateMsg.getName()); |
||||
|
existentAlarm.setOriginator(originatorId); |
||||
|
existentAlarm.setSeverity(AlarmSeverity.valueOf(alarmUpdateMsg.getSeverity())); |
||||
|
existentAlarm.setStartTs(alarmUpdateMsg.getStartTs()); |
||||
|
existentAlarm.setClearTs(alarmUpdateMsg.getClearTs()); |
||||
|
existentAlarm.setPropagate(alarmUpdateMsg.getPropagate()); |
||||
|
} |
||||
|
existentAlarm.setStatus(AlarmStatus.valueOf(alarmUpdateMsg.getStatus())); |
||||
|
existentAlarm.setAckTs(alarmUpdateMsg.getAckTs()); |
||||
|
existentAlarm.setEndTs(alarmUpdateMsg.getEndTs()); |
||||
|
existentAlarm.setDetails(JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails())); |
||||
|
alarmService.createOrUpdateAlarm(existentAlarm); |
||||
|
break; |
||||
|
case ALARM_ACK_RPC_MESSAGE: |
||||
|
if (existentAlarm != null) { |
||||
|
alarmService.ackAlarm(tenantId, existentAlarm.getId(), alarmUpdateMsg.getAckTs()); |
||||
|
} |
||||
|
break; |
||||
|
case ALARM_CLEAR_RPC_MESSAGE: |
||||
|
if (existentAlarm != null) { |
||||
|
alarmService.clearAlarm(tenantId, existentAlarm.getId(), |
||||
|
JacksonUtil.OBJECT_MAPPER.readTree(alarmUpdateMsg.getDetails()), alarmUpdateMsg.getAckTs()); |
||||
|
} |
||||
|
break; |
||||
|
case ENTITY_DELETED_RPC_MESSAGE: |
||||
|
if (existentAlarm != null) { |
||||
|
alarmService.deleteAlarm(tenantId, existentAlarm.getId()); |
||||
|
} |
||||
|
break; |
||||
|
} |
||||
|
return Futures.immediateFuture(null); |
||||
|
} catch (Exception e) { |
||||
|
log.error("[{}] Failed to process alarm update msg [{}]", tenantId, alarmUpdateMsg, e); |
||||
|
return Futures.immediateFailedFuture(e); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
private EntityId getAlarmOriginator(TenantId tenantId, String entityName, EntityType entityType) { |
||||
|
switch (entityType) { |
||||
|
case DEVICE: |
||||
|
return deviceService.findDeviceByTenantIdAndName(tenantId, entityName).getId(); |
||||
|
case ASSET: |
||||
|
return assetService.findAssetByTenantIdAndName(tenantId, entityName).getId(); |
||||
|
case ENTITY_VIEW: |
||||
|
return entityViewService.findEntityViewByTenantIdAndName(tenantId, entityName).getId(); |
||||
|
default: |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public AlarmUpdateMsg convertAlarmEventToAlarmMsg(TenantId tenantId, UUID entityId, EdgeEventActionType actionType, JsonNode body) { |
||||
|
AlarmId alarmId = new AlarmId(entityId); |
||||
|
UpdateMsgType msgType = getUpdateMsgType(actionType); |
||||
|
switch (actionType) { |
||||
|
case ADDED: |
||||
|
case UPDATED: |
||||
|
case ALARM_ACK: |
||||
|
case ALARM_CLEAR: |
||||
|
Alarm alarm = alarmService.findAlarmById(tenantId, alarmId); |
||||
|
if (alarm != null) { |
||||
|
return alarmMsgConstructor.constructAlarmUpdatedMsg(tenantId, msgType, alarm); |
||||
|
} |
||||
|
break; |
||||
|
case DELETED: |
||||
|
Alarm deletedAlarm = JacksonUtil.OBJECT_MAPPER.convertValue(body, Alarm.class); |
||||
|
return alarmMsgConstructor.constructAlarmUpdatedMsg(tenantId, msgType, deletedAlarm); |
||||
|
} |
||||
|
return null; |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,134 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.edge.rpc.processor.device; |
||||
|
|
||||
|
import com.datastax.oss.driver.api.core.uuid.Uuids; |
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.beans.factory.annotation.Autowired; |
||||
|
import org.springframework.data.util.Pair; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.Device; |
||||
|
import org.thingsboard.server.common.data.StringUtils; |
||||
|
import org.thingsboard.server.common.data.device.data.DeviceData; |
||||
|
import org.thingsboard.server.common.data.id.CustomerId; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.common.data.id.DeviceProfileId; |
||||
|
import org.thingsboard.server.common.data.id.OtaPackageId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.security.DeviceCredentials; |
||||
|
import org.thingsboard.server.common.data.security.DeviceCredentialsType; |
||||
|
import org.thingsboard.server.gen.edge.v1.DeviceCredentialsUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.DeviceUpdateMsg; |
||||
|
import org.thingsboard.server.queue.util.DataDecodingEncodingService; |
||||
|
import org.thingsboard.server.service.edge.rpc.processor.BaseEdgeProcessor; |
||||
|
|
||||
|
import java.util.Optional; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
@Slf4j |
||||
|
public abstract class BaseDeviceProcessor extends BaseEdgeProcessor { |
||||
|
|
||||
|
@Autowired |
||||
|
protected DataDecodingEncodingService dataDecodingEncodingService; |
||||
|
|
||||
|
protected Pair<Boolean, Boolean> saveOrUpdateDevice(TenantId tenantId, DeviceId deviceId, DeviceUpdateMsg deviceUpdateMsg, CustomerId customerId) { |
||||
|
boolean created = false; |
||||
|
boolean deviceNameUpdated = false; |
||||
|
deviceCreationLock.lock(); |
||||
|
try { |
||||
|
Device device = deviceService.findDeviceById(tenantId, deviceId); |
||||
|
String deviceName = deviceUpdateMsg.getName(); |
||||
|
if (device == null) { |
||||
|
created = true; |
||||
|
device = new Device(); |
||||
|
device.setTenantId(tenantId); |
||||
|
device.setCreatedTime(Uuids.unixTimestamp(deviceId.getId())); |
||||
|
Device deviceByName = deviceService.findDeviceByTenantIdAndName(tenantId, deviceName); |
||||
|
if (deviceByName != null) { |
||||
|
deviceName = deviceName + "_" + StringUtils.randomAlphabetic(15); |
||||
|
log.warn("Device with name {} already exists. Renaming device name to {}", |
||||
|
deviceUpdateMsg.getName(), deviceName); |
||||
|
deviceNameUpdated = true; |
||||
|
} |
||||
|
} |
||||
|
device.setName(deviceName); |
||||
|
device.setType(deviceUpdateMsg.getType()); |
||||
|
device.setLabel(deviceUpdateMsg.hasLabel() ? deviceUpdateMsg.getLabel() : null); |
||||
|
device.setAdditionalInfo(deviceUpdateMsg.hasAdditionalInfo() |
||||
|
? JacksonUtil.toJsonNode(deviceUpdateMsg.getAdditionalInfo()) : null); |
||||
|
|
||||
|
UUID deviceProfileUUID = safeGetUUID(deviceUpdateMsg.getDeviceProfileIdMSB(), deviceUpdateMsg.getDeviceProfileIdLSB()); |
||||
|
device.setDeviceProfileId(deviceProfileUUID != null ? new DeviceProfileId(deviceProfileUUID) : null); |
||||
|
|
||||
|
device.setCustomerId(customerId); |
||||
|
|
||||
|
Optional<DeviceData> deviceDataOpt = |
||||
|
dataDecodingEncodingService.decode(deviceUpdateMsg.getDeviceDataBytes().toByteArray()); |
||||
|
device.setDeviceData(deviceDataOpt.orElse(null)); |
||||
|
|
||||
|
UUID firmwareUUID = safeGetUUID(deviceUpdateMsg.getFirmwareIdMSB(), deviceUpdateMsg.getFirmwareIdLSB()); |
||||
|
device.setFirmwareId(firmwareUUID != null ? new OtaPackageId(firmwareUUID) : null); |
||||
|
|
||||
|
UUID softwareUUID = safeGetUUID(deviceUpdateMsg.getSoftwareIdMSB(), deviceUpdateMsg.getSoftwareIdLSB()); |
||||
|
device.setSoftwareId(softwareUUID != null ? new OtaPackageId(softwareUUID) : null); |
||||
|
deviceValidator.validate(device, Device::getTenantId); |
||||
|
if (created) { |
||||
|
device.setId(deviceId); |
||||
|
} |
||||
|
Device savedDevice = deviceService.saveDevice(device, false); |
||||
|
if (created) { |
||||
|
DeviceCredentials deviceCredentials = new DeviceCredentials(); |
||||
|
deviceCredentials.setDeviceId(new DeviceId(savedDevice.getUuidId())); |
||||
|
deviceCredentials.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN); |
||||
|
deviceCredentials.setCredentialsId(StringUtils.randomAlphanumeric(20)); |
||||
|
deviceCredentialsService.createDeviceCredentials(device.getTenantId(), deviceCredentials); |
||||
|
} |
||||
|
tbClusterService.onDeviceUpdated(savedDevice, created ? null : device, false); |
||||
|
} finally { |
||||
|
deviceCreationLock.unlock(); |
||||
|
} |
||||
|
return Pair.of(created, deviceNameUpdated); |
||||
|
} |
||||
|
|
||||
|
public ListenableFuture<Void> processDeviceCredentialsMsg(TenantId tenantId, DeviceCredentialsUpdateMsg deviceCredentialsUpdateMsg) { |
||||
|
log.debug("[{}] Executing processDeviceCredentialsMsg, deviceCredentialsUpdateMsg [{}]", tenantId, deviceCredentialsUpdateMsg); |
||||
|
DeviceId deviceId = new DeviceId(new UUID(deviceCredentialsUpdateMsg.getDeviceIdMSB(), deviceCredentialsUpdateMsg.getDeviceIdLSB())); |
||||
|
ListenableFuture<Device> deviceFuture = deviceService.findDeviceByIdAsync(tenantId, deviceId); |
||||
|
return Futures.transform(deviceFuture, device -> { |
||||
|
if (device != null) { |
||||
|
log.debug("Updating device credentials for device [{}]. New device credentials Id [{}], value [{}]", |
||||
|
device.getName(), deviceCredentialsUpdateMsg.getCredentialsId(), deviceCredentialsUpdateMsg.getCredentialsValue()); |
||||
|
try { |
||||
|
DeviceCredentials deviceCredentials = deviceCredentialsService.findDeviceCredentialsByDeviceId(tenantId, device.getId()); |
||||
|
deviceCredentials.setCredentialsType(DeviceCredentialsType.valueOf(deviceCredentialsUpdateMsg.getCredentialsType())); |
||||
|
deviceCredentials.setCredentialsId(deviceCredentialsUpdateMsg.getCredentialsId()); |
||||
|
deviceCredentials.setCredentialsValue(deviceCredentialsUpdateMsg.hasCredentialsValue() |
||||
|
? deviceCredentialsUpdateMsg.getCredentialsValue() : null); |
||||
|
deviceCredentialsService.updateDeviceCredentials(tenantId, deviceCredentials); |
||||
|
} catch (Exception e) { |
||||
|
log.error("Can't update device credentials for device [{}], deviceCredentialsUpdateMsg [{}]", |
||||
|
device.getName(), deviceCredentialsUpdateMsg, e); |
||||
|
throw new RuntimeException(e); |
||||
|
} |
||||
|
} else { |
||||
|
log.warn("Can't find device by id [{}], deviceCredentialsUpdateMsg [{}]", deviceId, deviceCredentialsUpdateMsg); |
||||
|
} |
||||
|
return null; |
||||
|
}, dbCallbackExecutorService); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,82 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.edge.rpc.processor.relation; |
||||
|
|
||||
|
import com.fasterxml.jackson.core.JsonProcessingException; |
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.EdgeUtils; |
||||
|
import org.thingsboard.server.common.data.EntityType; |
||||
|
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.EdgeId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.relation.EntityRelation; |
||||
|
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.RelationUpdateMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
|
||||
|
import java.util.ArrayList; |
||||
|
import java.util.HashSet; |
||||
|
import java.util.List; |
||||
|
import java.util.Set; |
||||
|
|
||||
|
@Component |
||||
|
@Slf4j |
||||
|
@TbCoreComponent |
||||
|
public class RelationEdgeProcessor extends BaseRelationProcessor { |
||||
|
|
||||
|
public DownlinkMsg convertRelationEventToDownlink(EdgeEvent edgeEvent) { |
||||
|
EntityRelation entityRelation = JacksonUtil.OBJECT_MAPPER.convertValue(edgeEvent.getBody(), EntityRelation.class); |
||||
|
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); |
||||
|
RelationUpdateMsg relationUpdateMsg = relationMsgConstructor.constructRelationUpdatedMsg(msgType, entityRelation); |
||||
|
return DownlinkMsg.newBuilder() |
||||
|
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
||||
|
.addRelationUpdateMsg(relationUpdateMsg) |
||||
|
.build(); |
||||
|
} |
||||
|
|
||||
|
public ListenableFuture<Void> processRelationNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) throws JsonProcessingException { |
||||
|
EntityRelation relation = JacksonUtil.OBJECT_MAPPER.readValue(edgeNotificationMsg.getBody(), EntityRelation.class); |
||||
|
if (relation.getFrom().getEntityType().equals(EntityType.EDGE) || |
||||
|
relation.getTo().getEntityType().equals(EntityType.EDGE)) { |
||||
|
return Futures.immediateFuture(null); |
||||
|
} |
||||
|
|
||||
|
Set<EdgeId> uniqueEdgeIds = new HashSet<>(); |
||||
|
uniqueEdgeIds.addAll(edgeService.findAllRelatedEdgeIds(tenantId, relation.getTo())); |
||||
|
uniqueEdgeIds.addAll(edgeService.findAllRelatedEdgeIds(tenantId, relation.getFrom())); |
||||
|
if (uniqueEdgeIds.isEmpty()) { |
||||
|
return Futures.immediateFuture(null); |
||||
|
} |
||||
|
List<ListenableFuture<Void>> futures = new ArrayList<>(); |
||||
|
for (EdgeId edgeId : uniqueEdgeIds) { |
||||
|
futures.add(saveEdgeEvent(tenantId, |
||||
|
edgeId, |
||||
|
EdgeEventType.RELATION, |
||||
|
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()), |
||||
|
null, |
||||
|
JacksonUtil.OBJECT_MAPPER.valueToTree(relation))); |
||||
|
} |
||||
|
return Futures.transform(Futures.allAsList(futures), voids -> null, dbCallbackExecutorService); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,48 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2022 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.service.edge.rpc.processor.telemetry; |
||||
|
|
||||
|
import com.fasterxml.jackson.core.JsonProcessingException; |
||||
|
import lombok.extern.slf4j.Slf4j; |
||||
|
import org.springframework.stereotype.Component; |
||||
|
import org.thingsboard.server.common.data.DataConstants; |
||||
|
import org.thingsboard.server.common.data.EdgeUtils; |
||||
|
import org.thingsboard.server.common.data.EntityType; |
||||
|
import org.thingsboard.server.common.data.edge.EdgeEvent; |
||||
|
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; |
||||
|
import org.thingsboard.server.gen.edge.v1.EntityDataProto; |
||||
|
import org.thingsboard.server.queue.util.TbCoreComponent; |
||||
|
|
||||
|
@Component |
||||
|
@Slf4j |
||||
|
@TbCoreComponent |
||||
|
public class TelemetryEdgeProcessor extends BaseTelemetryProcessor { |
||||
|
|
||||
|
@Override |
||||
|
protected String getMsgSourceKey() { |
||||
|
return DataConstants.EDGE_MSG_SOURCE; |
||||
|
} |
||||
|
|
||||
|
public DownlinkMsg convertTelemetryEventToDownlink(EdgeEvent edgeEvent) throws JsonProcessingException { |
||||
|
EntityType entityType = EntityType.valueOf(edgeEvent.getType().name()); |
||||
|
EntityDataProto entityDataProto = convertTelemetryEventToEntityDataProto(entityType, edgeEvent.getEntityId(), |
||||
|
edgeEvent.getAction(), edgeEvent.getBody()); |
||||
|
return DownlinkMsg.newBuilder() |
||||
|
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
||||
|
.addEntityData(entityDataProto) |
||||
|
.build(); |
||||
|
} |
||||
|
} |
||||
Loading…
Reference in new issue