Browse Source

Fixed alarm delete functionality on edge

pull/5038/head
Volodymyr Babak 5 years ago
parent
commit
1297fef550
  1. 8
      application/src/main/java/org/thingsboard/server/controller/AlarmController.java
  2. 14
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  3. 1
      application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java
  4. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  5. 108
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java
  6. 2
      rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java

8
application/src/main/java/org/thingsboard/server/controller/AlarmController.java

@ -38,6 +38,7 @@ import org.thingsboard.server.common.data.edge.EdgeEventActionType;
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode; import org.thingsboard.server.common.data.exception.ThingsboardErrorCode;
import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.AlarmId; 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.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.EntityIdFactory;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
@ -46,6 +47,8 @@ import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Operation;
import org.thingsboard.server.service.security.permission.Resource; import org.thingsboard.server.service.security.permission.Resource;
import java.util.List;
@RestController @RestController
@TbCoreComponent @TbCoreComponent
@RequestMapping("/api") @RequestMapping("/api")
@ -112,10 +115,13 @@ public class AlarmController extends BaseController {
AlarmId alarmId = new AlarmId(toUUID(strAlarmId)); AlarmId alarmId = new AlarmId(toUUID(strAlarmId));
Alarm alarm = checkAlarmId(alarmId, Operation.WRITE); Alarm alarm = checkAlarmId(alarmId, Operation.WRITE);
List<EdgeId> relatedEdgeIds = findRelatedEdgeIds(getTenantId(), alarm.getOriginator());
logEntityAction(alarm.getOriginator(), alarm, logEntityAction(alarm.getOriginator(), alarm,
getCurrentUser().getCustomerId(), getCurrentUser().getCustomerId(),
ActionType.ALARM_DELETE, null); ActionType.ALARM_DELETE, null);
sendEntityNotificationMsg(getTenantId(), alarmId, EdgeEventActionType.DELETED);
sendAlarmDeleteNotificationMsg(getTenantId(), alarmId, relatedEdgeIds, alarm);
return alarmService.deleteAlarm(getTenantId(), alarmId); return alarmService.deleteAlarm(getTenantId(), alarmId);
} catch (Exception e) { } catch (Exception e) {

14
application/src/main/java/org/thingsboard/server/controller/BaseController.java

@ -852,13 +852,25 @@ public abstract class BaseController {
} }
protected void sendDeleteNotificationMsg(TenantId tenantId, EntityId entityId, List<EdgeId> edgeIds) { protected void sendDeleteNotificationMsg(TenantId tenantId, EntityId entityId, List<EdgeId> edgeIds) {
sendDeleteNotificationMsg(tenantId, entityId, edgeIds, null);
}
protected void sendDeleteNotificationMsg(TenantId tenantId, EntityId entityId, List<EdgeId> edgeIds, String body) {
if (edgeIds != null && !edgeIds.isEmpty()) { if (edgeIds != null && !edgeIds.isEmpty()) {
for (EdgeId edgeId : edgeIds) { for (EdgeId edgeId : edgeIds) {
sendNotificationMsgToEdgeService(tenantId, edgeId, entityId, null, null, EdgeEventActionType.DELETED); sendNotificationMsgToEdgeService(tenantId, edgeId, entityId, body, null, EdgeEventActionType.DELETED);
} }
} }
} }
protected void sendAlarmDeleteNotificationMsg(TenantId tenantId, EntityId entityId, List<EdgeId> edgeIds, Alarm alarm) {
try {
sendDeleteNotificationMsg(tenantId, entityId, edgeIds, json.writeValueAsString(alarm));
} catch (Exception e) {
log.warn("Failed to push delete alarm msg to core: {}", alarm, e);
}
}
protected void sendEntityAssignToCustomerNotificationMsg(TenantId tenantId, EntityId entityId, CustomerId customerId, EdgeEventActionType action) { protected void sendEntityAssignToCustomerNotificationMsg(TenantId tenantId, EntityId entityId, CustomerId customerId, EdgeEventActionType action) {
try { try {
sendNotificationMsgToEdgeService(tenantId, null, entityId, json.writeValueAsString(customerId), null, action); sendNotificationMsgToEdgeService(tenantId, null, entityId, json.writeValueAsString(customerId), null, action);

1
application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java

@ -121,6 +121,7 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
@Override @Override
public void pushNotificationToEdge(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback) { public void pushNotificationToEdge(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback) {
log.trace("Pushing notification to edge {}", edgeNotificationMsg);
try { try {
TenantId tenantId = new TenantId(new UUID(edgeNotificationMsg.getTenantIdMSB(), edgeNotificationMsg.getTenantIdLSB())); TenantId tenantId = new TenantId(new UUID(edgeNotificationMsg.getTenantIdMSB(), edgeNotificationMsg.getTenantIdLSB()));
EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType()); EdgeEventType type = EdgeEventType.valueOf(edgeNotificationMsg.getType());

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

@ -527,7 +527,7 @@ public final class EdgeGrpcSession implements Closeable {
case RULE_CHAIN_METADATA: case RULE_CHAIN_METADATA:
return ctx.getRuleChainProcessor().processRuleChainMetadataToEdge(edgeEvent, msgType); return ctx.getRuleChainProcessor().processRuleChainMetadataToEdge(edgeEvent, msgType);
case ALARM: case ALARM:
return ctx.getAlarmProcessor().processAlarmToEdge(edge, edgeEvent, msgType); return ctx.getAlarmProcessor().processAlarmToEdge(edge, edgeEvent, msgType, action);
case USER: case USER:
return ctx.getUserProcessor().processUserToEdge(edge, edgeEvent, msgType, action); return ctx.getUserProcessor().processUserToEdge(edge, edgeEvent, msgType, action);
case RELATION: case RELATION:

108
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/AlarmEdgeProcessor.java

@ -15,6 +15,7 @@
*/ */
package org.thingsboard.server.service.edge.rpc.processor; package org.thingsboard.server.service.edge.rpc.processor;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
@ -114,59 +115,84 @@ public class AlarmEdgeProcessor extends BaseEdgeProcessor {
} }
} }
public DownlinkMsg processAlarmToEdge(Edge edge, EdgeEvent edgeEvent, UpdateMsgType msgType) { public DownlinkMsg processAlarmToEdge(Edge edge, EdgeEvent edgeEvent, UpdateMsgType msgType, EdgeEventActionType action) {
AlarmId alarmId = new AlarmId(edgeEvent.getEntityId());
DownlinkMsg downlinkMsg = null; DownlinkMsg downlinkMsg = null;
try { switch (action) {
AlarmId alarmId = new AlarmId(edgeEvent.getEntityId()); case ADDED:
Alarm alarm = alarmService.findAlarmByIdAsync(edgeEvent.getTenantId(), alarmId).get(); case UPDATED:
if (alarm != null) { 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(edge.getTenantId(), msgType, alarm))
.build();
}
} catch (Exception e) {
log.error("Can't process alarm msg [{}] [{}]", edgeEvent, msgType, e);
}
break;
case DELETED:
Alarm alarm = mapper.convertValue(edgeEvent.getBody(), Alarm.class);
AlarmUpdateMsg alarmUpdateMsg =
alarmMsgConstructor.constructAlarmUpdatedMsg(edge.getTenantId(), msgType, alarm);
downlinkMsg = DownlinkMsg.newBuilder() downlinkMsg = DownlinkMsg.newBuilder()
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) .setDownlinkMsgId(EdgeUtils.nextPositiveInt())
.addAlarmUpdateMsg(alarmMsgConstructor.constructAlarmUpdatedMsg(edge.getTenantId(), msgType, alarm)) .addAlarmUpdateMsg(alarmUpdateMsg)
.build(); .build();
} break;
} catch (Exception e) {
log.error("Can't process alarm msg [{}] [{}]", edgeEvent, msgType, e);
} }
return downlinkMsg; return downlinkMsg;
} }
public void processAlarmNotification(TenantId tenantId, TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg) { public 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())); AlarmId alarmId = new AlarmId(new UUID(edgeNotificationMsg.getEntityIdMSB(), edgeNotificationMsg.getEntityIdLSB()));
ListenableFuture<Alarm> alarmFuture = alarmService.findAlarmByIdAsync(tenantId, alarmId); switch (actionType) {
Futures.addCallback(alarmFuture, new FutureCallback<Alarm>() { case DELETED:
@Override EdgeId edgeId = new EdgeId(new UUID(edgeNotificationMsg.getEdgeIdMSB(), edgeNotificationMsg.getEdgeIdLSB()));
public void onSuccess(@Nullable Alarm alarm) { Alarm alarm = mapper.readValue(edgeNotificationMsg.getBody(), Alarm.class);
if (alarm != null) { saveEdgeEvent(tenantId, edgeId, EdgeEventType.ALARM, actionType, alarmId, mapper.valueToTree(alarm));
EdgeEventType type = EdgeUtils.getEdgeEventTypeByEntityType(alarm.getOriginator().getEntityType()); break;
if (type != null) { default:
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); ListenableFuture<Alarm> alarmFuture = alarmService.findAlarmByIdAsync(tenantId, alarmId);
PageData<EdgeId> pageData; Futures.addCallback(alarmFuture, new FutureCallback<Alarm>() {
do { @Override
pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator(), pageLink); public void onSuccess(@Nullable Alarm alarm) {
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { if (alarm != null) {
for (EdgeId edgeId : pageData.getData()) { EdgeEventType type = EdgeUtils.getEdgeEventTypeByEntityType(alarm.getOriginator().getEntityType());
saveEdgeEvent(tenantId, if (type != null) {
edgeId, PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE);
EdgeEventType.ALARM, PageData<EdgeId> pageData;
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()), do {
alarmId, pageData = edgeService.findRelatedEdgeIdsByEntityId(tenantId, alarm.getOriginator(), pageLink);
null); if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
} for (EdgeId edgeId : pageData.getData()) {
if (pageData.hasNext()) { saveEdgeEvent(tenantId,
pageLink = pageLink.nextPageLink(); edgeId,
} EdgeEventType.ALARM,
EdgeEventActionType.valueOf(edgeNotificationMsg.getAction()),
alarmId,
null);
}
if (pageData.hasNext()) {
pageLink = pageLink.nextPageLink();
}
}
} while (pageData != null && pageData.hasNext());
} }
} while (pageData != null && pageData.hasNext()); }
} }
}
}
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.warn("[{}] can't find alarm by id [{}] {}", tenantId.getId(), alarmId.getId(), t); log.warn("[{}] can't find alarm by id [{}] {}", tenantId.getId(), alarmId.getId(), t);
} }
}, dbCallbackExecutorService); }, dbCallbackExecutorService);
}
} }
} }

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

@ -337,7 +337,7 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable {
addTimePageLinkToParam(params, pageLink); addTimePageLinkToParam(params, pageLink);
return restTemplate.exchange( return restTemplate.exchange(
baseURL + urlSecondPart + getTimeUrlParams(pageLink), baseURL + urlSecondPart + "&" + getTimeUrlParams(pageLink),
HttpMethod.GET, HttpMethod.GET,
HttpEntity.EMPTY, HttpEntity.EMPTY,
new ParameterizedTypeReference<PageData<AlarmInfo>>() { new ParameterizedTypeReference<PageData<AlarmInfo>>() {

Loading…
Cancel
Save