Browse Source

Merge branch 'master' of github.com:thingsboard/thingsboard

pull/1359/head
nordmif 8 years ago
parent
commit
5a77d29be4
  1. 7
      application/src/main/data/upgrade/2.1.1/schema_update.cql
  2. 4
      application/src/main/data/upgrade/2.1.1/schema_update.sql
  3. 59
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  4. 4
      dao/src/main/resources/sql/schema-entities.sql
  5. 4
      dao/src/main/resources/sql/schema-ts.sql
  6. 9
      msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java
  7. 66
      msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java
  8. 108
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java
  9. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java
  10. 8
      rule-engine/rule-engine-components/src/main/resources/public/static/rulenode/rulenode-core-config.js
  11. 1
      tools/pom.xml
  12. 53
      tools/src/main/java/org/thingsboard/client/tools/RestClient.java
  13. 45
      ui/src/app/extension/extensions-forms/extension-form-modbus.directive.js
  14. 104
      ui/src/app/extension/extensions-forms/extension-form-modbus.tpl.html
  15. 2
      ui/src/app/locale/locale.constant-en_US.json
  16. 1528
      ui/src/app/locale/locale.constant-ru_RU.json

7
application/src/main/data/upgrade/2.1.1/schema_update.cql

@ -14,13 +14,6 @@
-- limitations under the License.
--
DROP MATERIALIZED VIEW IF EXISTS thingsboard.entity_view_by_tenant_and_name;
DROP MATERIALIZED VIEW IF EXISTS thingsboard.entity_view_by_tenant_and_search_text;
DROP MATERIALIZED VIEW IF EXISTS thingsboard.entity_view_by_tenant_and_customer;
DROP MATERIALIZED VIEW IF EXISTS thingsboard.entity_view_by_tenant_and_entity_id;
DROP TABLE IF EXISTS thingsboard.entity_views;
CREATE TABLE IF NOT EXISTS thingsboard.entity_views (
id timeuuid,
entity_id timeuuid,

4
application/src/main/data/upgrade/2.1.1/schema_update.sql

@ -14,10 +14,8 @@
-- limitations under the License.
--
DROP TABLE IF EXISTS entity_views;
CREATE TABLE IF NOT EXISTS entity_views (
id varchar(31) NOT NULL CONSTRAINT entity_view_pkey PRIMARY KEY,
id varchar(31) NOT NULL CONSTRAINT entity_views_pkey PRIMARY KEY,
entity_id varchar(31),
entity_type varchar(255),
tenant_id varchar(31),

59
application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java

@ -16,7 +16,6 @@
package org.thingsboard.server.actors.device;
import akka.actor.ActorContext;
import akka.event.LoggingAdapter;
import com.datastax.driver.core.utils.UUIDs;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
@ -26,12 +25,12 @@ import com.google.gson.JsonObject;
import com.google.gson.JsonParser;
import com.google.protobuf.InvalidProtocolBufferException;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
import org.thingsboard.rule.engine.api.RpcError;
import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg;
import org.thingsboard.rule.engine.api.msg.DeviceNameOrTypeUpdateMsg;
import org.thingsboard.server.actors.ActorSystemContext;
import org.thingsboard.server.actors.shared.AbstractContextAwareMsgProcessor;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.TenantId;
@ -42,7 +41,6 @@ import org.thingsboard.server.common.data.rpc.ToDeviceRpcRequestBody;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgDataType;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.cluster.ClusterEventMsg;
import org.thingsboard.server.common.msg.rpc.ToDeviceRpcRequest;
import org.thingsboard.server.common.msg.session.SessionMsgType;
import org.thingsboard.server.common.msg.timeout.DeviceActorClientSideRpcTimeoutMsg;
@ -81,12 +79,14 @@ import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.UUID;
import java.util.function.Consumer;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.DataConstants.CLIENT_SCOPE;
import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE;
/**
* @author Andrew Shvayka
*/
@ -263,10 +263,8 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
}
private void handleGetAttributesRequest(ActorContext context, SessionInfoProto sessionInfo, GetAttributeRequestMsg request) {
ListenableFuture<List<AttributeKvEntry>> clientAttributesFuture = getAttributeKvEntries(deviceId, DataConstants.CLIENT_SCOPE, toOptionalSet(request.getClientAttributeNamesList()));
ListenableFuture<List<AttributeKvEntry>> sharedAttributesFuture = getAttributeKvEntries(deviceId, DataConstants.SHARED_SCOPE, toOptionalSet(request.getSharedAttributeNamesList()));
int requestId = request.getRequestId();
Futures.addCallback(Futures.allAsList(Arrays.asList(clientAttributesFuture, sharedAttributesFuture)), new FutureCallback<List<List<AttributeKvEntry>>>() {
Futures.addCallback(getAttributesKvEntries(request), new FutureCallback<List<List<AttributeKvEntry>>>() {
@Override
public void onSuccess(@Nullable List<List<AttributeKvEntry>> result) {
GetAttributeResponseMsg responseMsg = GetAttributeResponseMsg.newBuilder()
@ -287,16 +285,35 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
});
}
private ListenableFuture<List<AttributeKvEntry>> getAttributeKvEntries(DeviceId deviceId, String scope, Optional<Set<String>> names) {
if (names.isPresent()) {
if (!names.get().isEmpty()) {
return systemContext.getAttributesService().find(tenantId, deviceId, scope, names.get());
} else {
return systemContext.getAttributesService().findAll(tenantId, deviceId, scope);
}
private ListenableFuture<List<List<AttributeKvEntry>>> getAttributesKvEntries(GetAttributeRequestMsg request) {
ListenableFuture<List<AttributeKvEntry>> clientAttributesFuture;
ListenableFuture<List<AttributeKvEntry>> sharedAttributesFuture;
if (CollectionUtils.isEmpty(request.getClientAttributeNamesList()) && CollectionUtils.isEmpty(request.getSharedAttributeNamesList())) {
clientAttributesFuture = findAllAttributesByScope(CLIENT_SCOPE);
sharedAttributesFuture = findAllAttributesByScope(SHARED_SCOPE);
} else if (!CollectionUtils.isEmpty(request.getClientAttributeNamesList()) && !CollectionUtils.isEmpty(request.getSharedAttributeNamesList())) {
clientAttributesFuture = findAttributesByScope(toSet(request.getClientAttributeNamesList()), CLIENT_SCOPE);
sharedAttributesFuture = findAttributesByScope(toSet(request.getSharedAttributeNamesList()), SHARED_SCOPE);
} else if (CollectionUtils.isEmpty(request.getClientAttributeNamesList()) && !CollectionUtils.isEmpty(request.getSharedAttributeNamesList())) {
clientAttributesFuture = Futures.immediateFuture(Collections.emptyList());
sharedAttributesFuture = findAttributesByScope(toSet(request.getSharedAttributeNamesList()), SHARED_SCOPE);
} else {
return Futures.immediateFuture(Collections.emptyList());
sharedAttributesFuture = Futures.immediateFuture(Collections.emptyList());
clientAttributesFuture = findAttributesByScope(toSet(request.getClientAttributeNamesList()), CLIENT_SCOPE);
}
return Futures.allAsList(Arrays.asList(clientAttributesFuture, sharedAttributesFuture));
}
private ListenableFuture<List<AttributeKvEntry>> findAllAttributesByScope(String scope) {
return systemContext.getAttributesService().findAll(tenantId, deviceId, scope);
}
private ListenableFuture<List<AttributeKvEntry>> findAttributesByScope(Set<String> attributesSet, String scope) {
return systemContext.getAttributesService().find(tenantId, deviceId, scope, attributesSet);
}
private Set<String> toSet(List<String> strings) {
return new HashSet<>(strings);
}
private void handlePostAttributesRequest(ActorContext context, SessionInfoProto sessionInfo, PostAttributeMsg postAttributes) {
@ -368,7 +385,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
AttributeUpdateNotificationMsg.Builder notification = AttributeUpdateNotificationMsg.newBuilder();
if (msg.isDeleted()) {
List<String> sharedKeys = msg.getDeletedKeys().stream()
.filter(key -> DataConstants.SHARED_SCOPE.equals(key.getScope()))
.filter(key -> SHARED_SCOPE.equals(key.getScope()))
.map(AttributeKey::getAttributeKey)
.collect(Collectors.toList());
if (!sharedKeys.isEmpty()) {
@ -376,7 +393,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
hasNotificationData = true;
}
} else {
if (DataConstants.SHARED_SCOPE.equals(msg.getScope())) {
if (SHARED_SCOPE.equals(msg.getScope())) {
List<AttributeKvEntry> attributes = new ArrayList<>(msg.getValues());
if (attributes.size() > 0) {
List<TsKvProto> sharedUpdated = msg.getValues().stream().map(this::toTsKvProto)
@ -545,14 +562,6 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
return json;
}
private Optional<Set<String>> toOptionalSet(List<String> strings) {
if (strings == null || strings.isEmpty()) {
return Optional.empty();
} else {
return Optional.of(new HashSet<>(strings));
}
}
private void sendToTransport(GetAttributeResponseMsg responseMsg, SessionInfoProto sessionInfo) {
DeviceActorToTransportMsg msg = DeviceActorToTransportMsg.newBuilder()
.setSessionIdMSB(sessionInfo.getSessionIdMSB())

4
dao/src/main/resources/sql/schema-entities.sql

@ -72,7 +72,7 @@ CREATE TABLE IF NOT EXISTS attribute_kv (
long_v bigint,
dbl_v double precision,
last_update_ts bigint,
CONSTRAINT attribute_kv_unq_key UNIQUE (entity_type, entity_id, attribute_type, attribute_key)
CONSTRAINT attribute_kv_pkey PRIMARY KEY (entity_type, entity_id, attribute_type, attribute_key)
);
CREATE TABLE IF NOT EXISTS component_descriptor (
@ -148,7 +148,7 @@ CREATE TABLE IF NOT EXISTS relation (
relation_type_group varchar(255),
relation_type varchar(255),
additional_info varchar,
CONSTRAINT relation_unq_key UNIQUE (from_id, from_type, relation_type_group, relation_type, to_id, to_type)
CONSTRAINT relation_pkey PRIMARY KEY (from_id, from_type, relation_type_group, relation_type, to_id, to_type)
);
CREATE TABLE IF NOT EXISTS tb_user (

4
dao/src/main/resources/sql/schema-ts.sql

@ -23,7 +23,7 @@ CREATE TABLE IF NOT EXISTS ts_kv (
str_v varchar(10000000),
long_v bigint,
dbl_v double precision,
CONSTRAINT ts_kv_unq_key UNIQUE (entity_type, entity_id, key, ts)
CONSTRAINT ts_kv_pkey PRIMARY KEY (entity_type, entity_id, key, ts)
);
CREATE TABLE IF NOT EXISTS ts_kv_latest (
@ -35,5 +35,5 @@ CREATE TABLE IF NOT EXISTS ts_kv_latest (
str_v varchar(10000000),
long_v bigint,
dbl_v double precision,
CONSTRAINT ts_kv_latest_unq_key UNIQUE (entity_type, entity_id, key)
CONSTRAINT ts_kv_latest_pkey PRIMARY KEY (entity_type, entity_id, key)
);

9
msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java

@ -32,7 +32,8 @@ import org.apache.http.impl.client.HttpClients;
import org.apache.http.impl.conn.PoolingHttpClientConnectionManager;
import org.apache.http.ssl.SSLContextBuilder;
import org.apache.http.ssl.SSLContexts;
import org.junit.*;
import org.junit.BeforeClass;
import org.junit.Rule;
import org.junit.rules.TestRule;
import org.junit.rules.TestWatcher;
import org.junit.runner.Description;
@ -43,7 +44,10 @@ import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.msa.mapper.WsTelemetryResponse;
import javax.net.ssl.*;
import javax.net.ssl.SSLContext;
import javax.net.ssl.SSLSession;
import javax.net.ssl.SSLSocket;
import java.net.URI;
import java.security.cert.X509Certificate;
import java.util.List;
@ -54,6 +58,7 @@ import java.util.Random;
public abstract class AbstractContainerTest {
protected static final String HTTPS_URL = "https://localhost";
protected static final String WSS_URL = "wss://localhost";
protected static String TB_TOKEN;
protected static RestClient restClient;
protected ObjectMapper mapper = new ObjectMapper();

66
msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/HttpClientTest.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.msa.connectivity;
import com.fasterxml.jackson.databind.JsonNode;
import com.google.common.collect.Sets;
import org.junit.Assert;
import org.junit.Test;
@ -25,6 +26,17 @@ import org.thingsboard.server.msa.AbstractContainerTest;
import org.thingsboard.server.msa.WsClient;
import org.thingsboard.server.msa.mapper.WsTelemetryResponse;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertTrue;
import static org.thingsboard.server.common.data.DataConstants.DEVICE;
import static org.thingsboard.server.common.data.DataConstants.SHARED_SCOPE;
public class HttpClientTest extends AbstractContainerTest {
@Test
@ -52,6 +64,58 @@ public class HttpClientTest extends AbstractContainerTest {
Assert.assertTrue(verify(actualLatestTelemetry, "doubleKey", Double.toString(42.0)));
Assert.assertTrue(verify(actualLatestTelemetry, "longKey", Long.toString(73)));
restClient.getRestTemplate().delete(HTTPS_URL + "/api/device/" + device.getId());
restClient.deleteDevice(device.getId());
}
@Test
public void getAttributes() throws Exception {
restClient.login("tenant@thingsboard.org", "tenant");
TB_TOKEN = restClient.getToken();
Device device = createDevice("test");
String accessToken = restClient.getCredentials(device.getId()).getCredentialsId();
assertNotNull(accessToken);
ResponseEntity deviceSharedAttributes = restClient.getRestTemplate()
.postForEntity(HTTPS_URL + "/api/plugins/telemetry/" + DEVICE + "/" + device.getId().toString() + "/attributes/" + SHARED_SCOPE, mapper.readTree(createPayload().toString()),
ResponseEntity.class,
accessToken);
Assert.assertTrue(deviceSharedAttributes.getStatusCode().is2xxSuccessful());
ResponseEntity deviceClientsAttributes = restClient.getRestTemplate()
.postForEntity(HTTPS_URL + "/api/v1/" + accessToken + "/attributes/", mapper.readTree(createPayload().toString()),
ResponseEntity.class,
accessToken);
Assert.assertTrue(deviceClientsAttributes.getStatusCode().is2xxSuccessful());
TimeUnit.SECONDS.sleep(3);
Optional<JsonNode> allOptional = restClient.getAttributes(accessToken, null, null);
assertTrue(allOptional.isPresent());
JsonNode all = allOptional.get();
assertEquals(2, all.size());
assertEquals(mapper.readTree(createPayload().toString()), all.get("shared"));
assertEquals(mapper.readTree(createPayload().toString()), all.get("client"));
Optional<JsonNode> sharedOptional = restClient.getAttributes(accessToken, null, "stringKey");
assertTrue(sharedOptional.isPresent());
JsonNode shared = sharedOptional.get();
assertEquals(shared.get("shared").get("stringKey"), mapper.readTree(createPayload().get("stringKey").toString()));
assertFalse(shared.has("client"));
Optional<JsonNode> clientOptional = restClient.getAttributes(accessToken, "longKey,stringKey", null);
assertTrue(clientOptional.isPresent());
JsonNode client = clientOptional.get();
assertFalse(client.has("shared"));
assertEquals(mapper.readTree(createPayload().get("longKey").toString()), client.get("client").get("longKey"));
assertEquals(client.get("client").get("stringKey"), mapper.readTree(createPayload().get("stringKey").toString()));
restClient.deleteDevice(device.getId());
}
}

108
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNode.java

@ -21,8 +21,15 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ArrayNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.Data;
import lombok.NoArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.rule.engine.api.*;
import org.apache.commons.lang3.math.NumberUtils;
import org.thingsboard.rule.engine.api.RuleNode;
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.rule.engine.api.util.DonAsynchron;
import org.thingsboard.rule.engine.api.util.TbNodeUtils;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
@ -37,7 +44,9 @@ import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import static org.thingsboard.rule.engine.api.TbRelationTypes.SUCCESS;
import static org.thingsboard.rule.engine.metadata.TbGetTelemetryNodeConfiguration.*;
import static org.thingsboard.rule.engine.metadata.TbGetTelemetryNodeConfiguration.FETCH_MODE_ALL;
import static org.thingsboard.rule.engine.metadata.TbGetTelemetryNodeConfiguration.FETCH_MODE_FIRST;
import static org.thingsboard.rule.engine.metadata.TbGetTelemetryNodeConfiguration.MAX_FETCH_SIZE;
import static org.thingsboard.server.common.data.kv.Aggregation.NONE;
/**
@ -57,8 +66,6 @@ public class TbGetTelemetryNode implements TbNode {
private TbGetTelemetryNodeConfiguration config;
private List<String> tsKeyNames;
private long startTsOffset;
private long endTsOffset;
private int limit;
private ObjectMapper mapper;
private String fetchMode;
@ -67,8 +74,6 @@ public class TbGetTelemetryNode implements TbNode {
public void init(TbContext ctx, TbNodeConfiguration configuration) throws TbNodeException {
this.config = TbNodeUtils.convert(configuration, TbGetTelemetryNodeConfiguration.class);
tsKeyNames = config.getLatestTsKeyNames();
startTsOffset = TimeUnit.valueOf(config.getStartIntervalTimeUnit()).toMillis(config.getStartInterval());
endTsOffset = TimeUnit.valueOf(config.getEndIntervalTimeUnit()).toMillis(config.getEndInterval());
limit = config.getFetchMode().equals(FETCH_MODE_ALL) ? MAX_FETCH_SIZE : 1;
fetchMode = config.getFetchMode();
mapper = new ObjectMapper();
@ -82,8 +87,8 @@ public class TbGetTelemetryNode implements TbNode {
ctx.tellFailure(msg, new IllegalStateException("Telemetry is not selected!"));
} else {
try {
List<ReadTsKvQuery> queries = buildQueries();
ListenableFuture<List<TsKvEntry>> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), queries);
checkMetadataKeyPatterns(msg);
ListenableFuture<List<TsKvEntry>> list = ctx.getTimeseriesService().findAll(ctx.getTenantId(), msg.getOriginator(), buildQueries(msg));
DonAsynchron.withCallback(list, data -> {
process(data, msg);
TbMsg newMsg = ctx.transformMsg(msg, msg.getType(), msg.getOriginator(), msg.getMetaData(), msg.getData());
@ -95,10 +100,12 @@ public class TbGetTelemetryNode implements TbNode {
}
}
private List<ReadTsKvQuery> buildQueries() {
long ts = System.currentTimeMillis();
long startTs = ts - startTsOffset;
long endTs = ts - endTsOffset;
@Override
public void destroy() {
}
private List<ReadTsKvQuery> buildQueries(TbMsg msg) {
String orderBy;
if (fetchMode.equals(FETCH_MODE_FIRST) || fetchMode.equals(FETCH_MODE_ALL)) {
orderBy = "ASC";
@ -106,7 +113,7 @@ public class TbGetTelemetryNode implements TbNode {
orderBy = "DESC";
}
return tsKeyNames.stream()
.map(key -> new BaseReadTsKvQuery(key, startTs, endTs, 1, limit, NONE, orderBy))
.map(key -> new BaseReadTsKvQuery(key, getInterval(msg).getStartTs(), getInterval(msg).getEndTs(), 1, limit, NONE, orderBy))
.collect(Collectors.toList());
}
@ -162,8 +169,79 @@ public class TbGetTelemetryNode implements TbNode {
return obj;
}
@Override
public void destroy() {
private Interval getInterval(TbMsg msg) {
Interval interval = new Interval();
if (config.isUseMetadataIntervalPatterns()) {
if (isParsable(msg, config.getStartIntervalPattern())) {
interval.setStartTs(Long.parseLong(TbNodeUtils.processPattern(config.getStartIntervalPattern(), msg.getMetaData())));
}
if (isParsable(msg, config.getEndIntervalPattern())) {
interval.setEndTs(Long.parseLong(TbNodeUtils.processPattern(config.getEndIntervalPattern(), msg.getMetaData())));
}
} else {
long ts = System.currentTimeMillis();
interval.setStartTs(ts - TimeUnit.valueOf(config.getStartIntervalTimeUnit()).toMillis(config.getStartInterval()));
interval.setEndTs(ts - TimeUnit.valueOf(config.getEndIntervalTimeUnit()).toMillis(config.getEndInterval()));
}
return interval;
}
private boolean isParsable(TbMsg msg, String pattern) {
return NumberUtils.isParsable(TbNodeUtils.processPattern(pattern, msg.getMetaData()));
}
private void checkMetadataKeyPatterns(TbMsg msg) {
isUndefined(msg, config.getStartIntervalPattern(), config.getEndIntervalPattern());
isInvalid(msg, config.getStartIntervalPattern(), config.getEndIntervalPattern());
}
private void isUndefined(TbMsg msg, String startIntervalPattern, String endIntervalPattern) {
if (getMetadataValue(msg, startIntervalPattern) == null && getMetadataValue(msg, endIntervalPattern) == null) {
throw new IllegalArgumentException("Message metadata values: '" +
replaceRegex(startIntervalPattern) + "' and '" +
replaceRegex(endIntervalPattern) + "' are undefined");
} else {
if (getMetadataValue(msg, startIntervalPattern) == null) {
throw new IllegalArgumentException("Message metadata value: '" +
replaceRegex(startIntervalPattern) + "' is undefined");
}
if (getMetadataValue(msg, endIntervalPattern) == null) {
throw new IllegalArgumentException("Message metadata value: '" +
replaceRegex(endIntervalPattern) + "' is undefined");
}
}
}
private void isInvalid(TbMsg msg, String startIntervalPattern, String endIntervalPattern) {
if (getInterval(msg).getStartTs() == null && getInterval(msg).getEndTs() == null) {
throw new IllegalArgumentException("Message metadata values: '" +
replaceRegex(startIntervalPattern) + "' and '" +
replaceRegex(endIntervalPattern) + "' have invalid format");
} else {
if (getInterval(msg).getStartTs() == null) {
throw new IllegalArgumentException("Message metadata value: '" +
replaceRegex(startIntervalPattern) + "' has invalid format");
}
if (getInterval(msg).getEndTs() == null) {
throw new IllegalArgumentException("Message metadata value: '" +
replaceRegex(endIntervalPattern) + "' has invalid format");
}
}
}
private String getMetadataValue(TbMsg msg, String pattern) {
return msg.getMetaData().getValue(replaceRegex(pattern));
}
private String replaceRegex(String pattern) {
return pattern.replaceAll("[${}]", "");
}
@Data
@NoArgsConstructor
private static class Interval {
private Long startTs;
private Long endTs;
}
}

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/metadata/TbGetTelemetryNodeConfiguration.java

@ -35,6 +35,12 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration<TbGetT
private int startInterval;
private int endInterval;
private String startIntervalPattern;
private String endIntervalPattern;
private boolean useMetadataIntervalPatterns;
private String startIntervalTimeUnit;
private String endIntervalTimeUnit;
private String fetchMode; //FIRST, LAST, LATEST
@ -52,6 +58,9 @@ public class TbGetTelemetryNodeConfiguration implements NodeConfiguration<TbGetT
configuration.setStartInterval(2);
configuration.setEndIntervalTimeUnit(TimeUnit.MINUTES.name());
configuration.setEndInterval(1);
configuration.setUseMetadataIntervalPatterns(false);
configuration.setStartIntervalPattern("");
configuration.setEndIntervalPattern("");
return configuration;
}
}

8
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

1
tools/pom.xml

@ -23,7 +23,6 @@
<version>2.2.1-SNAPSHOT</version>
<artifactId>thingsboard</artifactId>
</parent>
<groupId>org.thingsboard</groupId>
<artifactId>tools</artifactId>
<packaging>jar</packaging>

53
tools/src/main/java/org/thingsboard/client/tools/RestClient.java

@ -36,6 +36,7 @@ import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType;
import java.io.IOException;
import java.util.Collections;
@ -107,6 +108,27 @@ public class RestClient implements ClientHttpRequestInterceptor {
}
}
public Optional<JsonNode> getAttributes(String accessToken, String clientKeys, String sharedKeys) {
Map<String, String> params = new HashMap<>();
params.put("accessToken", accessToken);
params.put("clientKeys", clientKeys);
params.put("sharedKeys", sharedKeys);
try {
ResponseEntity<JsonNode> telemetryEntity = restTemplate.getForEntity(baseURL + "/api/v1/{accessToken}/attributes?clientKeys={clientKeys}&sharedKeys={sharedKeys}", JsonNode.class, params);
return Optional.of(telemetryEntity.getBody());
} catch (HttpClientErrorException exception) {
if (exception.getStatusCode() == HttpStatus.NOT_FOUND) {
return Optional.empty();
} else {
throw exception;
}
}
}
public Customer createCustomer(Customer customer) {
return restTemplate.postForEntity(baseURL + "/api/customer", customer, Customer.class).getBody();
}
public Customer createCustomer(String title) {
Customer customer = new Customer();
customer.setTitle(title);
@ -120,6 +142,25 @@ public class RestClient implements ClientHttpRequestInterceptor {
return restTemplate.postForEntity(baseURL + "/api/device", device, Device.class).getBody();
}
public DeviceCredentials updateDeviceCredentials(DeviceId deviceId, String token) {
DeviceCredentials deviceCredentials = getCredentials(deviceId);
deviceCredentials.setCredentialsType(DeviceCredentialsType.ACCESS_TOKEN);
deviceCredentials.setCredentialsId(token);
return saveDeviceCredentials(deviceCredentials);
}
public DeviceCredentials saveDeviceCredentials(DeviceCredentials deviceCredentials) {
return restTemplate.postForEntity(baseURL + "/api/device/credentials", deviceCredentials, DeviceCredentials.class).getBody();
}
public Device createDevice(Device device) {
return restTemplate.postForEntity(baseURL + "/api/device", device, Device.class).getBody();
}
public Asset createAsset(Asset asset) {
return restTemplate.postForEntity(baseURL + "/api/asset", asset, Asset.class).getBody();
}
public Asset createAsset(String name, String type) {
Asset asset = new Asset();
asset.setName(name);
@ -131,6 +172,18 @@ public class RestClient implements ClientHttpRequestInterceptor {
return restTemplate.postForEntity(baseURL + "/api/alarm", alarm, Alarm.class).getBody();
}
public void deleteCustomer(CustomerId customerId) {
restTemplate.delete(baseURL + "/api/customer/{customerId}", customerId);
}
public void deleteDevice(DeviceId deviceId) {
restTemplate.delete(baseURL + "/api/device/{deviceId}", deviceId);
}
public void deleteAsset(AssetId assetId) {
restTemplate.delete(baseURL + "/api/asset/{assetId}", assetId);
}
public Device assignDevice(CustomerId customerId, DeviceId deviceId) {
return restTemplate.postForEntity(baseURL + "/api/customer/{customerId}/device/{deviceId}", null, Device.class,
customerId.toString(), deviceId.toString()).getBody();

45
ui/src/app/extension/extensions-forms/extension-form-modbus.directive.js

@ -31,17 +31,38 @@ export default function ExtensionFormModbusDirective($compile, $templateCache, $
var linker = function(scope, element) {
function TcpTransport() {
this.type = "tcp",
this.host = "localhost",
this.port = 502,
this.timeout = 5000,
this.reconnect = true,
this.rtuOverTcp = false
}
function UdpTransport() {
this.type = "udp",
this.host = "localhost",
this.port = 502,
this.timeout = 5000
}
function RtuTransport() {
this.type = "rtu",
this.portName = "COM1",
this.encoding = "ascii",
this.timeout = 5000,
this.baudRate = 115200,
this.dataBits = 7,
this.stopBits = 1,
this.parity ="even"
}
function Server() {
this.transport = {
"type": "tcp",
"host": "localhost",
"port": 502,
"timeout": 3000
};
this.transport = new TcpTransport();
this.devices = []
}
function Device() {
this.unitId = 1;
this.deviceName = "";
@ -105,9 +126,13 @@ export default function ExtensionFormModbusDirective($compile, $templateCache, $
scope.onTransportChanged = function(server) {
var type = server.transport.type;
server.transport = {};
server.transport.type = type;
server.transport.timeout = 3000;
if (type === "tcp") {
server.transport = new TcpTransport();
} else if (type === "udp") {
server.transport = new UdpTransport();
} else if (type === "rtu") {
server.transport = new RtuTransport();
}
scope.theForm.$setDirty();
};

104
ui/src/app/extension/extensions-forms/extension-form-modbus.tpl.html

@ -73,49 +73,69 @@
</div>
<div layout="row" ng-if="server.transport.type == 'tcp' || server.transport.type == 'udp'">
<md-input-container flex="33" class="md-block">
<label translate>extension.host</label>
<input required name="transportHost_{{serverIndex}}" ng-model="server.transport.host">
<div ng-messages="theForm['transportHost_' + serverIndex].$error">
<div translate ng-message="required">extension.field-required</div>
</div>
</md-input-container>
<div ng-if="server.transport.type == 'tcp' || server.transport.type == 'udp'">
<div layout="row">
<md-input-container flex="33" class="md-block">
<label translate>extension.host</label>
<input required name="transportHost_{{serverIndex}}" ng-model="server.transport.host">
<div ng-messages="theForm['transportHost_' + serverIndex].$error">
<div translate ng-message="required">extension.field-required</div>
</div>
</md-input-container>
<md-input-container flex="33" class="md-block">
<label translate>extension.port</label>
<input type="number"
required
name="transportPort_{{serverIndex}}"
ng-model="server.transport.port"
min="1"
max="65535"
>
<div ng-messages="theForm['transportPort_' + serverIndex].$error">
<div translate
ng-message="required"
>extension.field-required</div>
<div translate
ng-message="min"
>extension.port-range</div>
<div translate
ng-message="max"
>extension.port-range</div>
</div>
</md-input-container>
<md-input-container flex="33" class="md-block">
<label translate>extension.timeout</label>
<input type="number"
required name="transportTimeout_{{serverIndex}}"
ng-model="server.transport.timeout"
>
<div ng-messages="theForm['transportTimeout_' + serverIndex].$error">
<div translate
ng-message="required"
>extension.field-required</div>
</div>
</md-input-container>
<md-input-container flex="33" class="md-block">
<label translate>extension.port</label>
<input type="number"
required
name="transportPort_{{serverIndex}}"
ng-model="server.transport.port"
min="1"
max="65535"
>
<div ng-messages="theForm['transportPort_' + serverIndex].$error">
<div translate
ng-message="required"
>extension.field-required</div>
<div translate
ng-message="min"
>extension.port-range</div>
<div translate
ng-message="max"
>extension.port-range</div>
</div>
</md-input-container>
<md-input-container flex="33" class="md-block">
<label translate>extension.timeout</label>
<input type="number"
required name="transportTimeout_{{serverIndex}}"
ng-model="server.transport.timeout"
>
<div ng-messages="theForm['transportTimeout_' + serverIndex].$error">
<div translate
ng-message="required"
>extension.field-required</div>
</div>
</md-input-container>
</div>
<div layout="row" ng-if="server.transport.type == 'tcp'">
<md-input-container flex="50" class="md-block">
<md-checkbox aria-label="{{ 'extension.modbus-tcp-reconnect' | translate }}"
ng-checked="server.transport.reconnect"
name="transportTcpReconnect_{{serverIndex}}"
ng-model="server.transport.reconnect">{{ 'extension.modbus-tcp-reconnect' | translate }}
</md-checkbox>
</md-input-container>
<md-input-container flex="50" class="md-block">
<md-checkbox aria-label="{{ 'extension.modbus-rtu-over-tcp' | translate }}"
ng-checked="server.transport.rtuOverTcp"
name="transportRtuOverTcp_{{serverIndex}}"
ng-model="server.transport.rtuOverTcp">{{ 'extension.modbus-rtu-over-tcp' | translate }}
</md-checkbox>
</md-input-container>
</div>
</div>
<div ng-if="server.transport.type == 'rtu'">

2
ui/src/app/locale/locale.constant-en_US.json

@ -1008,6 +1008,8 @@
"modbus-add-server": "Add server/slave",
"modbus-add-server-prompt": "Please add server/slave",
"modbus-transport": "Transport",
"modbus-tcp-reconnect": "Automatically reconnect",
"modbus-rtu-over-tcp": "RTU over TCP",
"modbus-port-name": "Serial port name",
"modbus-encoding": "Encoding",
"modbus-parity": "Parity",

1528
ui/src/app/locale/locale.constant-ru_RU.json

File diff suppressed because it is too large
Loading…
Cancel
Save