|
|
@ -316,64 +316,62 @@ public class DefaultSyncEdgeService implements SyncEdgeService { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public void processRuleChainMetadataRequestMsg(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg) { |
|
|
public ListenableFuture<Void> processRuleChainMetadataRequestMsg(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg) { |
|
|
if (ruleChainMetadataRequestMsg.getRuleChainIdMSB() != 0 && ruleChainMetadataRequestMsg.getRuleChainIdLSB() != 0) { |
|
|
if (ruleChainMetadataRequestMsg.getRuleChainIdMSB() != 0 && ruleChainMetadataRequestMsg.getRuleChainIdLSB() != 0) { |
|
|
RuleChainId ruleChainId = new RuleChainId(new UUID(ruleChainMetadataRequestMsg.getRuleChainIdMSB(), ruleChainMetadataRequestMsg.getRuleChainIdLSB())); |
|
|
RuleChainId ruleChainId = new RuleChainId(new UUID(ruleChainMetadataRequestMsg.getRuleChainIdMSB(), ruleChainMetadataRequestMsg.getRuleChainIdLSB())); |
|
|
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.RULE_CHAIN_METADATA, ActionType.ADDED, ruleChainId, null); |
|
|
ListenableFuture<EdgeEvent> future = saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.RULE_CHAIN_METADATA, ActionType.ADDED, ruleChainId, null); |
|
|
|
|
|
return Futures.transform(future, edgeEvent -> null, dbCallbackExecutorService); |
|
|
} |
|
|
} |
|
|
|
|
|
return Futures.immediateFuture(null); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public void processAttributesRequestMsg(Edge edge, AttributesRequestMsg attributesRequestMsg) { |
|
|
public ListenableFuture<Void> processAttributesRequestMsg(Edge edge, AttributesRequestMsg attributesRequestMsg) { |
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid( |
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid( |
|
|
EntityType.valueOf(attributesRequestMsg.getEntityType()), |
|
|
EntityType.valueOf(attributesRequestMsg.getEntityType()), |
|
|
new UUID(attributesRequestMsg.getEntityIdMSB(), attributesRequestMsg.getEntityIdLSB())); |
|
|
new UUID(attributesRequestMsg.getEntityIdMSB(), attributesRequestMsg.getEntityIdLSB())); |
|
|
final EdgeEventType edgeEventType = getEdgeQueueTypeByEntityType(entityId.getEntityType()); |
|
|
final EdgeEventType edgeEventType = getEdgeQueueTypeByEntityType(entityId.getEntityType()); |
|
|
if (edgeEventType != null) { |
|
|
if (edgeEventType != null) { |
|
|
ListenableFuture<List<AttributeKvEntry>> ssAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.SERVER_SCOPE); |
|
|
ListenableFuture<List<AttributeKvEntry>> ssAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.SERVER_SCOPE); |
|
|
Futures.addCallback(ssAttrFuture, new FutureCallback<List<AttributeKvEntry>>() { |
|
|
return Futures.transform(ssAttrFuture, ssAttributes -> { |
|
|
@Override |
|
|
if (ssAttributes != null && !ssAttributes.isEmpty()) { |
|
|
public void onSuccess(@Nullable List<AttributeKvEntry> ssAttributes) { |
|
|
try { |
|
|
if (ssAttributes != null && !ssAttributes.isEmpty()) { |
|
|
Map<String, Object> entityData = new HashMap<>(); |
|
|
try { |
|
|
ObjectNode attributes = mapper.createObjectNode(); |
|
|
Map<String, Object> entityData = new HashMap<>(); |
|
|
for (AttributeKvEntry attr : ssAttributes) { |
|
|
ObjectNode attributes = mapper.createObjectNode(); |
|
|
if (attr.getDataType() == DataType.BOOLEAN && attr.getBooleanValue().isPresent()) { |
|
|
for (AttributeKvEntry attr : ssAttributes) { |
|
|
attributes.put(attr.getKey(), attr.getBooleanValue().get()); |
|
|
if (attr.getDataType() == DataType.BOOLEAN && attr.getBooleanValue().isPresent()) { |
|
|
} else if (attr.getDataType() == DataType.DOUBLE && attr.getDoubleValue().isPresent()) { |
|
|
attributes.put(attr.getKey(), attr.getBooleanValue().get()); |
|
|
attributes.put(attr.getKey(), attr.getDoubleValue().get()); |
|
|
} else if (attr.getDataType() == DataType.DOUBLE && attr.getDoubleValue().isPresent()) { |
|
|
} else if (attr.getDataType() == DataType.LONG && attr.getLongValue().isPresent()) { |
|
|
attributes.put(attr.getKey(), attr.getDoubleValue().get()); |
|
|
attributes.put(attr.getKey(), attr.getLongValue().get()); |
|
|
} else if (attr.getDataType() == DataType.LONG && attr.getLongValue().isPresent()) { |
|
|
} else { |
|
|
attributes.put(attr.getKey(), attr.getLongValue().get()); |
|
|
attributes.put(attr.getKey(), attr.getValueAsString()); |
|
|
} else { |
|
|
|
|
|
attributes.put(attr.getKey(), attr.getValueAsString()); |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
} |
|
|
entityData.put("kv", attributes); |
|
|
|
|
|
entityData.put("scope", DataConstants.SERVER_SCOPE); |
|
|
|
|
|
JsonNode entityBody = mapper.valueToTree(entityData); |
|
|
|
|
|
log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, entityBody); |
|
|
|
|
|
saveEdgeEvent(edge.getTenantId(), |
|
|
|
|
|
edge.getId(), |
|
|
|
|
|
edgeEventType, |
|
|
|
|
|
ActionType.ATTRIBUTES_UPDATED, |
|
|
|
|
|
entityId, |
|
|
|
|
|
entityBody); |
|
|
|
|
|
} catch (Exception e) { |
|
|
|
|
|
log.error("[{}] Failed to send attribute updates to the edge", edge.getName(), e); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
entityData.put("kv", attributes); |
|
|
|
|
|
entityData.put("scope", DataConstants.SERVER_SCOPE); |
|
|
|
|
|
JsonNode entityBody = mapper.valueToTree(entityData); |
|
|
|
|
|
log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, entityBody); |
|
|
|
|
|
saveEdgeEvent(edge.getTenantId(), |
|
|
|
|
|
edge.getId(), |
|
|
|
|
|
edgeEventType, |
|
|
|
|
|
ActionType.ATTRIBUTES_UPDATED, |
|
|
|
|
|
entityId, |
|
|
|
|
|
entityBody); |
|
|
|
|
|
} catch (Exception e) { |
|
|
|
|
|
log.error("[{}] Failed to send attribute updates to the edge", edge.getName(), e); |
|
|
|
|
|
throw new RuntimeException("[" + edge.getName() + "] Failed to send attribute updates to the edge", e); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
return null; |
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
|
|
|
|
|
|
} |
|
|
|
|
|
}, dbCallbackExecutorService); |
|
|
}, dbCallbackExecutorService); |
|
|
|
|
|
|
|
|
// TODO: voba - push shared attributes to edge?
|
|
|
// TODO: voba - push shared attributes to edge?
|
|
|
ListenableFuture<List<AttributeKvEntry>> shAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.SHARED_SCOPE); |
|
|
// ListenableFuture<List<AttributeKvEntry>> shAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.SHARED_SCOPE);
|
|
|
ListenableFuture<List<AttributeKvEntry>> clAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.CLIENT_SCOPE); |
|
|
// ListenableFuture<List<AttributeKvEntry>> clAttrFuture = attributesService.findAll(edge.getTenantId(), entityId, DataConstants.CLIENT_SCOPE);
|
|
|
|
|
|
} else { |
|
|
|
|
|
return Futures.immediateFuture(null); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ -391,7 +389,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public void processRelationRequestMsg(Edge edge, RelationRequestMsg relationRequestMsg) { |
|
|
public ListenableFuture<Void> processRelationRequestMsg(Edge edge, RelationRequestMsg relationRequestMsg) { |
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid( |
|
|
EntityId entityId = EntityIdFactory.getByTypeAndUuid( |
|
|
EntityType.valueOf(relationRequestMsg.getEntityType()), |
|
|
EntityType.valueOf(relationRequestMsg.getEntityType()), |
|
|
new UUID(relationRequestMsg.getEntityIdMSB(), relationRequestMsg.getEntityIdLSB())); |
|
|
new UUID(relationRequestMsg.getEntityIdMSB(), relationRequestMsg.getEntityIdLSB())); |
|
|
@ -400,39 +398,33 @@ public class DefaultSyncEdgeService implements SyncEdgeService { |
|
|
futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.FROM)); |
|
|
futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.FROM)); |
|
|
futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.TO)); |
|
|
futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.TO)); |
|
|
ListenableFuture<List<List<EntityRelation>>> relationsListFuture = Futures.allAsList(futures); |
|
|
ListenableFuture<List<List<EntityRelation>>> relationsListFuture = Futures.allAsList(futures); |
|
|
Futures.addCallback(relationsListFuture, new FutureCallback<List<List<EntityRelation>>>() { |
|
|
return Futures.transform(relationsListFuture, relationsList -> { |
|
|
@Override |
|
|
try { |
|
|
public void onSuccess(@Nullable List<List<EntityRelation>> relationsList) { |
|
|
if (relationsList != null && !relationsList.isEmpty()) { |
|
|
try { |
|
|
for (List<EntityRelation> entityRelations : relationsList) { |
|
|
if (!relationsList.isEmpty()) { |
|
|
log.trace("[{}] [{}] [{}] relation(s) are going to be pushed to edge.", edge.getId(), entityId, entityRelations.size()); |
|
|
for (List<EntityRelation> entityRelations : relationsList) { |
|
|
for (EntityRelation relation : entityRelations) { |
|
|
log.trace("[{}] [{}] [{}] relation(s) are going to be pushed to edge.", edge.getId(), entityId, entityRelations.size()); |
|
|
try { |
|
|
for (EntityRelation relation : entityRelations) { |
|
|
if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) && |
|
|
try { |
|
|
!relation.getTo().getEntityType().equals(EntityType.EDGE)) { |
|
|
if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) && |
|
|
saveEdgeEvent(edge.getTenantId(), |
|
|
!relation.getTo().getEntityType().equals(EntityType.EDGE)) { |
|
|
edge.getId(), |
|
|
saveEdgeEvent(edge.getTenantId(), |
|
|
EdgeEventType.RELATION, |
|
|
edge.getId(), |
|
|
ActionType.ADDED, |
|
|
EdgeEventType.RELATION, |
|
|
null, |
|
|
ActionType.ADDED, |
|
|
mapper.valueToTree(relation)); |
|
|
null, |
|
|
|
|
|
mapper.valueToTree(relation)); |
|
|
|
|
|
} |
|
|
|
|
|
} catch (Exception e) { |
|
|
|
|
|
log.error("Exception during loading relation [{}] to edge on sync!", relation, e); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
} catch (Exception e) { |
|
|
|
|
|
log.error("Exception during loading relation [{}] to edge on sync!", relation, e); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} catch (Exception e) { |
|
|
|
|
|
log.error("Exception during loading relation(s) to edge on sync!", e); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
} catch (Exception e) { |
|
|
|
|
|
log.error("Exception during loading relation(s) to edge on sync!", e); |
|
|
|
|
|
throw new RuntimeException("Exception during loading relation(s) to edge on sync!", e); |
|
|
} |
|
|
} |
|
|
|
|
|
return null; |
|
|
@Override |
|
|
|
|
|
public void onFailure(Throwable t) { |
|
|
|
|
|
log.error("Exception during loading relation(s) to edge on sync!", t); |
|
|
|
|
|
} |
|
|
|
|
|
}, dbCallbackExecutorService); |
|
|
}, dbCallbackExecutorService); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ -443,22 +435,26 @@ public class DefaultSyncEdgeService implements SyncEdgeService { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public void processDeviceCredentialsRequestMsg(Edge edge, DeviceCredentialsRequestMsg deviceCredentialsRequestMsg) { |
|
|
public ListenableFuture<Void> processDeviceCredentialsRequestMsg(Edge edge, DeviceCredentialsRequestMsg deviceCredentialsRequestMsg) { |
|
|
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())); |
|
|
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.DEVICE, ActionType.CREDENTIALS_UPDATED, deviceId, null); |
|
|
ListenableFuture<EdgeEvent> future = saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.DEVICE, ActionType.CREDENTIALS_UPDATED, deviceId, null); |
|
|
|
|
|
return Futures.transform(future, edgeEvent -> null, dbCallbackExecutorService); |
|
|
} |
|
|
} |
|
|
|
|
|
return Futures.immediateFuture(null); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public void processUserCredentialsRequestMsg(Edge edge, UserCredentialsRequestMsg userCredentialsRequestMsg) { |
|
|
public ListenableFuture<Void> processUserCredentialsRequestMsg(Edge edge, UserCredentialsRequestMsg userCredentialsRequestMsg) { |
|
|
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())); |
|
|
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.USER, ActionType.CREDENTIALS_UPDATED, userId, null); |
|
|
ListenableFuture<EdgeEvent> future = saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.USER, ActionType.CREDENTIALS_UPDATED, userId, null); |
|
|
|
|
|
return Futures.transform(future, edgeEvent -> null, dbCallbackExecutorService); |
|
|
} |
|
|
} |
|
|
|
|
|
return Futures.immediateFuture(null); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void saveEdgeEvent(TenantId tenantId, |
|
|
private ListenableFuture<EdgeEvent> saveEdgeEvent(TenantId tenantId, |
|
|
EdgeId edgeId, |
|
|
EdgeId edgeId, |
|
|
EdgeEventType edgeEventType, |
|
|
EdgeEventType edgeEventType, |
|
|
ActionType edgeEventAction, |
|
|
ActionType edgeEventAction, |
|
|
@ -476,6 +472,6 @@ public class DefaultSyncEdgeService implements SyncEdgeService { |
|
|
edgeEvent.setEntityId(entityId.getId()); |
|
|
edgeEvent.setEntityId(entityId.getId()); |
|
|
} |
|
|
} |
|
|
edgeEvent.setEntityBody(entityBody); |
|
|
edgeEvent.setEntityBody(entityBody); |
|
|
edgeEventService.saveAsync(edgeEvent); |
|
|
return edgeEventService.saveAsync(edgeEvent); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|