|
|
@ -122,26 +122,13 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { |
|
|
@Override |
|
|
@Override |
|
|
public ListenableFuture<Void> processRuleChainMetadataRequestMsg(TenantId tenantId, Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg) { |
|
|
public ListenableFuture<Void> processRuleChainMetadataRequestMsg(TenantId tenantId, Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg) { |
|
|
log.trace("[{}] processRuleChainMetadataRequestMsg [{}][{}]", tenantId, edge.getName(), ruleChainMetadataRequestMsg); |
|
|
log.trace("[{}] processRuleChainMetadataRequestMsg [{}][{}]", tenantId, edge.getName(), ruleChainMetadataRequestMsg); |
|
|
SettableFuture<Void> futureToSet = SettableFuture.create(); |
|
|
|
|
|
if (ruleChainMetadataRequestMsg.getRuleChainIdMSB() != 0 && ruleChainMetadataRequestMsg.getRuleChainIdLSB() != 0) { |
|
|
if (ruleChainMetadataRequestMsg.getRuleChainIdMSB() != 0 && ruleChainMetadataRequestMsg.getRuleChainIdLSB() != 0) { |
|
|
RuleChainId ruleChainId = |
|
|
RuleChainId ruleChainId = |
|
|
new RuleChainId(new UUID(ruleChainMetadataRequestMsg.getRuleChainIdMSB(), ruleChainMetadataRequestMsg.getRuleChainIdLSB())); |
|
|
new RuleChainId(new UUID(ruleChainMetadataRequestMsg.getRuleChainIdMSB(), ruleChainMetadataRequestMsg.getRuleChainIdLSB())); |
|
|
ListenableFuture<EdgeEvent> future = saveEdgeEvent(tenantId, edge.getId(), |
|
|
saveEdgeEvent(tenantId, edge.getId(), |
|
|
EdgeEventType.RULE_CHAIN_METADATA, EdgeEventActionType.ADDED, ruleChainId, null); |
|
|
EdgeEventType.RULE_CHAIN_METADATA, EdgeEventActionType.ADDED, ruleChainId, null); |
|
|
Futures.addCallback(future, new FutureCallback<EdgeEvent>() { |
|
|
|
|
|
@Override |
|
|
|
|
|
public void onSuccess(@Nullable EdgeEvent result) { |
|
|
|
|
|
futureToSet.set(null); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
log.error("Can't save edge event [{}]", ruleChainMetadataRequestMsg, t); |
|
|
|
|
|
futureToSet.setException(t); |
|
|
|
|
|
} |
|
|
|
|
|
}, dbCallbackExecutorService); |
|
|
|
|
|
} |
|
|
} |
|
|
return futureToSet; |
|
|
return Futures.immediateFuture(null); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
@ -154,8 +141,8 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { |
|
|
if (type != null) { |
|
|
if (type != null) { |
|
|
SettableFuture<Void> futureToSet = SettableFuture.create(); |
|
|
SettableFuture<Void> futureToSet = SettableFuture.create(); |
|
|
String scope = attributesRequestMsg.getScope(); |
|
|
String scope = attributesRequestMsg.getScope(); |
|
|
ListenableFuture<List<AttributeKvEntry>> ssAttrFuture = attributesService.findAll(tenantId, entityId, scope); |
|
|
ListenableFuture<List<AttributeKvEntry>> findAttrFuture = attributesService.findAll(tenantId, entityId, scope); |
|
|
Futures.addCallback(ssAttrFuture, new FutureCallback<List<AttributeKvEntry>>() { |
|
|
Futures.addCallback(findAttrFuture, new FutureCallback<List<AttributeKvEntry>>() { |
|
|
@Override |
|
|
@Override |
|
|
public void onSuccess(@Nullable List<AttributeKvEntry> ssAttributes) { |
|
|
public void onSuccess(@Nullable List<AttributeKvEntry> ssAttributes) { |
|
|
if (ssAttributes != null && !ssAttributes.isEmpty()) { |
|
|
if (ssAttributes != null && !ssAttributes.isEmpty()) { |
|
|
@ -184,8 +171,9 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { |
|
|
entityId, |
|
|
entityId, |
|
|
body); |
|
|
body); |
|
|
} catch (Exception e) { |
|
|
} catch (Exception e) { |
|
|
log.error("[{}] Failed to send attribute updates to the edge", edge.getName(), e); |
|
|
log.error("[{}] Failed to save attribute updates to the edge", edge.getName(), e); |
|
|
throw new RuntimeException("[" + edge.getName() + "] Failed to send attribute updates to the edge", e); |
|
|
futureToSet.setException(new RuntimeException("[" + edge.getName() + "] Failed to send attribute updates to the edge", e)); |
|
|
|
|
|
return; |
|
|
} |
|
|
} |
|
|
} else { |
|
|
} else { |
|
|
log.trace("[{}][{}] No attributes found for entity {} [{}]", tenantId, |
|
|
log.trace("[{}][{}] No attributes found for entity {} [{}]", tenantId, |
|
|
@ -198,7 +186,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public void onFailure(Throwable t) { |
|
|
public void onFailure(Throwable t) { |
|
|
log.error("Can't save attributes [{}]", attributesRequestMsg, t); |
|
|
log.error("Can't find attributes [{}]", attributesRequestMsg, t); |
|
|
futureToSet.setException(t); |
|
|
futureToSet.setException(t); |
|
|
} |
|
|
} |
|
|
}, dbCallbackExecutorService); |
|
|
}, dbCallbackExecutorService); |
|
|
@ -273,82 +261,39 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { |
|
|
@Override |
|
|
@Override |
|
|
public ListenableFuture<Void> processDeviceCredentialsRequestMsg(TenantId tenantId, Edge edge, DeviceCredentialsRequestMsg deviceCredentialsRequestMsg) { |
|
|
public ListenableFuture<Void> processDeviceCredentialsRequestMsg(TenantId tenantId, Edge edge, DeviceCredentialsRequestMsg deviceCredentialsRequestMsg) { |
|
|
log.trace("[{}] processDeviceCredentialsRequestMsg [{}][{}]", tenantId, edge.getName(), deviceCredentialsRequestMsg); |
|
|
log.trace("[{}] processDeviceCredentialsRequestMsg [{}][{}]", tenantId, edge.getName(), deviceCredentialsRequestMsg); |
|
|
SettableFuture<Void> futureToSet = SettableFuture.create(); |
|
|
|
|
|
if (deviceCredentialsRequestMsg.getDeviceIdMSB() != 0 && deviceCredentialsRequestMsg.getDeviceIdLSB() != 0) { |
|
|
if (deviceCredentialsRequestMsg.getDeviceIdMSB() != 0 && deviceCredentialsRequestMsg.getDeviceIdLSB() != 0) { |
|
|
DeviceId deviceId = new DeviceId(new UUID(deviceCredentialsRequestMsg.getDeviceIdMSB(), deviceCredentialsRequestMsg.getDeviceIdLSB())); |
|
|
DeviceId deviceId = new DeviceId(new UUID(deviceCredentialsRequestMsg.getDeviceIdMSB(), deviceCredentialsRequestMsg.getDeviceIdLSB())); |
|
|
ListenableFuture<EdgeEvent> future = saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, |
|
|
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, |
|
|
EdgeEventActionType.CREDENTIALS_UPDATED, deviceId, null); |
|
|
EdgeEventActionType.CREDENTIALS_UPDATED, deviceId, null); |
|
|
Futures.addCallback(future, new FutureCallback<EdgeEvent>() { |
|
|
|
|
|
@Override |
|
|
|
|
|
public void onSuccess(@Nullable EdgeEvent result) { |
|
|
|
|
|
futureToSet.set(null); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
log.error("Can't save edge event [{}]", deviceCredentialsRequestMsg, t); |
|
|
|
|
|
futureToSet.setException(t); |
|
|
|
|
|
} |
|
|
|
|
|
}, dbCallbackExecutorService); |
|
|
|
|
|
} |
|
|
} |
|
|
return futureToSet; |
|
|
return Futures.immediateFuture(null); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public ListenableFuture<Void> processUserCredentialsRequestMsg(TenantId tenantId, Edge edge, UserCredentialsRequestMsg userCredentialsRequestMsg) { |
|
|
public ListenableFuture<Void> processUserCredentialsRequestMsg(TenantId tenantId, Edge edge, UserCredentialsRequestMsg userCredentialsRequestMsg) { |
|
|
log.trace("[{}] processUserCredentialsRequestMsg [{}][{}]", tenantId, edge.getName(), userCredentialsRequestMsg); |
|
|
log.trace("[{}] processUserCredentialsRequestMsg [{}][{}]", tenantId, edge.getName(), userCredentialsRequestMsg); |
|
|
SettableFuture<Void> futureToSet = SettableFuture.create(); |
|
|
|
|
|
if (userCredentialsRequestMsg.getUserIdMSB() != 0 && userCredentialsRequestMsg.getUserIdLSB() != 0) { |
|
|
if (userCredentialsRequestMsg.getUserIdMSB() != 0 && userCredentialsRequestMsg.getUserIdLSB() != 0) { |
|
|
UserId userId = new UserId(new UUID(userCredentialsRequestMsg.getUserIdMSB(), userCredentialsRequestMsg.getUserIdLSB())); |
|
|
UserId userId = new UserId(new UUID(userCredentialsRequestMsg.getUserIdMSB(), userCredentialsRequestMsg.getUserIdLSB())); |
|
|
ListenableFuture<EdgeEvent> future = saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.USER, |
|
|
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.USER, |
|
|
EdgeEventActionType.CREDENTIALS_UPDATED, userId, null); |
|
|
EdgeEventActionType.CREDENTIALS_UPDATED, userId, null); |
|
|
Futures.addCallback(future, new FutureCallback<>() { |
|
|
|
|
|
@Override |
|
|
|
|
|
public void onSuccess(@Nullable EdgeEvent result) { |
|
|
|
|
|
futureToSet.set(null); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
log.error("Can't save edge event [{}]", userCredentialsRequestMsg, t); |
|
|
|
|
|
futureToSet.setException(t); |
|
|
|
|
|
} |
|
|
|
|
|
}, dbCallbackExecutorService); |
|
|
|
|
|
} |
|
|
} |
|
|
return futureToSet; |
|
|
return Futures.immediateFuture(null); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public ListenableFuture<Void> processDeviceProfileDevicesRequestMsg(TenantId tenantId, Edge edge, DeviceProfileDevicesRequestMsg deviceProfileDevicesRequestMsg) { |
|
|
public ListenableFuture<Void> processDeviceProfileDevicesRequestMsg(TenantId tenantId, Edge edge, DeviceProfileDevicesRequestMsg deviceProfileDevicesRequestMsg) { |
|
|
log.trace("[{}] processDeviceProfileDevicesRequestMsg [{}][{}]", tenantId, edge.getName(), deviceProfileDevicesRequestMsg); |
|
|
log.trace("[{}] processDeviceProfileDevicesRequestMsg [{}][{}]", tenantId, edge.getName(), deviceProfileDevicesRequestMsg); |
|
|
SettableFuture<Void> futureToSet = SettableFuture.create(); |
|
|
|
|
|
if (deviceProfileDevicesRequestMsg.getDeviceProfileIdMSB() != 0 && deviceProfileDevicesRequestMsg.getDeviceProfileIdLSB() != 0) { |
|
|
if (deviceProfileDevicesRequestMsg.getDeviceProfileIdMSB() != 0 && deviceProfileDevicesRequestMsg.getDeviceProfileIdLSB() != 0) { |
|
|
DeviceProfileId deviceProfileId = new DeviceProfileId(new UUID(deviceProfileDevicesRequestMsg.getDeviceProfileIdMSB(), deviceProfileDevicesRequestMsg.getDeviceProfileIdLSB())); |
|
|
DeviceProfileId deviceProfileId = new DeviceProfileId(new UUID(deviceProfileDevicesRequestMsg.getDeviceProfileIdMSB(), deviceProfileDevicesRequestMsg.getDeviceProfileIdLSB())); |
|
|
DeviceProfile deviceProfileById = deviceProfileService.findDeviceProfileById(tenantId, deviceProfileId); |
|
|
DeviceProfile deviceProfileById = deviceProfileService.findDeviceProfileById(tenantId, deviceProfileId); |
|
|
List<ListenableFuture<EdgeEvent>> futures; |
|
|
|
|
|
if (deviceProfileById != null) { |
|
|
if (deviceProfileById != null) { |
|
|
futures = syncDevices(tenantId, edge, deviceProfileById.getName()); |
|
|
syncDevices(tenantId, edge, deviceProfileById.getName()); |
|
|
} else { |
|
|
|
|
|
futures = new ArrayList<>(); |
|
|
|
|
|
} |
|
|
} |
|
|
Futures.addCallback(Futures.allAsList(futures), new FutureCallback<>() { |
|
|
|
|
|
@Override |
|
|
|
|
|
public void onSuccess(@Nullable List<EdgeEvent> result) { |
|
|
|
|
|
futureToSet.set(null); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
log.error("Can't sync devices by device profile [{}]", deviceProfileDevicesRequestMsg, t); |
|
|
|
|
|
futureToSet.setException(t); |
|
|
|
|
|
} |
|
|
|
|
|
}, dbCallbackExecutorService); |
|
|
|
|
|
} |
|
|
} |
|
|
return futureToSet; |
|
|
return Futures.immediateFuture(null); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private List<ListenableFuture<EdgeEvent>> syncDevices(TenantId tenantId, Edge edge, String deviceType) { |
|
|
private void syncDevices(TenantId tenantId, Edge edge, String deviceType) { |
|
|
List<ListenableFuture<EdgeEvent>> futures = new ArrayList<>(); |
|
|
|
|
|
log.trace("[{}] syncDevices [{}][{}]", tenantId, edge.getName(), deviceType); |
|
|
log.trace("[{}] syncDevices [{}][{}]", tenantId, edge.getName(), deviceType); |
|
|
try { |
|
|
try { |
|
|
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); |
|
|
PageLink pageLink = new PageLink(DEFAULT_PAGE_SIZE); |
|
|
@ -358,7 +303,7 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { |
|
|
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { |
|
|
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { |
|
|
log.trace("[{}] [{}] device(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); |
|
|
log.trace("[{}] [{}] device(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); |
|
|
for (Device device : pageData.getData()) { |
|
|
for (Device device : pageData.getData()) { |
|
|
futures.add(saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ADDED, device.getId(), null)); |
|
|
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.ADDED, device.getId(), null); |
|
|
} |
|
|
} |
|
|
if (pageData.hasNext()) { |
|
|
if (pageData.hasNext()) { |
|
|
pageLink = pageLink.nextPageLink(); |
|
|
pageLink = pageLink.nextPageLink(); |
|
|
@ -368,40 +313,25 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { |
|
|
} catch (Exception e) { |
|
|
} catch (Exception e) { |
|
|
log.error("Exception during loading edge device(s) on sync!", e); |
|
|
log.error("Exception during loading edge device(s) on sync!", e); |
|
|
} |
|
|
} |
|
|
return futures; |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public ListenableFuture<Void> processWidgetBundleTypesRequestMsg(TenantId tenantId, Edge edge, |
|
|
public ListenableFuture<Void> processWidgetBundleTypesRequestMsg(TenantId tenantId, Edge edge, |
|
|
WidgetBundleTypesRequestMsg widgetBundleTypesRequestMsg) { |
|
|
WidgetBundleTypesRequestMsg widgetBundleTypesRequestMsg) { |
|
|
log.trace("[{}] processWidgetBundleTypesRequestMsg [{}][{}]", tenantId, edge.getName(), widgetBundleTypesRequestMsg); |
|
|
log.trace("[{}] processWidgetBundleTypesRequestMsg [{}][{}]", tenantId, edge.getName(), widgetBundleTypesRequestMsg); |
|
|
SettableFuture<Void> futureToSet = SettableFuture.create(); |
|
|
|
|
|
if (widgetBundleTypesRequestMsg.getWidgetBundleIdMSB() != 0 && widgetBundleTypesRequestMsg.getWidgetBundleIdLSB() != 0) { |
|
|
if (widgetBundleTypesRequestMsg.getWidgetBundleIdMSB() != 0 && widgetBundleTypesRequestMsg.getWidgetBundleIdLSB() != 0) { |
|
|
WidgetsBundleId widgetsBundleId = new WidgetsBundleId(new UUID(widgetBundleTypesRequestMsg.getWidgetBundleIdMSB(), widgetBundleTypesRequestMsg.getWidgetBundleIdLSB())); |
|
|
WidgetsBundleId widgetsBundleId = new WidgetsBundleId(new UUID(widgetBundleTypesRequestMsg.getWidgetBundleIdMSB(), widgetBundleTypesRequestMsg.getWidgetBundleIdLSB())); |
|
|
WidgetsBundle widgetsBundleById = widgetsBundleService.findWidgetsBundleById(tenantId, widgetsBundleId); |
|
|
WidgetsBundle widgetsBundleById = widgetsBundleService.findWidgetsBundleById(tenantId, widgetsBundleId); |
|
|
List<ListenableFuture<EdgeEvent>> futures = new ArrayList<>(); |
|
|
|
|
|
if (widgetsBundleById != null) { |
|
|
if (widgetsBundleById != null) { |
|
|
List<WidgetType> widgetTypesToPush = |
|
|
List<WidgetType> widgetTypesToPush = |
|
|
widgetTypeService.findWidgetTypesByTenantIdAndBundleAlias(widgetsBundleById.getTenantId(), widgetsBundleById.getAlias()); |
|
|
widgetTypeService.findWidgetTypesByTenantIdAndBundleAlias(widgetsBundleById.getTenantId(), widgetsBundleById.getAlias()); |
|
|
|
|
|
|
|
|
for (WidgetType widgetType : widgetTypesToPush) { |
|
|
for (WidgetType widgetType : widgetTypesToPush) { |
|
|
futures.add(saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.WIDGET_TYPE, EdgeEventActionType.ADDED, widgetType.getId(), null)); |
|
|
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.WIDGET_TYPE, EdgeEventActionType.ADDED, widgetType.getId(), null); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
Futures.addCallback(Futures.allAsList(futures), new FutureCallback<>() { |
|
|
|
|
|
@Override |
|
|
|
|
|
public void onSuccess(@Nullable List<EdgeEvent> result) { |
|
|
|
|
|
futureToSet.set(null); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
log.error("Can't sync widget types by widget bundle [{}]", widgetBundleTypesRequestMsg, t); |
|
|
|
|
|
futureToSet.setException(t); |
|
|
|
|
|
} |
|
|
|
|
|
}, dbCallbackExecutorService); |
|
|
|
|
|
} |
|
|
} |
|
|
return futureToSet; |
|
|
return Futures.immediateFuture(null); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
@ -416,9 +346,12 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { |
|
|
public void onSuccess(@Nullable List<EntityView> entityViews) { |
|
|
public void onSuccess(@Nullable List<EntityView> entityViews) { |
|
|
try { |
|
|
try { |
|
|
if (entityViews != null && !entityViews.isEmpty()) { |
|
|
if (entityViews != null && !entityViews.isEmpty()) { |
|
|
|
|
|
List<ListenableFuture<Boolean>> futures = new ArrayList<>(); |
|
|
for (EntityView entityView : entityViews) { |
|
|
for (EntityView entityView : entityViews) { |
|
|
Futures.addCallback(relationService.checkRelation(tenantId, edge.getId(), entityView.getId(), |
|
|
ListenableFuture<Boolean> future = relationService.checkRelation(tenantId, edge.getId(), entityView.getId(), |
|
|
EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE), new FutureCallback<>() { |
|
|
EntityRelation.CONTAINS_TYPE, RelationTypeGroup.EDGE); |
|
|
|
|
|
futures.add(future); |
|
|
|
|
|
Futures.addCallback(future, new FutureCallback<>() { |
|
|
@Override |
|
|
@Override |
|
|
public void onSuccess(@Nullable Boolean result) { |
|
|
public void onSuccess(@Nullable Boolean result) { |
|
|
if (Boolean.TRUE.equals(result)) { |
|
|
if (Boolean.TRUE.equals(result)) { |
|
|
@ -426,16 +359,27 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { |
|
|
EdgeEventActionType.ADDED, entityView.getId(), null); |
|
|
EdgeEventActionType.ADDED, entityView.getId(), null); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public void onFailure(Throwable t) { |
|
|
public void onFailure(Throwable t) { |
|
|
log.error("Exception during loading relation [{}] to edge on sync!", t, t); |
|
|
// Do nothing - error handles in allAsList
|
|
|
futureToSet.setException(t); |
|
|
|
|
|
} |
|
|
} |
|
|
}, dbCallbackExecutorService); |
|
|
}, dbCallbackExecutorService); |
|
|
} |
|
|
} |
|
|
|
|
|
Futures.addCallback(Futures.allAsList(futures), new FutureCallback<>() { |
|
|
|
|
|
@Override |
|
|
|
|
|
public void onSuccess(@Nullable List<Boolean> result) { |
|
|
|
|
|
futureToSet.set(null); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
log.error("Exception during loading relation [{}] to edge on sync!", t, t); |
|
|
|
|
|
futureToSet.setException(t); |
|
|
|
|
|
} |
|
|
|
|
|
}, dbCallbackExecutorService); |
|
|
|
|
|
} else { |
|
|
|
|
|
futureToSet.set(null); |
|
|
} |
|
|
} |
|
|
futureToSet.set(null); |
|
|
|
|
|
} catch (Exception e) { |
|
|
} catch (Exception e) { |
|
|
log.error("Exception during loading relation(s) to edge on sync!", e); |
|
|
log.error("Exception during loading relation(s) to edge on sync!", e); |
|
|
futureToSet.setException(e); |
|
|
futureToSet.setException(e); |
|
|
@ -451,30 +395,19 @@ public class DefaultEdgeRequestsService implements EdgeRequestsService { |
|
|
return futureToSet; |
|
|
return futureToSet; |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private ListenableFuture<EdgeEvent> saveEdgeEvent(TenantId tenantId, |
|
|
private void saveEdgeEvent(TenantId tenantId, |
|
|
EdgeId edgeId, |
|
|
EdgeId edgeId, |
|
|
EdgeEventType type, |
|
|
EdgeEventType type, |
|
|
EdgeEventActionType action, |
|
|
EdgeEventActionType action, |
|
|
EntityId entityId, |
|
|
EntityId entityId, |
|
|
JsonNode body) { |
|
|
JsonNode body) { |
|
|
log.trace("Pushing edge event to edge queue. tenantId [{}], edgeId [{}], type [{}], action[{}], entityId [{}], body [{}]", |
|
|
log.trace("Pushing edge event to edge queue. tenantId [{}], edgeId [{}], type [{}], action[{}], entityId [{}], body [{}]", |
|
|
tenantId, edgeId, type, action, entityId, body); |
|
|
tenantId, edgeId, type, action, entityId, body); |
|
|
|
|
|
|
|
|
EdgeEvent edgeEvent = EdgeEventUtils.constructEdgeEvent(tenantId, edgeId, type, action, entityId, body); |
|
|
EdgeEvent edgeEvent = EdgeEventUtils.constructEdgeEvent(tenantId, edgeId, type, action, entityId, body); |
|
|
|
|
|
|
|
|
ListenableFuture<EdgeEvent> future = edgeEventService.saveAsync(edgeEvent); |
|
|
edgeEventService.save(edgeEvent); |
|
|
Futures.addCallback(future, new FutureCallback<>() { |
|
|
tbClusterService.onEdgeEventUpdate(tenantId, edgeId); |
|
|
@Override |
|
|
|
|
|
public void onSuccess(@Nullable EdgeEvent result) { |
|
|
|
|
|
tbClusterService.onEdgeEventUpdate(tenantId, edgeId); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
log.warn("[{}] Can't save edge event [{}] for edge [{}]", tenantId.getId(), edgeEvent, edgeId.getId(), t); |
|
|
|
|
|
} |
|
|
|
|
|
}, dbCallbackExecutorService); |
|
|
|
|
|
return future; |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
} |
|
|
} |
|
|
|