Browse Source

Add latest debug input information to test script dialog. Add device attributes node.

pull/804/head
Igor Kulikov 8 years ago
parent
commit
cbceca66f8
  1. 2
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  2. 34
      application/src/main/java/org/thingsboard/server/controller/RuleChainController.java
  3. 5
      dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java
  4. 15
      dao/src/main/java/org/thingsboard/server/dao/event/CassandraBaseEventDao.java
  5. 12
      dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java
  6. 5
      dao/src/main/java/org/thingsboard/server/dao/event/EventService.java
  7. 16
      dao/src/main/java/org/thingsboard/server/dao/sql/event/EventRepository.java
  8. 11
      dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java
  9. 19
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/DonAsynchron.java
  10. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java
  11. 29
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/data/DeviceRelationsQuery.java
  12. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeConfiguration.java
  13. 103
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java
  14. 62
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java
  15. 52
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java
  16. 48
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java
  17. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java
  18. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java
  19. 56
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoader.java
  20. 6
      rule-engine/rule-engine-components/src/main/resources/public/static/rulenode/rulenode-core-config.js
  21. 14
      ui/src/app/api/rule-chain.service.js
  22. 1
      ui/src/app/rulechain/rulenode-config.directive.js
  23. 1
      ui/src/app/rulechain/rulenode-config.tpl.html
  24. 3
      ui/src/app/rulechain/rulenode-defined-config.directive.js
  25. 1
      ui/src/app/rulechain/rulenode-fieldset.tpl.html
  26. 67
      ui/src/app/rulechain/script/node-script-test.service.js
  27. 2
      ui/src/app/rulechain/script/node-script-test.tpl.html

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

@ -607,7 +607,7 @@ public abstract class BaseController {
} }
public static Exception toException(Throwable error) { public static Exception toException(Throwable error) {
return Exception.class.isInstance(error) ? (Exception) error : new Exception(error); return error != null ? (Exception.class.isInstance(error) ? (Exception) error : new Exception(error)) : null;
} }
} }

34
application/src/main/java/org/thingsboard/server/controller/RuleChainController.java

@ -21,14 +21,18 @@ import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode; import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpStatus; import org.springframework.http.HttpStatus;
import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.security.access.prepost.PreAuthorize;
import org.springframework.util.StringUtils; import org.springframework.util.StringUtils;
import org.springframework.web.bind.annotation.*; import org.springframework.web.bind.annotation.*;
import org.thingsboard.rule.engine.api.ScriptEngine; import org.thingsboard.rule.engine.api.ScriptEngine;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.Event;
import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.RuleNodeId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
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;
@ -38,6 +42,7 @@ import org.thingsboard.server.common.data.rule.RuleChainMetaData;
import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.dao.event.EventService;
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.service.script.NashornJsEngine; import org.thingsboard.server.service.script.NashornJsEngine;
@ -52,9 +57,13 @@ import java.util.Set;
public class RuleChainController extends BaseController { public class RuleChainController extends BaseController {
public static final String RULE_CHAIN_ID = "ruleChainId"; public static final String RULE_CHAIN_ID = "ruleChainId";
public static final String RULE_NODE_ID = "ruleNodeId";
private static final ObjectMapper objectMapper = new ObjectMapper(); private static final ObjectMapper objectMapper = new ObjectMapper();
@Autowired
private EventService eventService;
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')") @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')")
@RequestMapping(value = "/ruleChain/{ruleChainId}", method = RequestMethod.GET) @RequestMapping(value = "/ruleChain/{ruleChainId}", method = RequestMethod.GET)
@ResponseBody @ResponseBody
@ -217,6 +226,31 @@ public class RuleChainController extends BaseController {
} }
} }
@PreAuthorize("hasAnyAuthority('TENANT_ADMIN')")
@RequestMapping(value = "/ruleNode/{ruleNodeId}/debugIn", method = RequestMethod.GET)
@ResponseBody
public JsonNode getLatestRuleNodeDebugInput(@PathVariable(RULE_NODE_ID) String strRuleNodeId) throws ThingsboardException {
checkParameter(RULE_NODE_ID, strRuleNodeId);
try {
RuleNodeId ruleNodeId = new RuleNodeId(toUUID(strRuleNodeId));
TenantId tenantId = getCurrentUser().getTenantId();
List<Event> events = eventService.findLatestEvents(tenantId, ruleNodeId, DataConstants.DEBUG_RULE_NODE, 2);
JsonNode result = null;
if (events != null) {
for (Event event : events) {
JsonNode body = event.getBody();
if (body.has("type") && body.get("type").asText().equals("IN")) {
result = body;
break;
}
}
}
return result;
} catch (Exception e) {
throw handleException(e);
}
}
@PreAuthorize("hasAuthority('TENANT_ADMIN')") @PreAuthorize("hasAuthority('TENANT_ADMIN')")
@RequestMapping(value = "/ruleChain/testScript", method = RequestMethod.POST) @RequestMapping(value = "/ruleChain/testScript", method = RequestMethod.POST)
@ResponseBody @ResponseBody

