Browse Source

Added findMissingToRelatedRuleChains method

pull/3957/head
Volodymyr Babak 6 years ago
parent
commit
e58af8fb22
  1. 35
      application/src/main/java/org/thingsboard/server/controller/EdgeController.java
  2. 5
      application/src/main/java/org/thingsboard/server/service/edge/DefaultEdgeNotificationService.java
  3. 3
      application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java
  4. 7
      application/src/main/java/org/thingsboard/server/service/edge/EdgeNotificationService.java
  5. 4
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java
  6. 14
      application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java
  7. 165
      application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java
  8. 13
      application/src/main/java/org/thingsboard/server/service/edge/rpc/init/SyncEdgeService.java
  9. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java
  10. 58
      dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java
  11. 4
      rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java
  12. 4
      ui/src/app/api/edge.service.js
  13. 2
      ui/src/app/edge/edge.directive.js

35
application/src/main/java/org/thingsboard/server/controller/EdgeController.java

@ -389,10 +389,12 @@ public class EdgeController extends BaseController {
checkNotNull(query.getEdgeTypes()); checkNotNull(query.getEdgeTypes());
checkEntityId(query.getParameters().getEntityId(), Operation.READ); checkEntityId(query.getParameters().getEntityId(), Operation.READ);
try { try {
List<Edge> edges = checkNotNull(edgeService.findEdgesByQuery(getCurrentUser().getTenantId(), query).get()); SecurityUser user = getCurrentUser();
TenantId tenantId = user.getTenantId();
List<Edge> edges = checkNotNull(edgeService.findEdgesByQuery(tenantId, query).get());
edges = edges.stream().filter(edge -> { edges = edges.stream().filter(edge -> {
try { try {
accessControlService.checkPermission(getCurrentUser(), Resource.EDGE, Operation.READ, edge.getId(), edge); accessControlService.checkPermission(user, Resource.EDGE, Operation.READ, edge.getId(), edge);
return true; return true;
} catch (ThingsboardException e) { } catch (ThingsboardException e) {
return false; return false;
@ -419,14 +421,18 @@ public class EdgeController extends BaseController {
} }
@PreAuthorize("hasAuthority('TENANT_ADMIN')") @PreAuthorize("hasAuthority('TENANT_ADMIN')")
@RequestMapping(value = "/edge/sync", method = RequestMethod.POST) @RequestMapping(value = "/edge/sync/{edgeId}", method = RequestMethod.POST)
public void syncEdge(@RequestBody EdgeId edgeId) throws ThingsboardException { public void syncEdge(@PathVariable("edgeId") String strEdgeId) throws ThingsboardException {
checkParameter("edgeId", strEdgeId);
try { try {
edgeId = checkNotNull(edgeId);
if (isEdgesEnabled()) { if (isEdgesEnabled()) {
EdgeGrpcSession session = edgeGrpcService.getEdgeGrpcSessionById(edgeId); EdgeId edgeId = new EdgeId(toUUID(strEdgeId));
edgeId = checkNotNull(edgeId);
SecurityUser user = getCurrentUser();
TenantId tenantId = user.getTenantId();
EdgeGrpcSession session = edgeGrpcService.getEdgeGrpcSessionById(tenantId, edgeId);
Edge edge = session.getEdge(); Edge edge = session.getEdge();
syncEdgeService.sync(edge); syncEdgeService.sync(tenantId, edge);
} else { } else {
throw new ThingsboardException("Edges support disabled", ThingsboardErrorCode.GENERAL); throw new ThingsboardException("Edges support disabled", ThingsboardErrorCode.GENERAL);
} }
@ -455,4 +461,19 @@ public class EdgeController extends BaseController {
throw new ThingsboardException(e, ThingsboardErrorCode.SUBSCRIPTION_VIOLATION); throw new ThingsboardException(e, ThingsboardErrorCode.SUBSCRIPTION_VIOLATION);
} }
} }
@PreAuthorize("hasAuthority('TENANT_ADMIN')")
@RequestMapping(value = "/edge/missingToRelatedRuleChains/{edgeId}", method = RequestMethod.GET)
@ResponseBody
public String findMissingToRelatedRuleChains(@PathVariable("edgeId") String strEdgeId) throws ThingsboardException {
try {
EdgeId edgeId = new EdgeId(toUUID(strEdgeId));
edgeId = checkNotNull(edgeId);
SecurityUser user = getCurrentUser();
TenantId tenantId = user.getTenantId();
return edgeService.findMissingToRelatedRuleChains(tenantId, edgeId);
} catch (Exception e) {
throw handleException(e);
}
}
} }

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

@ -113,11 +113,6 @@ public class DefaultEdgeNotificationService implements EdgeNotificationService {
} }
} }
@Override
public TimePageData<EdgeEvent> findEdgeEvents(TenantId tenantId, EdgeId edgeId, TimePageLink pageLink) {
return edgeEventService.findEdgeEvents(tenantId, edgeId, pageLink, true);
}
@Override @Override
public Edge setEdgeRootRuleChain(TenantId tenantId, Edge edge, RuleChainId ruleChainId) throws IOException { public Edge setEdgeRootRuleChain(TenantId tenantId, Edge edge, RuleChainId ruleChainId) throws IOException {
edge.setRootRuleChainId(ruleChainId); edge.setRootRuleChainId(ruleChainId);

3
application/src/main/java/org/thingsboard/server/service/edge/EdgeContextComponent.java

@ -28,6 +28,7 @@ import org.thingsboard.server.dao.customer.CustomerService;
import org.thingsboard.server.dao.dashboard.DashboardService; import org.thingsboard.server.dao.dashboard.DashboardService;
import org.thingsboard.server.dao.device.DeviceCredentialsService; import org.thingsboard.server.dao.device.DeviceCredentialsService;
import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.edge.EdgeEventService;
import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.edge.EdgeService;
import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.entityview.EntityViewService;
import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.relation.RelationService;
@ -74,7 +75,7 @@ public class EdgeContextComponent {
@Lazy @Lazy
@Autowired @Autowired
private EdgeNotificationService edgeNotificationService; private EdgeEventService edgeEventService;
@Lazy @Lazy
@Autowired @Autowired

7
application/src/main/java/org/thingsboard/server/service/edge/EdgeNotificationService.java

@ -15,14 +15,9 @@
*/ */
package org.thingsboard.server.service.edge; package org.thingsboard.server.service.edge;
import org.thingsboard.server.common.data.Event;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TimePageData;
import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos;
@ -30,8 +25,6 @@ import java.io.IOException;
public interface EdgeNotificationService { public interface EdgeNotificationService {
TimePageData<EdgeEvent> findEdgeEvents(TenantId tenantId, EdgeId edgeId, TimePageLink pageLink);
Edge setEdgeRootRuleChain(TenantId tenantId, Edge edge, RuleChainId ruleChainId) throws IOException; Edge setEdgeRootRuleChain(TenantId tenantId, Edge edge, RuleChainId ruleChainId) throws IOException;
void pushNotificationToEdge(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback); void pushNotificationToEdge(TransportProtos.EdgeNotificationMsgProto edgeNotificationMsg, TbCallback callback);

4
application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcService.java

@ -28,6 +28,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -186,11 +187,12 @@ public class EdgeGrpcService extends EdgeRpcServiceGrpc.EdgeRpcServiceImplBase i
scheduleEdgeEventsCheck(edgeGrpcSession); scheduleEdgeEventsCheck(edgeGrpcSession);
} }
public EdgeGrpcSession getEdgeGrpcSessionById(EdgeId edgeId) { public EdgeGrpcSession getEdgeGrpcSessionById(TenantId tenantId, EdgeId edgeId) {
EdgeGrpcSession session = sessions.get(edgeId); EdgeGrpcSession session = sessions.get(edgeId);
if (session != null && session.isConnected()) { if (session != null && session.isConnected()) {
return session; return session;
} else { } else {
log.error("[{}] Edge is not connected [{}]", tenantId, edgeId);
throw new RuntimeException("Edge is not connected"); throw new RuntimeException("Edge is not connected");
} }
} }

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

@ -165,7 +165,7 @@ public final class EdgeGrpcSession implements Closeable {
} }
} }
if (connected && requestMsg.getMsgType().equals(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)) { if (connected && requestMsg.getMsgType().equals(RequestMsgType.SYNC_REQUEST_RPC_MESSAGE)) {
ctx.getSyncEdgeService().sync(edge); ctx.getSyncEdgeService().sync(edge.getTenantId(), edge);
} }
if (connected) { if (connected) {
if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE) && requestMsg.hasUplinkMsg()) { if (requestMsg.getMsgType().equals(RequestMsgType.UPLINK_RPC_MESSAGE) && requestMsg.hasUplinkMsg()) {
@ -267,7 +267,7 @@ public final class EdgeGrpcSession implements Closeable {
UUID ifOffset = null; UUID ifOffset = null;
boolean success = true; boolean success = true;
do { do {
pageData = ctx.getEdgeNotificationService().findEdgeEvents(edge.getTenantId(), edge.getId(), pageLink); pageData = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), pageLink, true);
if (isConnected() && !pageData.getData().isEmpty()) { if (isConnected() && !pageData.getData().isEmpty()) {
log.trace("[{}] [{}] event(s) are going to be processed.", this.sessionId, pageData.getData().size()); log.trace("[{}] [{}] event(s) are going to be processed.", this.sessionId, pageData.getData().size());
List<DownlinkMsg> downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData()); List<DownlinkMsg> downlinkMsgsPack = convertToDownlinkMsgsPack(pageData.getData());
@ -899,27 +899,27 @@ public final class EdgeGrpcSession implements Closeable {
} }
if (uplinkMsg.getRuleChainMetadataRequestMsgList() != null && !uplinkMsg.getRuleChainMetadataRequestMsgList().isEmpty()) { if (uplinkMsg.getRuleChainMetadataRequestMsgList() != null && !uplinkMsg.getRuleChainMetadataRequestMsgList().isEmpty()) {
for (RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg : uplinkMsg.getRuleChainMetadataRequestMsgList()) { for (RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg : uplinkMsg.getRuleChainMetadataRequestMsgList()) {
result.add(ctx.getSyncEdgeService().processRuleChainMetadataRequestMsg(edge, ruleChainMetadataRequestMsg)); result.add(ctx.getSyncEdgeService().processRuleChainMetadataRequestMsg(edge.getTenantId(), edge, ruleChainMetadataRequestMsg));
} }
} }
if (uplinkMsg.getAttributesRequestMsgList() != null && !uplinkMsg.getAttributesRequestMsgList().isEmpty()) { if (uplinkMsg.getAttributesRequestMsgList() != null && !uplinkMsg.getAttributesRequestMsgList().isEmpty()) {
for (AttributesRequestMsg attributesRequestMsg : uplinkMsg.getAttributesRequestMsgList()) { for (AttributesRequestMsg attributesRequestMsg : uplinkMsg.getAttributesRequestMsgList()) {
result.add(ctx.getSyncEdgeService().processAttributesRequestMsg(edge, attributesRequestMsg)); result.add(ctx.getSyncEdgeService().processAttributesRequestMsg(edge.getTenantId(), edge, attributesRequestMsg));
} }
} }
if (uplinkMsg.getRelationRequestMsgList() != null && !uplinkMsg.getRelationRequestMsgList().isEmpty()) { if (uplinkMsg.getRelationRequestMsgList() != null && !uplinkMsg.getRelationRequestMsgList().isEmpty()) {
for (RelationRequestMsg relationRequestMsg : uplinkMsg.getRelationRequestMsgList()) { for (RelationRequestMsg relationRequestMsg : uplinkMsg.getRelationRequestMsgList()) {
result.add(ctx.getSyncEdgeService().processRelationRequestMsg(edge, relationRequestMsg)); result.add(ctx.getSyncEdgeService().processRelationRequestMsg(edge.getTenantId(), edge, relationRequestMsg));
} }
} }
if (uplinkMsg.getUserCredentialsRequestMsgList() != null && !uplinkMsg.getUserCredentialsRequestMsgList().isEmpty()) { if (uplinkMsg.getUserCredentialsRequestMsgList() != null && !uplinkMsg.getUserCredentialsRequestMsgList().isEmpty()) {
for (UserCredentialsRequestMsg userCredentialsRequestMsg : uplinkMsg.getUserCredentialsRequestMsgList()) { for (UserCredentialsRequestMsg userCredentialsRequestMsg : uplinkMsg.getUserCredentialsRequestMsgList()) {
result.add(ctx.getSyncEdgeService().processUserCredentialsRequestMsg(edge, userCredentialsRequestMsg)); result.add(ctx.getSyncEdgeService().processUserCredentialsRequestMsg(edge.getTenantId(), edge, userCredentialsRequestMsg));
} }
} }
if (uplinkMsg.getDeviceCredentialsRequestMsgList() != null && !uplinkMsg.getDeviceCredentialsRequestMsgList().isEmpty()) { if (uplinkMsg.getDeviceCredentialsRequestMsgList() != null && !uplinkMsg.getDeviceCredentialsRequestMsgList().isEmpty()) {
for (DeviceCredentialsRequestMsg deviceCredentialsRequestMsg : uplinkMsg.getDeviceCredentialsRequestMsgList()) { for (DeviceCredentialsRequestMsg deviceCredentialsRequestMsg : uplinkMsg.getDeviceCredentialsRequestMsgList()) {
result.add(ctx.getSyncEdgeService().processDeviceCredentialsRequestMsg(edge, deviceCredentialsRequestMsg)); result.add(ctx.getSyncEdgeService().processDeviceCredentialsRequestMsg(edge.getTenantId(), edge, deviceCredentialsRequestMsg));
} }
} }
if (uplinkMsg.getDeviceRpcCallMsgList() != null && !uplinkMsg.getDeviceRpcCallMsgList().isEmpty()) { if (uplinkMsg.getDeviceRpcCallMsgList() != null && !uplinkMsg.getDeviceRpcCallMsgList().isEmpty()) {

165
application/src/main/java/org/thingsboard/server/service/edge/rpc/init/DefaultSyncEdgeService.java

@ -33,7 +33,6 @@ import org.springframework.core.io.DefaultResourceLoader;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.AdminSettings; import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.DashboardInfo; import org.thingsboard.server.common.data.DashboardInfo;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EdgeUtils; import org.thingsboard.server.common.data.EdgeUtils;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
@ -146,37 +145,37 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
private TbClusterService tbClusterService; private TbClusterService tbClusterService;
@Override @Override
public void sync(Edge edge) { public void sync(TenantId tenantId, Edge edge) {
log.trace("[{}][{}] Staring edge sync process", edge.getTenantId(), edge.getId()); log.trace("[{}][{}] Staring edge sync process", tenantId, edge.getId());
try { try {
syncWidgetsBundleAndWidgetTypes(edge); syncWidgetsBundleAndWidgetTypes(tenantId, edge);
syncAdminSettings(edge); syncAdminSettings(tenantId, edge);
syncRuleChains(edge, new TimePageLink(DEFAULT_LIMIT)); syncRuleChains(tenantId, edge, new TimePageLink(DEFAULT_LIMIT));
syncUsers(edge, new TextPageLink(DEFAULT_LIMIT)); syncUsers(tenantId, edge, new TextPageLink(DEFAULT_LIMIT));
syncDevices(edge, new TimePageLink(DEFAULT_LIMIT)); syncDevices(tenantId, edge, new TimePageLink(DEFAULT_LIMIT));
syncAssets(edge, new TimePageLink(DEFAULT_LIMIT)); syncAssets(tenantId, edge, new TimePageLink(DEFAULT_LIMIT));
syncEntityViews(edge, new TimePageLink(DEFAULT_LIMIT)); syncEntityViews(tenantId, edge, new TimePageLink(DEFAULT_LIMIT));
syncDashboards(edge, new TimePageLink(DEFAULT_LIMIT)); syncDashboards(tenantId, edge, new TimePageLink(DEFAULT_LIMIT));
} catch (Exception e) { } catch (Exception e) {
log.error("[{}][{}] Exception during sync process", edge.getTenantId(), edge.getId(), e); log.error("[{}][{}] Exception during sync process", tenantId, edge.getId(), e);
} }
} }
private void syncRuleChains(Edge edge, TimePageLink pageLink) { private void syncRuleChains(TenantId tenantId, Edge edge, TimePageLink pageLink) {
log.trace("[{}] syncRuleChains [{}] [{}]", edge.getTenantId(), edge.getName(), pageLink); log.trace("[{}] syncRuleChains [{}] [{}]", tenantId, edge.getName(), pageLink);
try { try {
ListenableFuture<TimePageData<RuleChain>> future = ListenableFuture<TimePageData<RuleChain>> future =
ruleChainService.findRuleChainsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink); ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink);
Futures.addCallback(future, new FutureCallback<TimePageData<RuleChain>>() { Futures.addCallback(future, new FutureCallback<TimePageData<RuleChain>>() {
@Override @Override
public void onSuccess(@Nullable TimePageData<RuleChain> pageData) { public void onSuccess(@Nullable TimePageData<RuleChain> pageData) {
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
log.trace("[{}] [{}] rule chains(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); log.trace("[{}] [{}] rule chains(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size());
for (RuleChain ruleChain : pageData.getData()) { for (RuleChain ruleChain : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.RULE_CHAIN, EdgeEventActionType.ADDED, ruleChain.getId(), null); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.RULE_CHAIN, EdgeEventActionType.ADDED, ruleChain.getId(), null);
} }
if (pageData.hasNext()) { if (pageData.hasNext()) {
syncRuleChains(edge, pageData.getNextPageLink()); syncRuleChains(tenantId, edge, pageData.getNextPageLink());
} }
} }
} }
@ -191,21 +190,21 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
} }
private void syncDevices(Edge edge, TimePageLink pageLink) { private void syncDevices(TenantId tenantId, Edge edge, TimePageLink pageLink) {
log.trace("[{}] syncDevices [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncDevices [{}]", tenantId, edge.getName());
try { try {
ListenableFuture<TimePageData<Device>> future = ListenableFuture<TimePageData<Device>> future =
deviceService.findDevicesByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink); deviceService.findDevicesByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink);
Futures.addCallback(future, new FutureCallback<TimePageData<Device>>() { Futures.addCallback(future, new FutureCallback<TimePageData<Device>>() {
@Override @Override
public void onSuccess(@Nullable TimePageData<Device> pageData) { public void onSuccess(@Nullable TimePageData<Device> pageData) {
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()) {
saveEdgeEvent(edge.getTenantId(), 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()) {
syncDevices(edge, pageData.getNextPageLink()); syncDevices(tenantId, edge, pageData.getNextPageLink());
} }
} }
} }
@ -220,20 +219,20 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
} }
private void syncAssets(Edge edge, TimePageLink pageLink) { private void syncAssets(TenantId tenantId, Edge edge, TimePageLink pageLink) {
log.trace("[{}] syncAssets [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncAssets [{}]", tenantId, edge.getName());
try { try {
ListenableFuture<TimePageData<Asset>> future = assetService.findAssetsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink); ListenableFuture<TimePageData<Asset>> future = assetService.findAssetsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink);
Futures.addCallback(future, new FutureCallback<TimePageData<Asset>>() { Futures.addCallback(future, new FutureCallback<TimePageData<Asset>>() {
@Override @Override
public void onSuccess(@Nullable TimePageData<Asset> pageData) { public void onSuccess(@Nullable TimePageData<Asset> pageData) {
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
log.trace("[{}] [{}] asset(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); log.trace("[{}] [{}] asset(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size());
for (Asset asset : pageData.getData()) { for (Asset asset : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ASSET, EdgeEventActionType.ADDED, asset.getId(), null); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.ASSET, EdgeEventActionType.ADDED, asset.getId(), null);
} }
if (pageData.hasNext()) { if (pageData.hasNext()) {
syncAssets(edge, pageData.getNextPageLink()); syncAssets(tenantId, edge, pageData.getNextPageLink());
} }
} }
} }
@ -248,20 +247,20 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
} }
private void syncEntityViews(Edge edge, TimePageLink pageLink) { private void syncEntityViews(TenantId tenantId, Edge edge, TimePageLink pageLink) {
log.trace("[{}] syncEntityViews [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncEntityViews [{}]", tenantId, edge.getName());
try { try {
ListenableFuture<TimePageData<EntityView>> future = entityViewService.findEntityViewsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink); ListenableFuture<TimePageData<EntityView>> future = entityViewService.findEntityViewsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink);
Futures.addCallback(future, new FutureCallback<TimePageData<EntityView>>() { Futures.addCallback(future, new FutureCallback<TimePageData<EntityView>>() {
@Override @Override
public void onSuccess(@Nullable TimePageData<EntityView> pageData) { public void onSuccess(@Nullable TimePageData<EntityView> pageData) {
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
log.trace("[{}] [{}] entity view(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); log.trace("[{}] [{}] entity view(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size());
for (EntityView entityView : pageData.getData()) { for (EntityView entityView : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ENTITY_VIEW, EdgeEventActionType.ADDED, entityView.getId(), null); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.ENTITY_VIEW, EdgeEventActionType.ADDED, entityView.getId(), null);
} }
if (pageData.hasNext()) { if (pageData.hasNext()) {
syncEntityViews(edge, pageData.getNextPageLink()); syncEntityViews(tenantId, edge, pageData.getNextPageLink());
} }
} }
} }
@ -276,20 +275,20 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
} }
private void syncDashboards(Edge edge, TimePageLink pageLink) { private void syncDashboards(TenantId tenantId, Edge edge, TimePageLink pageLink) {
log.trace("[{}] syncDashboards [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncDashboards [{}]", tenantId, edge.getName());
try { try {
ListenableFuture<TimePageData<DashboardInfo>> future = dashboardService.findDashboardsByTenantIdAndEdgeId(edge.getTenantId(), edge.getId(), pageLink); ListenableFuture<TimePageData<DashboardInfo>> future = dashboardService.findDashboardsByTenantIdAndEdgeId(tenantId, edge.getId(), pageLink);
Futures.addCallback(future, new FutureCallback<TimePageData<DashboardInfo>>() { Futures.addCallback(future, new FutureCallback<TimePageData<DashboardInfo>>() {
@Override @Override
public void onSuccess(@Nullable TimePageData<DashboardInfo> pageData) { public void onSuccess(@Nullable TimePageData<DashboardInfo> pageData) {
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
log.trace("[{}] [{}] dashboard(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); log.trace("[{}] [{}] dashboard(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size());
for (DashboardInfo dashboardInfo : pageData.getData()) { for (DashboardInfo dashboardInfo : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.DASHBOARD, EdgeEventActionType.ADDED, dashboardInfo.getId(), null); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DASHBOARD, EdgeEventActionType.ADDED, dashboardInfo.getId(), null);
} }
if (pageData.hasNext()) { if (pageData.hasNext()) {
syncDashboards(edge, pageData.getNextPageLink()); syncDashboards(tenantId, edge, pageData.getNextPageLink());
} }
} }
} }
@ -304,31 +303,31 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
} }
private void syncUsers(Edge edge, TextPageLink pageLink) { private void syncUsers(TenantId tenantId, Edge edge, TextPageLink pageLink) {
log.trace("[{}] syncUsers [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncUsers [{}]", tenantId, edge.getName());
try { try {
TextPageData<User> pageData; TextPageData<User> pageData;
do { do {
pageData = userService.findTenantAdmins(edge.getTenantId(), pageLink); pageData = userService.findTenantAdmins(tenantId, pageLink);
pushUsersToEdge(pageData, edge); pushUsersToEdge(tenantId, pageData, edge);
if (pageData != null && pageData.hasNext()) { if (pageData != null && pageData.hasNext()) {
pageLink = pageData.getNextPageLink(); pageLink = pageData.getNextPageLink();
} }
} while (pageData != null && pageData.hasNext()); } while (pageData != null && pageData.hasNext());
syncCustomerUsers(edge); syncCustomerUsers(tenantId, edge);
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during loading edge user(s) on sync!", e); log.error("Exception during loading edge user(s) on sync!", e);
} }
} }
private void syncCustomerUsers(Edge edge) { private void syncCustomerUsers(TenantId tenantId, Edge edge) {
if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) { if (edge.getCustomerId() != null && !EntityId.NULL_UUID.equals(edge.getCustomerId().getId())) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, edge.getCustomerId(), null); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.CUSTOMER, EdgeEventActionType.ADDED, edge.getCustomerId(), null);
TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT); TextPageLink pageLink = new TextPageLink(DEFAULT_LIMIT);
TextPageData<User> pageData; TextPageData<User> pageData;
do { do {
pageData = userService.findCustomerUsers(edge.getTenantId(), edge.getCustomerId(), pageLink); pageData = userService.findCustomerUsers(tenantId, edge.getCustomerId(), pageLink);
pushUsersToEdge(pageData, edge); pushUsersToEdge(tenantId, pageData, edge);
if (pageData != null && pageData.hasNext()) { if (pageData != null && pageData.hasNext()) {
pageLink = pageData.getNextPageLink(); pageLink = pageData.getNextPageLink();
} }
@ -336,45 +335,45 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
} }
private void pushUsersToEdge(TextPageData<User> pageData, Edge edge) { private void pushUsersToEdge(TenantId tenantId, TextPageData<User> pageData, Edge edge) {
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) { if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
log.trace("[{}] [{}] user(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size()); log.trace("[{}] [{}] user(s) are going to be pushed to edge.", edge.getId(), pageData.getData().size());
for (User user : pageData.getData()) { for (User user : pageData.getData()) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.USER, EdgeEventActionType.ADDED, user.getId(), null); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.USER, EdgeEventActionType.ADDED, user.getId(), null);
} }
} }
} }
private void syncWidgetsBundleAndWidgetTypes(Edge edge) { private void syncWidgetsBundleAndWidgetTypes(TenantId tenantId, Edge edge) {
log.trace("[{}] syncWidgetsBundleAndWidgetTypes [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncWidgetsBundleAndWidgetTypes [{}]", tenantId, edge.getName());
List<WidgetsBundle> widgetsBundlesToPush = new ArrayList<>(); List<WidgetsBundle> widgetsBundlesToPush = new ArrayList<>();
List<WidgetType> widgetTypesToPush = new ArrayList<>(); List<WidgetType> widgetTypesToPush = new ArrayList<>();
widgetsBundlesToPush.addAll(widgetsBundleService.findAllTenantWidgetsBundlesByTenantId(edge.getTenantId())); widgetsBundlesToPush.addAll(widgetsBundleService.findAllTenantWidgetsBundlesByTenantId(tenantId));
widgetsBundlesToPush.addAll(widgetsBundleService.findSystemWidgetsBundles(edge.getTenantId())); widgetsBundlesToPush.addAll(widgetsBundleService.findSystemWidgetsBundles(tenantId));
try { try {
for (WidgetsBundle widgetsBundle: widgetsBundlesToPush) { for (WidgetsBundle widgetsBundle: widgetsBundlesToPush) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.WIDGETS_BUNDLE, EdgeEventActionType.ADDED, widgetsBundle.getId(), null); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.WIDGETS_BUNDLE, EdgeEventActionType.ADDED, widgetsBundle.getId(), null);
widgetTypesToPush.addAll(widgetTypeService.findWidgetTypesByTenantIdAndBundleAlias(widgetsBundle.getTenantId(), widgetsBundle.getAlias())); widgetTypesToPush.addAll(widgetTypeService.findWidgetTypesByTenantIdAndBundleAlias(widgetsBundle.getTenantId(), widgetsBundle.getAlias()));
} }
for (WidgetType widgetType: widgetTypesToPush) { for (WidgetType widgetType: widgetTypesToPush) {
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.WIDGET_TYPE, EdgeEventActionType.ADDED, widgetType.getId(), null); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.WIDGET_TYPE, EdgeEventActionType.ADDED, widgetType.getId(), null);
} }
} catch (Exception e) { } catch (Exception e) {
log.error("Exception during loading widgets bundle(s) and widget type(s) on sync!", e); log.error("Exception during loading widgets bundle(s) and widget type(s) on sync!", e);
} }
} }
private void syncAdminSettings(Edge edge) { private void syncAdminSettings(TenantId tenantId, Edge edge) {
log.trace("[{}] syncAdminSettings [{}]", edge.getTenantId(), edge.getName()); log.trace("[{}] syncAdminSettings [{}]", tenantId, edge.getName());
try { try {
AdminSettings systemMailSettings = adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, "mail"); AdminSettings systemMailSettings = adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, "mail");
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ADMIN_SETTINGS, EdgeEventActionType.UPDATED, null, mapper.valueToTree(systemMailSettings)); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS, EdgeEventActionType.UPDATED, null, mapper.valueToTree(systemMailSettings));
AdminSettings tenantMailSettings = convertToTenantAdminSettings(systemMailSettings.getKey(), (ObjectNode) systemMailSettings.getJsonValue()); AdminSettings tenantMailSettings = convertToTenantAdminSettings(systemMailSettings.getKey(), (ObjectNode) systemMailSettings.getJsonValue());
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ADMIN_SETTINGS, EdgeEventActionType.UPDATED, null, mapper.valueToTree(tenantMailSettings)); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS, EdgeEventActionType.UPDATED, null, mapper.valueToTree(tenantMailSettings));
AdminSettings systemMailTemplates = loadMailTemplates(); AdminSettings systemMailTemplates = loadMailTemplates();
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ADMIN_SETTINGS, EdgeEventActionType.UPDATED, null, mapper.valueToTree(systemMailTemplates)); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS, EdgeEventActionType.UPDATED, null, mapper.valueToTree(systemMailTemplates));
AdminSettings tenantMailTemplates = convertToTenantAdminSettings(systemMailTemplates.getKey(), (ObjectNode) systemMailTemplates.getJsonValue()); AdminSettings tenantMailTemplates = convertToTenantAdminSettings(systemMailTemplates.getKey(), (ObjectNode) systemMailTemplates.getJsonValue());
saveEdgeEvent(edge.getTenantId(), edge.getId(), EdgeEventType.ADMIN_SETTINGS, EdgeEventActionType.UPDATED, null, mapper.valueToTree(tenantMailTemplates)); saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.ADMIN_SETTINGS, EdgeEventActionType.UPDATED, null, mapper.valueToTree(tenantMailTemplates));
} catch (Exception e) { } catch (Exception e) {
log.error("Can't load admin settings", e); log.error("Can't load admin settings", e);
} }
@ -433,13 +432,13 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
@Override @Override
public ListenableFuture<Void> processRuleChainMetadataRequestMsg(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg) { public ListenableFuture<Void> processRuleChainMetadataRequestMsg(TenantId tenantId, Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg) {
log.trace("[{}] processRuleChainMetadataRequestMsg [{}][{}]", edge.getTenantId(), edge.getName(), ruleChainMetadataRequestMsg); log.trace("[{}] processRuleChainMetadataRequestMsg [{}][{}]", tenantId, edge.getName(), ruleChainMetadataRequestMsg);
SettableFuture<Void> futureToSet = SettableFuture.create(); 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(edge.getTenantId(), edge.getId(), EdgeEventType.RULE_CHAIN_METADATA, EdgeEventActionType.ADDED, ruleChainId, null); ListenableFuture<EdgeEvent> future = saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.RULE_CHAIN_METADATA, EdgeEventActionType.ADDED, ruleChainId, null);
Futures.addCallback(future, new FutureCallback<EdgeEvent>() { Futures.addCallback(future, new FutureCallback<EdgeEvent>() {
@Override @Override
public void onSuccess(@Nullable EdgeEvent result) { public void onSuccess(@Nullable EdgeEvent result) {
@ -457,8 +456,8 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
@Override @Override
public ListenableFuture<Void> processAttributesRequestMsg(Edge edge, AttributesRequestMsg attributesRequestMsg) { public ListenableFuture<Void> processAttributesRequestMsg(TenantId tenantId, Edge edge, AttributesRequestMsg attributesRequestMsg) {
log.trace("[{}] processAttributesRequestMsg [{}][{}]", edge.getTenantId(), edge.getName(), attributesRequestMsg); log.trace("[{}] processAttributesRequestMsg [{}][{}]", tenantId, edge.getName(), 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()));
@ -466,7 +465,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
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(edge.getTenantId(), entityId, scope); ListenableFuture<List<AttributeKvEntry>> ssAttrFuture = attributesService.findAll(tenantId, entityId, scope);
Futures.addCallback(ssAttrFuture, new FutureCallback<List<AttributeKvEntry>>() { Futures.addCallback(ssAttrFuture, new FutureCallback<List<AttributeKvEntry>>() {
@Override @Override
public void onSuccess(@Nullable List<AttributeKvEntry> ssAttributes) { public void onSuccess(@Nullable List<AttributeKvEntry> ssAttributes) {
@ -489,7 +488,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
entityData.put("scope", scope); entityData.put("scope", scope);
JsonNode body = mapper.valueToTree(entityData); JsonNode body = mapper.valueToTree(entityData);
log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, body); log.debug("Sending attributes data msg, entityId [{}], attributes [{}]", entityId, body);
saveEdgeEvent(edge.getTenantId(), saveEdgeEvent(tenantId,
edge.getId(), edge.getId(),
type, type,
EdgeEventActionType.ATTRIBUTES_UPDATED, EdgeEventActionType.ATTRIBUTES_UPDATED,
@ -500,7 +499,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
throw new RuntimeException("[" + edge.getName() + "] Failed to send attribute updates to the edge", e); throw new RuntimeException("[" + edge.getName() + "] Failed to send attribute updates to the edge", e);
} }
} else { } else {
log.trace("[{}][{}] No attributes found for entity {} [{}]", edge.getTenantId(), log.trace("[{}][{}] No attributes found for entity {} [{}]", tenantId,
edge.getName(), edge.getName(),
entityId.getEntityType(), entityId.getEntityType(),
entityId.getId()); entityId.getId());
@ -516,21 +515,21 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
}, dbCallbackExecutorService); }, dbCallbackExecutorService);
return futureToSet; return futureToSet;
} else { } else {
log.warn("[{}] Type doesn't supported {}", edge.getTenantId(), entityId.getEntityType()); log.warn("[{}] Type doesn't supported {}", tenantId, entityId.getEntityType());
return Futures.immediateFuture(null); return Futures.immediateFuture(null);
} }
} }
@Override @Override
public ListenableFuture<Void> processRelationRequestMsg(Edge edge, RelationRequestMsg relationRequestMsg) { public ListenableFuture<Void> processRelationRequestMsg(TenantId tenantId, Edge edge, RelationRequestMsg relationRequestMsg) {
log.trace("[{}] processRelationRequestMsg [{}][{}]", edge.getTenantId(), edge.getName(), relationRequestMsg); log.trace("[{}] processRelationRequestMsg [{}][{}]", tenantId, edge.getName(), 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()));
List<ListenableFuture<List<EntityRelation>>> futures = new ArrayList<>(); List<ListenableFuture<List<EntityRelation>>> futures = new ArrayList<>();
futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.FROM)); futures.add(findRelationByQuery(tenantId, edge, entityId, EntitySearchDirection.FROM));
futures.add(findRelationByQuery(edge, entityId, EntitySearchDirection.TO)); futures.add(findRelationByQuery(tenantId, edge, entityId, EntitySearchDirection.TO));
ListenableFuture<List<List<EntityRelation>>> relationsListFuture = Futures.allAsList(futures); ListenableFuture<List<List<EntityRelation>>> relationsListFuture = Futures.allAsList(futures);
SettableFuture<Void> futureToSet = SettableFuture.create(); SettableFuture<Void> futureToSet = SettableFuture.create();
Futures.addCallback(relationsListFuture, new FutureCallback<List<List<EntityRelation>>>() { Futures.addCallback(relationsListFuture, new FutureCallback<List<List<EntityRelation>>>() {
@ -544,7 +543,7 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
try { try {
if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) && if (!relation.getFrom().getEntityType().equals(EntityType.EDGE) &&
!relation.getTo().getEntityType().equals(EntityType.EDGE)) { !relation.getTo().getEntityType().equals(EntityType.EDGE)) {
saveEdgeEvent(edge.getTenantId(), saveEdgeEvent(tenantId,
edge.getId(), edge.getId(),
EdgeEventType.RELATION, EdgeEventType.RELATION,
EdgeEventActionType.ADDED, EdgeEventActionType.ADDED,
@ -568,26 +567,26 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
log.error("[{}] Can't find relation by query. Entity id [{}]", edge.getTenantId(), entityId, t); log.error("[{}] Can't find relation by query. Entity id [{}]", tenantId, entityId, t);
futureToSet.setException(t); futureToSet.setException(t);
} }
}, dbCallbackExecutorService); }, dbCallbackExecutorService);
return futureToSet; return futureToSet;
} }
private ListenableFuture<List<EntityRelation>> findRelationByQuery(Edge edge, EntityId entityId, EntitySearchDirection direction) { private ListenableFuture<List<EntityRelation>> findRelationByQuery(TenantId tenantId, Edge edge, EntityId entityId, EntitySearchDirection direction) {
EntityRelationsQuery query = new EntityRelationsQuery(); EntityRelationsQuery query = new EntityRelationsQuery();
query.setParameters(new RelationsSearchParameters(entityId, direction, -1, false)); query.setParameters(new RelationsSearchParameters(entityId, direction, -1, false));
return relationService.findByQuery(edge.getTenantId(), query); return relationService.findByQuery(tenantId, query);
} }
@Override @Override
public ListenableFuture<Void> processDeviceCredentialsRequestMsg(Edge edge, DeviceCredentialsRequestMsg deviceCredentialsRequestMsg) { public ListenableFuture<Void> processDeviceCredentialsRequestMsg(TenantId tenantId, Edge edge, DeviceCredentialsRequestMsg deviceCredentialsRequestMsg) {
log.trace("[{}] processDeviceCredentialsRequestMsg [{}][{}]", edge.getTenantId(), edge.getName(), deviceCredentialsRequestMsg); log.trace("[{}] processDeviceCredentialsRequestMsg [{}][{}]", tenantId, edge.getName(), deviceCredentialsRequestMsg);
SettableFuture<Void> futureToSet = SettableFuture.create(); 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(edge.getTenantId(), edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_UPDATED, deviceId, null); ListenableFuture<EdgeEvent> future = saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.DEVICE, EdgeEventActionType.CREDENTIALS_UPDATED, deviceId, null);
Futures.addCallback(future, new FutureCallback<EdgeEvent>() { Futures.addCallback(future, new FutureCallback<EdgeEvent>() {
@Override @Override
public void onSuccess(@Nullable EdgeEvent result) { public void onSuccess(@Nullable EdgeEvent result) {
@ -605,12 +604,12 @@ public class DefaultSyncEdgeService implements SyncEdgeService {
} }
@Override @Override
public ListenableFuture<Void> processUserCredentialsRequestMsg(Edge edge, UserCredentialsRequestMsg userCredentialsRequestMsg) { public ListenableFuture<Void> processUserCredentialsRequestMsg(TenantId tenantId, Edge edge, UserCredentialsRequestMsg userCredentialsRequestMsg) {
log.trace("[{}] processUserCredentialsRequestMsg [{}][{}]", edge.getTenantId(), edge.getName(), userCredentialsRequestMsg); log.trace("[{}] processUserCredentialsRequestMsg [{}][{}]", tenantId, edge.getName(), userCredentialsRequestMsg);
SettableFuture<Void> futureToSet = SettableFuture.create(); 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(edge.getTenantId(), edge.getId(), EdgeEventType.USER, EdgeEventActionType.CREDENTIALS_UPDATED, userId, null); ListenableFuture<EdgeEvent> future = saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.USER, EdgeEventActionType.CREDENTIALS_UPDATED, userId, null);
Futures.addCallback(future, new FutureCallback<EdgeEvent>() { Futures.addCallback(future, new FutureCallback<EdgeEvent>() {
@Override @Override
public void onSuccess(@Nullable EdgeEvent result) { public void onSuccess(@Nullable EdgeEvent result) {

13
application/src/main/java/org/thingsboard/server/service/edge/rpc/init/SyncEdgeService.java

@ -17,6 +17,7 @@ package org.thingsboard.server.service.edge.rpc.init;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.gen.edge.AttributesRequestMsg; import org.thingsboard.server.gen.edge.AttributesRequestMsg;
import org.thingsboard.server.gen.edge.DeviceCredentialsRequestMsg; import org.thingsboard.server.gen.edge.DeviceCredentialsRequestMsg;
import org.thingsboard.server.gen.edge.RelationRequestMsg; import org.thingsboard.server.gen.edge.RelationRequestMsg;
@ -25,15 +26,15 @@ import org.thingsboard.server.gen.edge.UserCredentialsRequestMsg;
public interface SyncEdgeService { public interface SyncEdgeService {
void sync(Edge edge); void sync(TenantId tenantId, Edge edge);
ListenableFuture<Void> processRuleChainMetadataRequestMsg(Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg); ListenableFuture<Void> processRuleChainMetadataRequestMsg(TenantId tenantId, Edge edge, RuleChainMetadataRequestMsg ruleChainMetadataRequestMsg);
ListenableFuture<Void> processAttributesRequestMsg(Edge edge, AttributesRequestMsg attributesRequestMsg); ListenableFuture<Void> processAttributesRequestMsg(TenantId tenantId, Edge edge, AttributesRequestMsg attributesRequestMsg);
ListenableFuture<Void> processRelationRequestMsg(Edge edge, RelationRequestMsg relationRequestMsg); ListenableFuture<Void> processRelationRequestMsg(TenantId tenantId, Edge edge, RelationRequestMsg relationRequestMsg);
ListenableFuture<Void> processDeviceCredentialsRequestMsg(Edge edge, DeviceCredentialsRequestMsg deviceCredentialsRequestMsg); ListenableFuture<Void> processDeviceCredentialsRequestMsg(TenantId tenantId, Edge edge, DeviceCredentialsRequestMsg deviceCredentialsRequestMsg);
ListenableFuture<Void> processUserCredentialsRequestMsg(Edge edge, UserCredentialsRequestMsg userCredentialsRequestMsg); ListenableFuture<Void> processUserCredentialsRequestMsg(TenantId tenantId, Edge edge, UserCredentialsRequestMsg userCredentialsRequestMsg);
} }

2
common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java

@ -80,4 +80,6 @@ public interface EdgeService {
Object checkInstance(Object request); Object checkInstance(Object request);
Object activateInstance(String licenseSecret, String releaseDate); Object activateInstance(String licenseSecret, String releaseDate);
String findMissingToRelatedRuleChains(TenantId tenantId, EdgeId edgeId);
} }

58
dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java

@ -15,6 +15,9 @@
*/ */
package org.thingsboard.server.dao.edge; package org.thingsboard.server.dao.edge;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.base.Function; import com.google.common.base.Function;
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;
@ -54,10 +57,13 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.id.UserId;
import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageData;
import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.page.TimePageData;
import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntitySearchDirection; import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChain;
import org.thingsboard.server.common.data.rule.RuleChainConnectionInfo;
import org.thingsboard.server.dao.customer.CustomerDao; import org.thingsboard.server.dao.customer.CustomerDao;
import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.entity.AbstractEntityService;
import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.exception.DataValidationException;
@ -100,6 +106,8 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
public static final String INCORRECT_CUSTOMER_ID = "Incorrect customerId "; public static final String INCORRECT_CUSTOMER_ID = "Incorrect customerId ";
public static final String INCORRECT_EDGE_ID = "Incorrect edgeId "; public static final String INCORRECT_EDGE_ID = "Incorrect edgeId ";
private static final ObjectMapper mapper = new ObjectMapper();
private static final int DEFAULT_LIMIT = 100; private static final int DEFAULT_LIMIT = 100;
private RestTemplate restTemplate; private RestTemplate restTemplate;
@ -575,6 +583,56 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic
return this.restTemplate.postForEntity(EDGE_LICENSE_SERVER_ENDPOINT + "/api/license/activateInstance?licenseSecret={licenseSecret}&releaseDate={releaseDate}", (Object) null, Object.class, params); return this.restTemplate.postForEntity(EDGE_LICENSE_SERVER_ENDPOINT + "/api/license/activateInstance?licenseSecret={licenseSecret}&releaseDate={releaseDate}", (Object) null, Object.class, params);
} }
@Override
public String findMissingToRelatedRuleChains(TenantId tenantId, EdgeId edgeId) {
List<RuleChain> edgeRuleChains = findEdgeRuleChains(tenantId, edgeId);
List<RuleChainId> edgeRuleChainIds = edgeRuleChains.stream().map(IdBased::getId).collect(Collectors.toList());
ObjectNode result = mapper.createObjectNode();
for (RuleChain edgeRuleChain : edgeRuleChains) {
List<RuleChainConnectionInfo> connectionInfos =
ruleChainService.loadRuleChainMetaData(edgeRuleChain.getTenantId(), edgeRuleChain.getId()).getRuleChainConnections();
if (connectionInfos != null && !connectionInfos.isEmpty()) {
List<RuleChainId> connectedRuleChains =
connectionInfos.stream().map(RuleChainConnectionInfo::getTargetRuleChainId).collect(Collectors.toList());
List<String> missingRuleChains = new ArrayList<>();
for (RuleChainId connectedRuleChain : connectedRuleChains) {
if (!edgeRuleChainIds.contains(connectedRuleChain)) {
RuleChain ruleChainById = ruleChainService.findRuleChainById(tenantId, connectedRuleChain);
missingRuleChains.add(ruleChainById.getName());
}
}
if (!missingRuleChains.isEmpty()) {
ArrayNode array = mapper.createArrayNode();
for (String missingRuleChain : missingRuleChains) {
array.add(missingRuleChain);
}
result.set(edgeRuleChain.getName(), array);
}
}
}
return result.toString();
}
private List<RuleChain> findEdgeRuleChains(TenantId tenantId, EdgeId edgeId) {
List<RuleChain> result = new ArrayList<>();
TimePageLink pageLink = new TimePageLink(DEFAULT_LIMIT);
TimePageData<RuleChain> pageData;
try {
do {
pageData = ruleChainService.findRuleChainsByTenantIdAndEdgeId(tenantId, edgeId, pageLink).get();
if (pageData != null && pageData.getData() != null && !pageData.getData().isEmpty()) {
result.addAll(pageData.getData());
if (pageData.hasNext()) {
pageLink = pageData.getNextPageLink();
}
}
} while (pageData != null && pageData.hasNext());
} catch (Exception e) {
log.error("[{}] Can't find edge rule chains [{}]", tenantId, edgeId, e);
}
return result;
}
private void initRestTemplate() { private void initRestTemplate() {
boolean jdkHttpClientEnabled = isNotEmpty(System.getProperty("tb.proxy.jdk")) && System.getProperty("tb.proxy.jdk").equalsIgnoreCase("true"); boolean jdkHttpClientEnabled = isNotEmpty(System.getProperty("tb.proxy.jdk")) && System.getProperty("tb.proxy.jdk").equalsIgnoreCase("true");
boolean systemProxyEnabled = isNotEmpty(System.getProperty("tb.proxy.system")) && System.getProperty("tb.proxy.system").equalsIgnoreCase("true"); boolean systemProxyEnabled = isNotEmpty(System.getProperty("tb.proxy.system")) && System.getProperty("tb.proxy.system").equalsIgnoreCase("true");

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

@ -2384,7 +2384,9 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable {
} }
public void syncEdge(EdgeId edgeId) { public void syncEdge(EdgeId edgeId) {
restTemplate.postForEntity(baseURL + "/api/edge/sync", edgeId, EdgeId.class); Map<String, String> params = new HashMap<>();
params.put("edgeId", edgeId.toString());
restTemplate.postForEntity(baseURL + "/api/edge/sync/{edgeId}", null, EdgeId.class, params);
} }
@Deprecated @Deprecated

4
ui/src/app/api/edge.service.js

@ -297,8 +297,8 @@ function EdgeService($http, $q, customerService) {
function syncEdge(edgeId) { function syncEdge(edgeId) {
var deferred = $q.defer(); var deferred = $q.defer();
var url = '/api/edge/sync'; var url = '/api/edge/sync/' + edgeId;
$http.post(url, edgeId).then(function success(response) { $http.post(url, null).then(function success(response) {
deferred.resolve(response); deferred.resolve(response);
}, function fail(response) { }, function fail(response) {
deferred.reject(response.data); deferred.reject(response.data);

2
ui/src/app/edge/edge.directive.js

@ -70,7 +70,7 @@ export default function EdgeDirective($compile, $templateCache, $translate, $mdD
}; };
scope.onEdgeSync = function (edgeId) { scope.onEdgeSync = function (edgeId) {
edgeService.syncEdge(edgeId).then( edgeService.syncEdge(edgeId.id).then(
function success() { function success() {
toast.showSuccess($translate.instant('edge.sync-message'), 750, angular.element(element).parent().parent(), 'bottom left'); toast.showSuccess($translate.instant('edge.sync-message'), 750, angular.element(element).parent().parent(), 'bottom left');
}, },

Loading…
Cancel
Save