Browse Source

Added try/catch for processDeviceFromEdge

pull/7982/head
Volodymyr Babak 4 years ago
parent
commit
69285ee596
  1. 89
      application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessor.java

89
application/src/main/java/org/thingsboard/server/service/edge/rpc/processor/device/DeviceEdgeProcessor.java

@ -81,56 +81,61 @@ public class DeviceEdgeProcessor extends BaseEdgeProcessor {
private DataDecodingEncodingService dataDecodingEncodingService; private DataDecodingEncodingService dataDecodingEncodingService;
public ListenableFuture<Void> processDeviceFromEdge(TenantId tenantId, Edge edge, DeviceUpdateMsg deviceUpdateMsg) { public ListenableFuture<Void> processDeviceFromEdge(TenantId tenantId, Edge edge, DeviceUpdateMsg deviceUpdateMsg) {
log.trace("[{}] processDeviceFromEdge [{}] from edge [{}]", tenantId, deviceUpdateMsg, edge.getName()); try {
DeviceId deviceId = new DeviceId(new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB())); log.trace("[{}] processDeviceFromEdge [{}] from edge [{}]", tenantId, deviceUpdateMsg, edge.getName());
switch (deviceUpdateMsg.getMsgType()) { DeviceId deviceId = new DeviceId(new UUID(deviceUpdateMsg.getIdMSB(), deviceUpdateMsg.getIdLSB()));
case ENTITY_CREATED_RPC_MESSAGE: switch (deviceUpdateMsg.getMsgType()) {
String deviceName = deviceUpdateMsg.getName(); case ENTITY_CREATED_RPC_MESSAGE:
Device device = deviceService.findDeviceByTenantIdAndName(tenantId, deviceName); String deviceName = deviceUpdateMsg.getName();
if (device != null) { Device device = deviceService.findDeviceByTenantIdAndName(tenantId, deviceName);
boolean deviceAlreadyExistsForThisEdge = isDeviceAlreadyExistsOnCloudForThisEdge(tenantId, edge, device); if (device != null) {
if (deviceAlreadyExistsForThisEdge) { boolean deviceAlreadyExistsForThisEdge = isDeviceAlreadyExistsOnCloudForThisEdge(tenantId, edge, device);
log.info("[{}] Device with name '{}' already exists on the cloud, and related to this edge [{}]. " + if (deviceAlreadyExistsForThisEdge) {
"deviceUpdateMsg [{}], Updating device", tenantId, deviceName, edge.getId(), deviceUpdateMsg); log.info("[{}] Device with name '{}' already exists on the cloud, and related to this edge [{}]. " +
return updateDevice(tenantId, edge, deviceUpdateMsg); "deviceUpdateMsg [{}], Updating device", tenantId, deviceName, edge.getId(), deviceUpdateMsg);
return updateDevice(tenantId, edge, deviceUpdateMsg);
} else {
log.info("[{}] Device with name '{}' already exists on the cloud, but not related to this edge [{}]. deviceUpdateMsg [{}]." +
"Creating a new device with random prefix and relate to this edge", tenantId, deviceName, edge.getId(), deviceUpdateMsg);
String newDeviceName = deviceUpdateMsg.getName() + "_" + StringUtils.randomAlphabetic(15);
try {
createDevice(tenantId, deviceId, edge, deviceUpdateMsg, newDeviceName);
} catch (DataValidationException e) {
log.error("[{}] Device update msg can't be processed due to data validation [{}]", tenantId, deviceUpdateMsg, e);
return Futures.immediateFuture(null);
}
ObjectNode body = JacksonUtil.OBJECT_MAPPER.createObjectNode();
body.put("conflictName", deviceName);
ListenableFuture<Void> input = saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ENTITY_MERGE_REQUEST, deviceId, body);
return Futures.transformAsync(input, unused ->
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_REQUEST, deviceId, null),
dbCallbackExecutorService);
}
} else { } else {
log.info("[{}] Device with name '{}' already exists on the cloud, but not related to this edge [{}]. deviceUpdateMsg [{}]." + log.info("[{}] Creating new device on the cloud [{}]", tenantId, deviceUpdateMsg);
"Creating a new device with random prefix and relate to this edge", tenantId, deviceName, edge.getId(), deviceUpdateMsg);
String newDeviceName = deviceUpdateMsg.getName() + "_" + StringUtils.randomAlphabetic(15);
try { try {
createDevice(tenantId, deviceId, edge, deviceUpdateMsg, newDeviceName); createDevice(tenantId, deviceId, edge, deviceUpdateMsg, deviceUpdateMsg.getName());
} catch (DataValidationException e) { } catch (DataValidationException e) {
log.error("[{}] Device update msg can't be processed due to data validation [{}]", tenantId, deviceUpdateMsg, e); log.error("[{}] Device update msg can't be processed due to data validation [{}]", tenantId, deviceUpdateMsg, e);
return Futures.immediateFuture(null); return Futures.immediateFuture(null);
} }
ObjectNode body = JacksonUtil.OBJECT_MAPPER.createObjectNode(); return saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_REQUEST, deviceId, null);
body.put("conflictName", deviceName);
ListenableFuture<Void> input = saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ENTITY_MERGE_REQUEST, deviceId, body);
return Futures.transformAsync(input, unused ->
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_REQUEST, deviceId, null),
dbCallbackExecutorService);
} }
} else { case ENTITY_UPDATED_RPC_MESSAGE:
log.info("[{}] Creating new device on the cloud [{}]", tenantId, deviceUpdateMsg); return updateDevice(tenantId, edge, deviceUpdateMsg);
try { case ENTITY_DELETED_RPC_MESSAGE:
createDevice(tenantId, deviceId, edge, deviceUpdateMsg, deviceUpdateMsg.getName()); Device deviceToDelete = deviceService.findDeviceById(tenantId, deviceId);
} catch (DataValidationException e) { if (deviceToDelete != null) {
log.error("[{}] Device update msg can't be processed due to data validation [{}]", tenantId, deviceUpdateMsg, e); deviceService.unassignDeviceFromEdge(tenantId, deviceId, edge.getId());
return Futures.immediateFuture(null);
} }
return saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_REQUEST, deviceId, null); return Futures.immediateFuture(null);
} case UNRECOGNIZED:
case ENTITY_UPDATED_RPC_MESSAGE: default:
return updateDevice(tenantId, edge, deviceUpdateMsg); return handleUnsupportedMsgType(deviceUpdateMsg.getMsgType());
case ENTITY_DELETED_RPC_MESSAGE: }
Device deviceToDelete = deviceService.findDeviceById(tenantId, deviceId); } catch (Exception e) {
if (deviceToDelete != null) { log.error("Failed to process device message from edge, {}", deviceUpdateMsg, e);
deviceService.unassignDeviceFromEdge(tenantId, deviceId, edge.getId()); return Futures.immediateFailedFuture(e);
}
return Futures.immediateFuture(null);
case UNRECOGNIZED:
default:
return handleUnsupportedMsgType(deviceUpdateMsg.getMsgType());
} }
} }

Loading…
Cancel
Save