5
dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java

@ -82,6 +82,11 @@ public class BaseEventService implements EventService {
return new TimePageData<>(events, pageLink); return new TimePageData<>(events, pageLink);
} }
@Override
public List<Event> findLatestEvents(TenantId tenantId, EntityId entityId, String eventType, int limit) {
return eventDao.findLatestEvents(tenantId.getId(), entityId, eventType, limit);
}
private DataValidator<Event> eventValidator = private DataValidator<Event> eventValidator =
new DataValidator<Event>() { new DataValidator<Event>() {
@Override @Override

15
dao/src/main/java/org/thingsboard/server/dao/event/CassandraBaseEventDao.java

@ -134,6 +134,21 @@ public class CassandraBaseEventDao extends CassandraAbstractSearchTimeDao<EventE
return DaoUtil.convertDataList(entities); return DaoUtil.convertDataList(entities);
} }
@Override
public List<Event> findLatestEvents(UUID tenantId, EntityId entityId, String eventType, int limit) {
log.trace("Try to find latest events by tenant [{}], entity [{}], type [{}] and limit [{}]", tenantId, entityId, eventType, limit);
Select select = select().from(EVENT_BY_TYPE_AND_ID_VIEW_NAME);
Select.Where query = select.where();
query.and(eq(ModelConstants.EVENT_TENANT_ID_PROPERTY, tenantId));
query.and(eq(ModelConstants.EVENT_ENTITY_TYPE_PROPERTY, entityId.getEntityType()));
query.and(eq(ModelConstants.EVENT_ENTITY_ID_PROPERTY, entityId.getId()));
query.and(eq(ModelConstants.EVENT_TYPE_PROPERTY, eventType));
query.limit(limit);
query.orderBy(QueryBuilder.desc(ModelConstants.EVENT_TYPE_PROPERTY), QueryBuilder.desc(ModelConstants.ID_PROPERTY));
List<EventEntity> entities = findListByStatement(query);
return DaoUtil.convertDataList(entities);
}
private Optional<Event> save(EventEntity entity, boolean ifNotExists) { private Optional<Event> save(EventEntity entity, boolean ifNotExists) {
if (entity.getId() == null) { if (entity.getId() == null) {
entity.setId(UUIDs.timeBased()); entity.setId(UUIDs.timeBased());

12
dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java

@ -76,4 +76,16 @@ public interface EventDao extends Dao<Event> {
* @return the event list * @return the event list
*/ */
List<Event> findEvents(UUID tenantId, EntityId entityId, String eventType, TimePageLink pageLink); List<Event> findEvents(UUID tenantId, EntityId entityId, String eventType, TimePageLink pageLink);
/**
* Find latest events by tenantId, entityId and eventType.
*
* @param tenantId the tenantId
* @param entityId the entityId
* @param eventType the eventType
* @param limit the limit
* @return the event list
*/
List<Event> findLatestEvents(UUID tenantId, EntityId entityId, String eventType, int limit);
} }

5
dao/src/main/java/org/thingsboard/server/dao/event/EventService.java

@ -21,7 +21,9 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TimePageData; import org.thingsboard.server.common.data.page.TimePageData;
import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.common.data.page.TimePageLink;
import java.util.List;
import java.util.Optional; import java.util.Optional;
import java.util.UUID;
public interface EventService { public interface EventService {
@ -34,4 +36,7 @@ public interface EventService {
TimePageData<Event> findEvents(TenantId tenantId, EntityId entityId, TimePageLink pageLink); TimePageData<Event> findEvents(TenantId tenantId, EntityId entityId, TimePageLink pageLink);
TimePageData<Event> findEvents(TenantId tenantId, EntityId entityId, String eventType, TimePageLink pageLink); TimePageData<Event> findEvents(TenantId tenantId, EntityId entityId, String eventType, TimePageLink pageLink);
List<Event> findLatestEvents(TenantId tenantId, EntityId entityId, String eventType, int limit);
} }

16
dao/src/main/java/org/thingsboard/server/dao/sql/event/EventRepository.java

@ -15,12 +15,18 @@
*/ */
package org.thingsboard.server.dao.sql.event; package org.thingsboard.server.dao.sql.event;
import org.springframework.data.domain.Pageable;
import org.springframework.data.jpa.repository.JpaSpecificationExecutor; import org.springframework.data.jpa.repository.JpaSpecificationExecutor;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.CrudRepository; import org.springframework.data.repository.CrudRepository;
import org.springframework.data.repository.query.Param;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.dao.model.sql.AlarmEntity;
import org.thingsboard.server.dao.model.sql.EventEntity; import org.thingsboard.server.dao.model.sql.EventEntity;
import org.thingsboard.server.dao.util.SqlDao; import org.thingsboard.server.dao.util.SqlDao;
import java.util.List;
/** /**
* Created by Valerii Sosliuk on 5/3/2017. * Created by Valerii Sosliuk on 5/3/2017.
*/ */
@ -36,4 +42,14 @@ public interface EventRepository extends CrudRepository<EventEntity, String>, Jp
EventEntity findByTenantIdAndEntityTypeAndEntityId(String tenantId, EventEntity findByTenantIdAndEntityTypeAndEntityId(String tenantId,
EntityType entityType, EntityType entityType,
String entityId); String entityId);
@Query("SELECT e FROM EventEntity e WHERE e.tenantId = :tenantId AND e.entityType = :entityType " +
"AND e.entityId = :entityId AND e.eventType = :eventType ORDER BY e.eventType DESC, e.id DESC")
List<EventEntity> findLatestByTenantIdAndEntityTypeAndEntityIdAndEventType(
@Param("tenantId") String tenantId,
@Param("entityType") EntityType entityType,
@Param("entityId") String entityId,
@Param("eventType") String eventType,
Pageable pageable);
} }

11
dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java

@ -109,6 +109,17 @@ public class JpaBaseEventDao extends JpaAbstractSearchTimeDao<EventEntity, Event
return DaoUtil.convertDataList(eventRepository.findAll(where(timeSearchSpec).and(fieldsSpec), pageable).getContent()); return DaoUtil.convertDataList(eventRepository.findAll(where(timeSearchSpec).and(fieldsSpec), pageable).getContent());
} }
@Override
public List<Event> findLatestEvents(UUID tenantId, EntityId entityId, String eventType, int limit) {
List<EventEntity> latest = eventRepository.findLatestByTenantIdAndEntityTypeAndEntityIdAndEventType(
UUIDConverter.fromTimeUUID(tenantId),
entityId.getEntityType(),
UUIDConverter.fromTimeUUID(entityId.getId()),
eventType,
new PageRequest(0, limit));
return DaoUtil.convertDataList(latest);
}
public Optional<Event> save(EventEntity entity, boolean ifNotExists) { public Optional<Event> save(EventEntity entity, boolean ifNotExists) {
log.debug("Save event [{}] ", entity); log.debug("Save event [{}] ", entity);
if (entity.getTenantId() == null) { if (entity.getTenantId() == null) {

19
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/DonAsynchron.java

@ -20,12 +20,19 @@ import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import javax.annotation.Nullable; import javax.annotation.Nullable;
import java.util.concurrent.Executor;
import java.util.function.Consumer; import java.util.function.Consumer;
public class DonAsynchron { public class DonAsynchron {
public static <T> void withCallback(ListenableFuture<T> future, Consumer<T> onSuccess, Consumer<Throwable> onFailure) { public static <T> void withCallback(ListenableFuture<T> future, Consumer<T> onSuccess,
Futures.addCallback(future, new FutureCallback<T>() { Consumer<Throwable> onFailure) {
withCallback(future, onSuccess, onFailure, null);
}
public static <T> void withCallback(ListenableFuture<T> future, Consumer<T> onSuccess,
Consumer<Throwable> onFailure, Executor executor) {
FutureCallback<T> callback = new FutureCallback<T>() {
@Override @Override
public void onSuccess(@Nullable T result) { public void onSuccess(@Nullable T result) {
try { try {
@ -33,13 +40,17 @@ public class DonAsynchron {
} catch (Throwable th) { } catch (Throwable th) {
onFailure(th); onFailure(th);
} }
} }
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
onFailure.accept(t); onFailure.accept(t);
} }
}); };
if (executor != null) {
Futures.addCallback(future, callback, executor);
} else {
Futures.addCallback(future, callback);
}
} }
} }

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java

@ -44,7 +44,7 @@ import org.thingsboard.server.common.msg.TbMsg;
"Message metadata can be accessed via <code>metadata</code> property. For example <code>'name = ' + metadata.customerName;</code>", "Message metadata can be accessed via <code>metadata</code> property. For example <code>'name = ' + metadata.customerName;</code>",
uiResources = {"static/rulenode/rulenode-core-config.js"}, uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbActionNodeClearAlarmConfig", configDirective = "tbActionNodeClearAlarmConfig",
icon = "notifications_active" icon = "notifications_off"
) )
public class TbClearAlarmNode extends TbAbstractAlarmNode<TbClearAlarmNodeConfiguration> { public class TbClearAlarmNode extends TbAbstractAlarmNode<TbClearAlarmNodeConfiguration> {

29
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/data/DeviceRelationsQuery.java

@ -0,0 +1,29 @@
/**
* Copyright © 2016-2018 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.rule.engine.data;
import lombok.Data;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import java.util.List;
@Data
public class DeviceRelationsQuery {
private EntitySearchDirection direction;
private int maxLevel = 1;
private String relationType;
private List<String> deviceTypes;
}

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/filter/TbJsSwitchNodeConfiguration.java

@ -32,7 +32,7 @@ public class TbJsSwitchNodeConfiguration implements NodeConfiguration<TbJsSwitch
configuration.setJsScript("function nextRelation(metadata, msg) {\n" + configuration.setJsScript("function nextRelation(metadata, msg) {\n" +
" return ['one','nine'];\n" + " return ['one','nine'];\n" +
"}\n" + "}\n" +
"if(msgType === 'POST_TELEMETRY') {\n" + "if(msgType === 'POST_TELEMETRY_REQUEST') {\n" +
" return ['two'];\n" + " return ['two'];\n" +
"}\n" + "}\n" +
"return nextRelation(metadata, msg);"); "return nextRelation(metadata, msg);");

103
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbAbstractGetAttributesNode.java

@ -0,0 +1,103 @@
/**
* Copyright © 2016-2018 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.rule.engine.metadata;
import com.google.common.base.Function;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import org.apache.commons.collections.CollectionUtils;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNode;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.msg.TbMsg;
import java.util.List;
import static org.thingsboard.rule.engine.DonAsynchron.withCallback;
import static org.thingsboard.rule.engine.api.TbRelationTypes.FAILURE;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.server.common.data.DataConstants.*;
public abstract class TbAbstractGetAttributesNode<C extends TbGetAttributesNodeConfiguration, T extends EntityId> implements TbNode {
protected C config;
@Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = loadGetAttributesNodeConfig(configuration);
}
protected abstract C loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException;
@Override
public void onMsg(TbContext ctx, TbMsg msg) throws TbNodeException {
try {
withCallback(
findEntityAsync(ctx, msg.getOriginator()),
entityId -> safePutAttributes(ctx, msg, entityId),
t -> ctx.tellError(msg, t), ctx.getDbCallbackExecutor());
} catch (Throwable th) {
ctx.tellError(msg, th);
}
}
private void safePutAttributes(TbContext ctx, TbMsg msg, T entityId) {
if(entityId == null || entityId.isNullUid()) {
ctx.tellNext(msg, FAILURE);
return;
}
ListenableFuture<List<Void>> allFutures = Futures.allAsList(
putLatestTelemetry(ctx, entityId, msg, config.getLatestTsKeyNames()),
putAttrAsync(ctx, entityId, msg, CLIENT_SCOPE, config.getClientAttributeNames(), "cs_"),
putAttrAsync(ctx, entityId, msg, SHARED_SCOPE, config.getSharedAttributeNames(), "shared_"),
putAttrAsync(ctx, entityId, msg, SERVER_SCOPE, config.getServerAttributeNames(), "ss_")
);
withCallback(allFutures, i -> ctx.tellNext(msg, SUCCESS), t -> ctx.tellError(msg, t));
}
private ListenableFuture<Void> putAttrAsync(TbContext ctx, EntityId entityId, TbMsg msg, String scope, List<String> keys, String prefix) {
if (CollectionUtils.isEmpty(keys)) {
return Futures.immediateFuture(null);
}
ListenableFuture<List<AttributeKvEntry>> latest = ctx.getAttributesService().find(entityId, scope, keys);
return Futures.transform(latest, (Function<? super List<AttributeKvEntry>, Void>) l -> {
l.forEach(r -> msg.getMetaData().putValue(prefix + r.getKey(), r.getValueAsString()));
return null;
});
}
private ListenableFuture<Void> putLatestTelemetry(TbContext ctx, EntityId entityId, TbMsg msg, List<String> keys) {
if (CollectionUtils.isEmpty(keys)) {
return Futures.immediateFuture(null);
}
ListenableFuture<List<TsKvEntry>> latest = ctx.getTimeseriesService().findLatest(entityId, keys);
return Futures.transform(latest, (Function<? super List<TsKvEntry>, Void>) l -> {
l.forEach(r -> msg.getMetaData().putValue(r.getKey(), r.getValueAsString()));
return null;
});
}
@Override
public void destroy() {
}
protected abstract ListenableFuture<T> findEntityAsync(TbContext ctx, EntityId originator);
}

62
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetAttributesNode.java

@ -15,23 +15,16 @@
*/ */
package org.thingsboard.rule.engine.metadata; package org.thingsboard.rule.engine.metadata;
import com.google.common.base.Function;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
import org.thingsboard.rule.engine.TbNodeUtils; import org.thingsboard.rule.engine.TbNodeUtils;
import org.thingsboard.rule.engine.api.*; import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.data.plugin.ComponentType;
import org.thingsboard.server.common.msg.TbMsg;
import java.util.List;
import static org.thingsboard.rule.engine.DonAsynchron.withCallback;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.server.common.data.DataConstants.*;
/** /**
* Created by ashvayka on 19.01.18. * Created by ashvayka on 19.01.18.
@ -47,50 +40,15 @@ import static org.thingsboard.server.common.data.DataConstants.*;
"<code>metadata.cs_temperature</code> or <code>metadata.shared_limit</code> ", "<code>metadata.cs_temperature</code> or <code>metadata.shared_limit</code> ",
uiResources = {"static/rulenode/rulenode-core-config.js"}, uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeOriginatorAttributesConfig") configDirective = "tbEnrichmentNodeOriginatorAttributesConfig")
public class TbGetAttributesNode implements TbNode { public class TbGetAttributesNode extends TbAbstractGetAttributesNode<TbGetAttributesNodeConfiguration, EntityId> {
private TbGetAttributesNodeConfiguration config;
@Override @Override
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException { protected TbGetAttributesNodeConfiguration loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbGetAttributesNodeConfiguration.class); return TbNodeUtils.convert(configuration, TbGetAttributesNodeConfiguration.class);
} }
@Override @Override
public void onMsg(TbContext ctx, TbMsg msg) throws TbNodeException { protected ListenableFuture<EntityId> findEntityAsync(TbContext ctx, EntityId originator) {
ListenableFuture<List<Void>> allFutures = Futures.allAsList( return Futures.immediateFuture(originator);
putLatestTelemetry(ctx, msg, config.getLatestTsKeyNames()),
putAttrAsync(ctx, msg, CLIENT_SCOPE, config.getClientAttributeNames(), "cs_"),
putAttrAsync(ctx, msg, SHARED_SCOPE, config.getSharedAttributeNames(), "shared_"),
putAttrAsync(ctx, msg, SERVER_SCOPE, config.getServerAttributeNames(), "ss_")
);
withCallback(allFutures, i -> ctx.tellNext(msg, SUCCESS), t -> ctx.tellError(msg, t));
}
private ListenableFuture<Void> putAttrAsync(TbContext ctx, TbMsg msg, String scope, List<String> keys, String prefix) {
if (CollectionUtils.isEmpty(keys)) {
return Futures.immediateFuture(null);
}
ListenableFuture<List<AttributeKvEntry>> latest = ctx.getAttributesService().find(msg.getOriginator(), scope, keys);
return Futures.transform(latest, (Function<? super List<AttributeKvEntry>, Void>) l -> {
l.forEach(r -> msg.getMetaData().putValue(prefix + r.getKey(), r.getValueAsString()));
return null;
});
}
private ListenableFuture<Void> putLatestTelemetry(TbContext ctx, TbMsg msg, List<String> keys) {
if (CollectionUtils.isEmpty(keys)) {
return Futures.immediateFuture(null);
}
ListenableFuture<List<TsKvEntry>> latest = ctx.getTimeseriesService().findLatest(msg.getOriginator(), keys);
return Futures.transform(latest, (Function<? super List<TsKvEntry>, Void>) l -> {
l.forEach(r -> msg.getMetaData().putValue(r.getKey(), r.getValueAsString()));
return null;
});
}
@Override
public void destroy() {
} }
} }

52
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNode.java

@ -0,0 +1,52 @@
/**
* Copyright © 2016-2018 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.rule.engine.metadata;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.rule.engine.TbNodeUtils;
import org.thingsboard.rule.engine.api.RuleNode;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.rule.engine.util.EntitiesRelatedDeviceIdAsyncLoader;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.plugin.ComponentType;
@Slf4j
@RuleNode(type = ComponentType.ENRICHMENT,
name = "device attributes",
configClazz = TbGetDeviceAttrNodeConfiguration.class,
nodeDescription = "Add Originators Related Device Attributes or Latest Telemetry into Message Metadata",
nodeDetails = "If Attributes enrichment configured, <b>CLIENT/SHARED/SERVER</b> attributes are added into Message metadata " +
"with specific prefix: <i>cs/shared/ss</i>. Latest telemetry value added into metadata without prefix. " +
"To access those attributes in other nodes this template can be used " +
"<code>metadata.cs_temperature</code> or <code>metadata.shared_limit</code> ",
uiResources = {"static/rulenode/rulenode-core-config.js"},
configDirective = "tbEnrichmentNodeDeviceAttributesConfig")
public class TbGetDeviceAttrNode extends TbAbstractGetAttributesNode<TbGetDeviceAttrNodeConfiguration, DeviceId> {
@Override
protected TbGetDeviceAttrNodeConfiguration loadGetAttributesNodeConfig(TbNodeConfiguration configuration) throws TbNodeException {
return TbNodeUtils.convert(configuration, TbGetDeviceAttrNodeConfiguration.class);
}
@Override
protected ListenableFuture<DeviceId> findEntityAsync(TbContext ctx, EntityId originator) {
return EntitiesRelatedDeviceIdAsyncLoader.findDeviceAsync(ctx, originator, config.getDeviceRelationsQuery());
}
}

48
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetDeviceAttrNodeConfiguration.java

@ -0,0 +1,48 @@
/**
* Copyright © 2016-2018 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.rule.engine.metadata;
import lombok.Data;
import org.thingsboard.rule.engine.data.DeviceRelationsQuery;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import java.util.Collections;
@Data
public class TbGetDeviceAttrNodeConfiguration extends TbGetAttributesNodeConfiguration {
private DeviceRelationsQuery deviceRelationsQuery;
@Override
public TbGetDeviceAttrNodeConfiguration defaultConfiguration() {
TbGetDeviceAttrNodeConfiguration configuration = new TbGetDeviceAttrNodeConfiguration();
configuration.setClientAttributeNames(Collections.emptyList());
configuration.setSharedAttributeNames(Collections.emptyList());
configuration.setServerAttributeNames(Collections.emptyList());
configuration.setLatestTsKeyNames(Collections.emptyList());
DeviceRelationsQuery deviceRelationsQuery = new DeviceRelationsQuery();
deviceRelationsQuery.setDirection(EntitySearchDirection.FROM);
deviceRelationsQuery.setMaxLevel(1);
deviceRelationsQuery.setRelationType(EntityRelation.CONTAINS_TYPE);
deviceRelationsQuery.setDeviceTypes(Collections.singletonList("default"));
configuration.setDeviceRelationsQuery(deviceRelationsQuery);
return configuration;
}
}

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetEntityAttrNodeConfiguration.java

@ -34,7 +34,7 @@ public class TbGetEntityAttrNodeConfiguration implements NodeConfiguration<TbGet
Map<String, String> attrMapping = new HashMap<>(); Map<String, String> attrMapping = new HashMap<>();
attrMapping.putIfAbsent("temperature", "tempo"); attrMapping.putIfAbsent("temperature", "tempo");
configuration.setAttrMapping(attrMapping); configuration.setAttrMapping(attrMapping);
configuration.setTelemetry(true); configuration.setTelemetry(false);
return configuration; return configuration;
} }
} }

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetRelatedAttrNodeConfiguration.java

@ -36,7 +36,7 @@ public class TbGetRelatedAttrNodeConfiguration extends TbGetEntityAttrNodeConfig
Map<String, String> attrMapping = new HashMap<>(); Map<String, String> attrMapping = new HashMap<>();
attrMapping.putIfAbsent("temperature", "tempo"); attrMapping.putIfAbsent("temperature", "tempo");
configuration.setAttrMapping(attrMapping); configuration.setAttrMapping(attrMapping);
configuration.setTelemetry(true); configuration.setTelemetry(false);
RelationsQuery relationsQuery = new RelationsQuery(); RelationsQuery relationsQuery = new RelationsQuery();
relationsQuery.setDirection(EntitySearchDirection.FROM); relationsQuery.setDirection(EntitySearchDirection.FROM);

56
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/EntitiesRelatedDeviceIdAsyncLoader.java

@ -0,0 +1,56 @@
/**
* Copyright © 2016-2018 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.rule.engine.util;
import com.google.common.util.concurrent.AsyncFunction;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import org.apache.commons.collections.CollectionUtils;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.data.DeviceRelationsQuery;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.device.DeviceSearchQuery;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.relation.RelationsSearchParameters;
import org.thingsboard.server.dao.device.DeviceService;
import java.util.List;
public class EntitiesRelatedDeviceIdAsyncLoader {
public static ListenableFuture<DeviceId> findDeviceAsync(TbContext ctx, EntityId originator,
DeviceRelationsQuery deviceRelationsQuery) {
DeviceService deviceService = ctx.getDeviceService();
DeviceSearchQuery query = buildQuery(originator, deviceRelationsQuery);
ListenableFuture<List<Device>> asyncDevices = deviceService.findDevicesByQuery(query);
return Futures.transform(asyncDevices, (AsyncFunction<List<Device>, DeviceId>)
d -> CollectionUtils.isNotEmpty(d) ? Futures.immediateFuture(d.get(0).getId())
: Futures.immediateFuture(null));
}
private static DeviceSearchQuery buildQuery(EntityId originator, DeviceRelationsQuery deviceRelationsQuery) {
DeviceSearchQuery query = new DeviceSearchQuery();
RelationsSearchParameters parameters = new RelationsSearchParameters(originator,
deviceRelationsQuery.getDirection(), deviceRelationsQuery.getMaxLevel());
query.setParameters(parameters);
query.setRelationType(deviceRelationsQuery.getRelationType());
query.setDeviceTypes(deviceRelationsQuery.getDeviceTypes());
return query;
}
}

6
rule-engine/rule-engine-components/src/main/resources/public/static/rulenode/rulenode-core-config.js

File diff suppressed because one or more lines are too long

14
ui/src/app/api/rule-chain.service.js

@ -34,7 +34,8 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co
getRuleNodeComponentByClazz: getRuleNodeComponentByClazz, getRuleNodeComponentByClazz: getRuleNodeComponentByClazz,
getRuleNodeSupportedLinks: getRuleNodeSupportedLinks, getRuleNodeSupportedLinks: getRuleNodeSupportedLinks,
resolveTargetRuleChains: resolveTargetRuleChains, resolveTargetRuleChains: resolveTargetRuleChains,
testScript: testScript testScript: testScript,
getLatestRuleNodeDebugInput: getLatestRuleNodeDebugInput
}; };
return service; return service;
@ -313,4 +314,15 @@ function RuleChainService($http, $q, $filter, $ocLazyLoad, $translate, types, co
return deferred.promise; return deferred.promise;
} }
function getLatestRuleNodeDebugInput(ruleNodeId) {
var deferred = $q.defer();
var url = '/api/ruleNode/' + ruleNodeId + '/debugIn';
$http.get(url).then(function success(response) {
deferred.resolve(response.data);
}, function fail() {
deferred.reject();
});
return deferred.promise;
}
} }

1
ui/src/app/rulechain/rulenode-config.directive.js

@ -68,6 +68,7 @@ export default function RuleNodeConfigDirective($compile, $templateCache, $injec
restrict: "E", restrict: "E",
require: "^ngModel", require: "^ngModel",
scope: { scope: {
ruleNodeId:'=',
nodeDefinition:'=', nodeDefinition:'=',
required:'=ngRequired', required:'=ngRequired',
readonly:'=ngReadonly' readonly:'=ngReadonly'

1
ui/src/app/rulechain/rulenode-config.tpl.html

@ -17,6 +17,7 @@
--> -->
<tb-rule-node-defined-config ng-if="useDefinedDirective()" <tb-rule-node-defined-config ng-if="useDefinedDirective()"
rule-node-id="ruleNodeId"
ng-model="configuration" ng-model="configuration"
rule-node-directive="{{nodeDefinition.configDirective}}" rule-node-directive="{{nodeDefinition.configDirective}}"
ng-required="required" ng-required="required"

3
ui/src/app/rulechain/rulenode-defined-config.directive.js

@ -40,7 +40,7 @@ export default function RuleNodeDefinedConfigDirective($compile) {
scope.ruleNodeConfigScope.$destroy(); scope.ruleNodeConfigScope.$destroy();
} }
var directive = snake_case(attrs.ruleNodeDirective, '-'); var directive = snake_case(attrs.ruleNodeDirective, '-');
var template = `<${directive} ng-model="configuration" ng-required="required" ng-readonly="readonly"></${directive}>`; var template = `<${directive} rule-node-id="ruleNodeId" ng-model="configuration" ng-required="required" ng-readonly="readonly"></${directive}>`;
element.html(template); element.html(template);
scope.ruleNodeConfigScope = scope.$new(); scope.ruleNodeConfigScope = scope.$new();
$compile(element.contents())(scope.ruleNodeConfigScope); $compile(element.contents())(scope.ruleNodeConfigScope);
@ -58,6 +58,7 @@ export default function RuleNodeDefinedConfigDirective($compile) {
restrict: "E", restrict: "E",
require: "^ngModel", require: "^ngModel",
scope: { scope: {
ruleNodeId:'=',
required:'=ngRequired', required:'=ngRequired',
readonly:'=ngReadonly' readonly:'=ngReadonly'
}, },

1
ui/src/app/rulechain/rulenode-fieldset.tpl.html

@ -37,6 +37,7 @@
</md-input-container> </md-input-container>
</section> </section>
<tb-rule-node-config ng-model="ruleNode.configuration" <tb-rule-node-config ng-model="ruleNode.configuration"
rule-node-id="ruleNode.ruleNodeId.id"
ng-required="true" ng-required="true"
node-definition="ruleNode.component.configurationDescriptor.nodeDefinition" node-definition="ruleNode.component.configurationDescriptor.nodeDefinition"
ng-readonly="$root.loading || !isEdit || isReadOnly"> ng-readonly="$root.loading || !isEdit || isReadOnly">

67
ui/src/app/rulechain/script/node-script-test.service.js

@ -21,7 +21,7 @@ import nodeScriptTestTemplate from './node-script-test.tpl.html';
/* eslint-enable import/no-unresolved, import/default */ /* eslint-enable import/no-unresolved, import/default */
/*@ngInject*/ /*@ngInject*/
export default function NodeScriptTest($q, $mdDialog, $document) { export default function NodeScriptTest($q, $mdDialog, $document, ruleChainService) {
var service = { var service = {
testNodeScript: testNodeScript testNodeScript: testNodeScript
@ -29,12 +29,72 @@ export default function NodeScriptTest($q, $mdDialog, $document) {
return service; return service;
function testNodeScript($event, script, scriptType, functionTitle, functionName, argNames, msg, metadata, msgType) { function testNodeScript($event, script, scriptType, functionTitle, functionName, argNames, ruleNodeId) {
var deferred = $q.defer(); var deferred = $q.defer();
if ($event) { if ($event) {
$event.stopPropagation(); $event.stopPropagation();
} }
var msg, metadata, msgType;
if (ruleNodeId) {
ruleChainService.getLatestRuleNodeDebugInput(ruleNodeId).then(
(debugIn) => {
if (debugIn) {
if (debugIn.data) {
msg = angular.fromJson(debugIn.data);
}
if (debugIn.metadata) {
metadata = angular.fromJson(debugIn.metadata);
}
msgType = debugIn.msgType;
}
openTestScriptDialog($event, script, scriptType, functionTitle,
functionName, argNames, msg, metadata, msgType).then(
(script) => {
deferred.resolve(script);
},
() => {
deferred.reject();
}
);
},
() => {
deferred.reject();
}
);
} else {
openTestScriptDialog($event, script, scriptType, functionTitle,
functionName, argNames).then(
(script) => {
deferred.resolve(script);
},
() => {
deferred.reject();
}
);
}
return deferred.promise;
}
function openTestScriptDialog($event, script, scriptType, functionTitle, functionName, argNames, msg, metadata, msgType) {
var deferred = $q.defer();
if (!msg) {
msg = {
temperature: 22.4,
humidity: 78
};
}
if (!metadata) {
metadata = {
deviceType: "default",
deviceName: "Test Device",
ts: new Date().getTime() + ""
};
}
if (!msgType) {
msgType = "POST_TELEMETRY_REQUEST";
}
var onShowingCallback = { var onShowingCallback = {
onShowed: () => { onShowed: () => {
} }
@ -74,7 +134,6 @@ export default function NodeScriptTest($q, $mdDialog, $document) {
deferred.reject(); deferred.reject();
} }
); );
return deferred.promise; return deferred.promise;
} }

2
ui/src/app/rulechain/script/node-script-test.tpl.html

@ -38,7 +38,7 @@
<ng-form name="payloadForm"> <ng-form name="payloadForm">
<div layout="column" style="height: 100%;"> <div layout="column" style="height: 100%;">
<div layout="row"> <div layout="row">
<md-input-container class="md-block" style="margin-bottom: 0px; min-width: 200px;"> <md-input-container class="md-block" style="margin-bottom: 0px; min-width: 300px;">
<label translate>rulenode.message-type</label> <label translate>rulenode.message-type</label>
<input required name="msgType" ng-model="vm.inputParams.msgType"> <input required name="msgType" ng-model="vm.inputParams.msgType">
<div ng-messages="payloadForm.msgType.$error"> <div ng-messages="payloadForm.msgType.$error">

Loading…
Cancel
Save