Browse Source

Merge with master

pull/3688/head
Andrii Shvaika 6 years ago
parent
commit
3518a3d927
  1. 1
      application/src/main/data/upgrade/3.1.1/schema_update_before.sql
  2. 16
      application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
  3. 2
      application/src/main/java/org/thingsboard/server/controller/AlarmController.java
  4. 7
      application/src/main/java/org/thingsboard/server/controller/AuthController.java
  5. 10
      application/src/main/java/org/thingsboard/server/controller/UserController.java
  6. 1
      application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java
  7. 12
      application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java
  8. 10
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java
  9. 2
      application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationFailureHandler.java
  10. 1
      application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandler.java
  11. 24
      application/src/main/java/org/thingsboard/server/service/security/system/DefaultSystemSecurityService.java
  12. 5
      application/src/main/java/org/thingsboard/server/service/security/system/SystemSecurityService.java
  13. 2
      common/data/src/main/java/org/thingsboard/server/common/data/DeviceProfile.java
  14. 10
      common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java
  15. 29
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  16. 1
      dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
  17. 5
      dao/src/main/java/org/thingsboard/server/dao/model/sql/DeviceProfileEntity.java
  18. 33
      dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java
  19. 1
      dao/src/main/resources/sql/schema-entities-hsql.sql
  20. 1
      dao/src/main/resources/sql/schema-entities.sql
  21. 90
      dao/src/test/java/org/thingsboard/server/dao/service/BaseRuleChainServiceTest.java
  22. 4
      dao/src/test/resources/sql/system-test.sql
  23. 52
      netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientImpl.java
  24. 2
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java
  25. 13
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java
  26. 9
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractRelationActionNode.java
  27. 2
      rule-engine/rule-engine-components/src/main/resources/public/static/rulenode/rulenode-core-config.js
  28. 7
      rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbAlarmNodeTest.java
  29. 4
      ui-ngx/src/app/modules/home/components/profile/add-device-profile-dialog.component.html
  30. 4
      ui-ngx/src/app/modules/home/components/profile/add-device-profile-dialog.component.ts
  31. 4
      ui-ngx/src/app/modules/home/components/profile/device-profile.component.html
  32. 5
      ui-ngx/src/app/modules/home/components/profile/device-profile.component.ts
  33. 4
      ui-ngx/src/app/modules/home/components/rule-chain/rule-chain-autocomplete.component.ts
  34. 9
      ui-ngx/src/app/modules/home/components/widget/lib/alarms-table-widget.component.ts
  35. 2
      ui-ngx/src/app/modules/home/components/widget/lib/table-widget.models.ts
  36. 22
      ui-ngx/src/app/modules/home/components/widget/lib/timeseries-table-widget.component.ts
  37. 6
      ui-ngx/src/app/modules/home/components/wizard/device-wizard-dialog.component.html
  38. 20
      ui-ngx/src/app/modules/home/components/wizard/device-wizard-dialog.component.ts
  39. 4
      ui-ngx/src/app/modules/home/pages/admin/general-settings.component.html
  40. 3
      ui-ngx/src/app/modules/home/pages/admin/general-settings.component.ts
  41. 18
      ui-ngx/src/app/modules/home/pages/admin/oauth2-settings.component.ts
  42. 1
      ui-ngx/src/app/shared/components/queue/queue-type-list.component.html
  43. 1
      ui-ngx/src/app/shared/models/device.models.ts
  44. 6
      ui-ngx/src/assets/locale/locale.constant-en_US.json
  45. 2048
      ui-ngx/src/assets/locale/locale.constant-pl_BR.json

1
application/src/main/data/upgrade/3.1.1/schema_update_before.sql

@ -96,6 +96,7 @@ CREATE TABLE IF NOT EXISTS device_profile (
is_default boolean,
tenant_id uuid,
default_rule_chain_id uuid,
default_queue_name varchar(255),
provision_device_key varchar,
CONSTRAINT device_profile_name_unq_key UNIQUE (tenant_id, name),
CONSTRAINT device_provision_key_unq_key UNIQUE (provision_device_key),

16
application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java

@ -259,26 +259,26 @@ class DefaultTbContext implements TbContext {
}
public TbMsg customerCreatedMsg(Customer customer, RuleNodeId ruleNodeId) {
return entityCreatedMsg(customer, customer.getId(), ruleNodeId);
return entityActionMsg(customer, customer.getId(), ruleNodeId, DataConstants.ENTITY_CREATED);
}
public TbMsg deviceCreatedMsg(Device device, RuleNodeId ruleNodeId) {
return entityCreatedMsg(device, device.getId(), ruleNodeId);
return entityActionMsg(device, device.getId(), ruleNodeId, DataConstants.ENTITY_CREATED);
}
public TbMsg assetCreatedMsg(Asset asset, RuleNodeId ruleNodeId) {
return entityCreatedMsg(asset, asset.getId(), ruleNodeId);
return entityActionMsg(asset, asset.getId(), ruleNodeId, DataConstants.ENTITY_CREATED);
}
public TbMsg alarmCreatedMsg(Alarm alarm, RuleNodeId ruleNodeId) {
return entityCreatedMsg(alarm, alarm.getId(), ruleNodeId);
public TbMsg alarmActionMsg(Alarm alarm, RuleNodeId ruleNodeId, String action) {
return entityActionMsg(alarm, alarm.getId(), ruleNodeId, action);
}
public <E, I extends EntityId> TbMsg entityCreatedMsg(E entity, I id, RuleNodeId ruleNodeId) {
public <E, I extends EntityId> TbMsg entityActionMsg(E entity, I id, RuleNodeId ruleNodeId, String action) {
try {
return TbMsg.newMsg(DataConstants.ENTITY_CREATED, id, getActionMetaData(ruleNodeId), mapper.writeValueAsString(mapper.valueToTree(entity)));
return TbMsg.newMsg(action, id, getActionMetaData(ruleNodeId), mapper.writeValueAsString(mapper.valueToTree(entity)));
} catch (JsonProcessingException | IllegalArgumentException e) {
throw new RuntimeException("Failed to process " + id.getEntityType().name().toLowerCase() + " created msg: " + e);
throw new RuntimeException("Failed to process " + id.getEntityType().name().toLowerCase() + " " + action + " msg: " + e);
}
}

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

@ -126,6 +126,7 @@ public class AlarmController extends BaseController {
long ackTs = System.currentTimeMillis();
alarmService.ackAlarm(getCurrentUser().getTenantId(), alarmId, ackTs).get();
alarm.setAckTs(ackTs);
alarm.setStatus(alarm.getStatus().isCleared() ? AlarmStatus.CLEARED_ACK : AlarmStatus.ACTIVE_ACK);
logEntityAction(alarm.getOriginator(), alarm, getCurrentUser().getCustomerId(), ActionType.ALARM_ACK, null);
} catch (Exception e) {
throw handleException(e);
@ -143,6 +144,7 @@ public class AlarmController extends BaseController {
long clearTs = System.currentTimeMillis();
alarmService.clearAlarm(getCurrentUser().getTenantId(), alarmId, null, clearTs).get();
alarm.setClearTs(clearTs);
alarm.setStatus(alarm.getStatus().isAck() ? AlarmStatus.CLEARED_ACK : AlarmStatus.CLEARED_UNACK);
logEntityAction(alarm.getOriginator(), alarm, getCurrentUser().getCustomerId(), ActionType.ALARM_CLEAR, null);
} catch (Exception e) {
throw handleException(e);

7
application/src/main/java/org/thingsboard/server/controller/AuthController.java

@ -166,7 +166,8 @@ public class AuthController extends BaseController {
try {
String email = resetPasswordByEmailRequest.get("email").asText();
UserCredentials userCredentials = userService.requestPasswordReset(TenantId.SYS_TENANT_ID, email);
String baseUrl = MiscUtils.constructBaseUrl(request);
User user = userService.findUserById(TenantId.SYS_TENANT_ID, userCredentials.getUserId());
String baseUrl = systemSecurityService.getBaseUrl(user.getTenantId(), user.getCustomerId(), request);
String resetUrl = String.format("%s/api/noauth/resetPassword?resetToken=%s", baseUrl,
userCredentials.getResetToken());
@ -214,7 +215,7 @@ public class AuthController extends BaseController {
User user = userService.findUserById(TenantId.SYS_TENANT_ID, credentials.getUserId());
UserPrincipal principal = new UserPrincipal(UserPrincipal.Type.USER_NAME, user.getEmail());
SecurityUser securityUser = new SecurityUser(user, credentials.isEnabled(), principal);
String baseUrl = MiscUtils.constructBaseUrl(request);
String baseUrl = systemSecurityService.getBaseUrl(user.getTenantId(), user.getCustomerId(), request);
String loginUrl = String.format("%s/login", baseUrl);
String email = user.getEmail();
@ -261,7 +262,7 @@ public class AuthController extends BaseController {
User user = userService.findUserById(TenantId.SYS_TENANT_ID, userCredentials.getUserId());
UserPrincipal principal = new UserPrincipal(UserPrincipal.Type.USER_NAME, user.getEmail());
SecurityUser securityUser = new SecurityUser(user, userCredentials.isEnabled(), principal);
String baseUrl = MiscUtils.constructBaseUrl(request);
String baseUrl = systemSecurityService.getBaseUrl(user.getTenantId(), user.getCustomerId(), request);
String loginUrl = String.format("%s/login", baseUrl);
String email = user.getEmail();
mailService.sendPasswordWasResetEmail(loginUrl, email);

10
application/src/main/java/org/thingsboard/server/controller/UserController.java

@ -52,6 +52,7 @@ import org.thingsboard.server.service.security.model.token.JwtToken;
import org.thingsboard.server.service.security.model.token.JwtTokenFactory;
import org.thingsboard.server.service.security.permission.Operation;
import org.thingsboard.server.service.security.permission.Resource;
import org.thingsboard.server.service.security.system.SystemSecurityService;
import org.thingsboard.server.utils.MiscUtils;
import javax.servlet.http.HttpServletRequest;
@ -78,6 +79,9 @@ public class UserController extends BaseController {
@Autowired
private RefreshTokenRepository refreshTokenRepository;
@Autowired
private SystemSecurityService systemSecurityService;
@PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')")
@RequestMapping(value = "/user/{userId}", method = RequestMethod.GET)
@ -146,7 +150,7 @@ public class UserController extends BaseController {
if (sendEmail) {
SecurityUser authUser = getCurrentUser();
UserCredentials userCredentials = userService.findUserCredentialsByUserId(authUser.getTenantId(), savedUser.getId());
String baseUrl = MiscUtils.constructBaseUrl(request);
String baseUrl = systemSecurityService.getBaseUrl(getTenantId(), getCurrentUser().getCustomerId(), request);
String activateUrl = String.format(ACTIVATE_URL_PATTERN, baseUrl,
userCredentials.getActivateToken());
String email = savedUser.getEmail();
@ -186,7 +190,7 @@ public class UserController extends BaseController {
UserCredentials userCredentials = userService.findUserCredentialsByUserId(getCurrentUser().getTenantId(), user.getId());
if (!userCredentials.isEnabled()) {
String baseUrl = MiscUtils.constructBaseUrl(request);
String baseUrl = systemSecurityService.getBaseUrl(getTenantId(), getCurrentUser().getCustomerId(), request);
String activateUrl = String.format(ACTIVATE_URL_PATTERN, baseUrl,
userCredentials.getActivateToken());
mailService.sendActivationEmail(activateUrl, email);
@ -211,7 +215,7 @@ public class UserController extends BaseController {
SecurityUser authUser = getCurrentUser();
UserCredentials userCredentials = userService.findUserCredentialsByUserId(authUser.getTenantId(), user.getId());
if (!userCredentials.isEnabled()) {
String baseUrl = MiscUtils.constructBaseUrl(request);
String baseUrl = systemSecurityService.getBaseUrl(getTenantId(), getCurrentUser().getCustomerId(), request);
String activateUrl = String.format(ACTIVATE_URL_PATTERN, baseUrl,
userCredentials.getActivateToken());
return activateUrl;

1
application/src/main/java/org/thingsboard/server/service/install/DefaultSystemDataLoaderService.java

@ -182,6 +182,7 @@ public class DefaultSystemDataLoaderService implements SystemDataLoaderService {
generalSettings.setKey("general");
ObjectNode node = objectMapper.createObjectNode();
node.put("baseUrl", "http://localhost:8080");
node.put("prohibitDifferentUrl", true);
generalSettings.setJsonValue(node);
adminSettingsService.saveAdminSettings(TenantId.SYS_TENANT_ID, generalSettings);

12
application/src/main/java/org/thingsboard/server/service/profile/DefaultTbDeviceProfileCache.java

@ -104,7 +104,7 @@ public class DefaultTbDeviceProfileCache implements TbDeviceProfileCache {
if (old != null) {
DeviceProfile newProfile = get(tenantId, deviceId);
if (newProfile == null || !old.equals(newProfile.getId())) {
notifyDeviceListeners(deviceId, newProfile);
notifyDeviceListeners(tenantId, deviceId, newProfile);
}
}
}
@ -150,10 +150,12 @@ public class DefaultTbDeviceProfileCache implements TbDeviceProfileCache {
}
}
private void notifyDeviceListeners(DeviceId deviceId, DeviceProfile profile) {
ConcurrentMap<EntityId, BiConsumer<DeviceId, DeviceProfile>> tenantListeners = deviceProfileListeners.get(profile.getTenantId());
if (tenantListeners != null) {
tenantListeners.forEach((id, listener) -> listener.accept(deviceId, profile));
private void notifyDeviceListeners(TenantId tenantId, DeviceId deviceId, DeviceProfile profile) {
if (profile != null) {
ConcurrentMap<EntityId, BiConsumer<DeviceId, DeviceProfile>> tenantListeners = deviceProfileListeners.get(tenantId);
if (tenantListeners != null) {
tenantListeners.forEach((id, listener) -> listener.accept(deviceId, profile));
}
}
}

10
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbClusterService.java

@ -158,8 +158,16 @@ public class DefaultTbClusterService implements TbClusterService {
private TbMsg transformMsg(TbMsg tbMsg, DeviceProfile deviceProfile) {
if (deviceProfile != null) {
RuleChainId targetRuleChainId = deviceProfile.getDefaultRuleChainId();
if (targetRuleChainId != null && !targetRuleChainId.equals(tbMsg.getRuleChainId())) {
String targetQueueName = deviceProfile.getDefaultQueueName();
boolean isRuleChainTransform = targetRuleChainId != null && !targetRuleChainId.equals(tbMsg.getRuleChainId());
boolean isQueueTransform = targetQueueName != null && !targetQueueName.equals(tbMsg.getQueueName());
if (isRuleChainTransform && isQueueTransform) {
tbMsg = TbMsg.transformMsg(tbMsg, targetRuleChainId, targetQueueName);
} else if (isRuleChainTransform) {
tbMsg = TbMsg.transformMsg(tbMsg, targetRuleChainId);
} else if (isQueueTransform) {
tbMsg = TbMsg.transformMsg(tbMsg, targetQueueName);
}
}
return tbMsg;

2
application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationFailureHandler.java

@ -30,7 +30,7 @@ import java.nio.charset.StandardCharsets;
@Component(value = "oauth2AuthenticationFailureHandler")
@ConditionalOnProperty(prefix = "security.oauth2", value = "enabled", havingValue = "true")
public class Oauth2AuthenticationFailureHandler extends SimpleUrlAuthenticationFailureHandler {
public class Oauth2AuthenticationFailureHandler extends SimpleUrlAuthenticationFailureHandler {
@Override
public void onAuthenticationFailure(HttpServletRequest request,

1
application/src/main/java/org/thingsboard/server/service/security/auth/oauth2/Oauth2AuthenticationSuccessHandler.java

@ -63,7 +63,6 @@ public class Oauth2AuthenticationSuccessHandler extends SimpleUrlAuthenticationS
public void onAuthenticationSuccess(HttpServletRequest request,
HttpServletResponse response,
Authentication authentication) throws IOException {
String baseUrl = MiscUtils.constructBaseUrl(request);
try {
OAuth2AuthenticationToken token = (OAuth2AuthenticationToken) authentication;

24
application/src/main/java/org/thingsboard/server/service/security/system/DefaultSystemSecurityService.java

@ -40,17 +40,20 @@ import org.thingsboard.rule.engine.api.MailService;
import org.thingsboard.server.common.data.AdminSettings;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.security.UserCredentials;
import org.thingsboard.server.common.data.security.model.SecuritySettings;
import org.thingsboard.server.common.data.security.model.UserPasswordPolicy;
import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.dao.settings.AdminSettingsService;
import org.thingsboard.server.dao.user.UserService;
import org.thingsboard.server.dao.user.UserServiceImpl;
import org.thingsboard.server.service.security.exception.UserPasswordExpiredException;
import org.thingsboard.server.common.data.security.model.SecuritySettings;
import org.thingsboard.server.common.data.security.model.UserPasswordPolicy;
import org.thingsboard.server.utils.MiscUtils;
import javax.annotation.Resource;
import javax.servlet.http.HttpServletRequest;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
@ -146,7 +149,7 @@ public class DefaultSystemSecurityService implements SystemSecurityService {
if (isPositiveInteger(securitySettings.getPasswordPolicy().getPasswordExpirationPeriodDays())) {
if ((userCredentials.getCreatedTime()
+ TimeUnit.DAYS.toMillis(securitySettings.getPasswordPolicy().getPasswordExpirationPeriodDays()))
< System.currentTimeMillis()) {
< System.currentTimeMillis()) {
userCredentials = userService.requestExpiredPasswordReset(tenantId, userCredentials.getId());
throw new UserPasswordExpiredException("User password expired!", userCredentials.getResetToken());
}
@ -197,6 +200,21 @@ public class DefaultSystemSecurityService implements SystemSecurityService {
}
}
@Override
public String getBaseUrl(TenantId tenantId, CustomerId customerId, HttpServletRequest httpServletRequest) {
String baseUrl;
AdminSettings generalSettings = adminSettingsService.findAdminSettingsByKey(TenantId.SYS_TENANT_ID, "general");
JsonNode prohibitDifferentUrl = generalSettings.getJsonValue().get("prohibitDifferentUrl");
if (prohibitDifferentUrl != null && prohibitDifferentUrl.asBoolean()) {
baseUrl = generalSettings.getJsonValue().get("baseUrl").asText();
} else {
baseUrl = MiscUtils.constructBaseUrl(httpServletRequest);
}
return baseUrl;
}
private static boolean isPositiveInteger(Integer val) {
return val != null && val.intValue() > 0;
}

5
application/src/main/java/org/thingsboard/server/service/security/system/SystemSecurityService.java

@ -16,11 +16,14 @@
package org.thingsboard.server.service.security.system;
import org.springframework.security.core.AuthenticationException;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.security.UserCredentials;
import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.common.data.security.model.SecuritySettings;
import javax.servlet.http.HttpServletRequest;
public interface SystemSecurityService {
SecuritySettings getSecuritySettings(TenantId tenantId);
@ -31,4 +34,6 @@ public interface SystemSecurityService {
void validatePassword(TenantId tenantId, String password, UserCredentials userCredentials) throws DataValidationException;
String getBaseUrl(TenantId tenantId, CustomerId customerId, HttpServletRequest httpServletRequest);
}

2
common/data/src/main/java/org/thingsboard/server/common/data/DeviceProfile.java

@ -43,6 +43,7 @@ public class DeviceProfile extends SearchTextBased<DeviceProfileId> implements H
private DeviceTransportType transportType;
private DeviceProfileProvisionType provisionType;
private RuleChainId defaultRuleChainId;
private String defaultQueueName;
private transient DeviceProfileData profileData;
@JsonIgnore
private byte[] profileDataBytes;
@ -63,6 +64,7 @@ public class DeviceProfile extends SearchTextBased<DeviceProfileId> implements H
this.description = deviceProfile.getDescription();
this.isDefault = deviceProfile.isDefault();
this.defaultRuleChainId = deviceProfile.getDefaultRuleChainId();
this.defaultQueueName = deviceProfile.getDefaultQueueName();
this.setProfileData(deviceProfile.getProfileData());
this.provisionDeviceKey = deviceProfile.getProvisionDeviceKey();
}

10
common/message/src/main/java/org/thingsboard/server/common/msg/TbMsg.java

@ -102,6 +102,16 @@ public final class TbMsg implements Serializable {
tbMsg.data, ruleChainId, null, tbMsg.ruleNodeExecCounter.get(), tbMsg.getCallback());
}
public static TbMsg transformMsg(TbMsg tbMsg, String queueName) {
return new TbMsg(queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, tbMsg.metaData, tbMsg.dataType,
tbMsg.data, tbMsg.getRuleChainId(), null, tbMsg.ruleNodeExecCounter.get(), tbMsg.getCallback());
}
public static TbMsg transformMsg(TbMsg tbMsg, RuleChainId ruleChainId, String queueName) {
return new TbMsg(queueName, tbMsg.id, tbMsg.ts, tbMsg.type, tbMsg.originator, tbMsg.metaData, tbMsg.dataType,
tbMsg.data, ruleChainId, null, tbMsg.ruleNodeExecCounter.get(), tbMsg.getCallback());
}
public static TbMsg newMsg(TbMsg tbMsg, RuleChainId ruleChainId, RuleNodeId ruleNodeId) {
return new TbMsg(tbMsg.getQueueName(), UUID.randomUUID(), tbMsg.getTs(), tbMsg.getType(), tbMsg.getOriginator(), tbMsg.getMetaData().copy(),
tbMsg.getDataType(), tbMsg.getData(), ruleChainId, ruleNodeId, tbMsg.ruleNodeExecCounter.get(), TbMsgCallback.EMPTY);

29
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

@ -373,10 +373,8 @@ public class DefaultTransportService implements TransportService {
metaData.putValue("deviceType", sessionInfo.getDeviceType());
metaData.putValue("ts", tsKv.getTs() + "");
JsonObject json = JsonUtils.getJsonObject(tsKv.getKvList());
RuleChainId ruleChainId = resolveRuleChainId(sessionInfo);
TbMsg tbMsg = TbMsg.newMsg(ServiceQueue.MAIN, SessionMsgType.POST_TELEMETRY_REQUEST.name(),
deviceId, metaData, gson.toJson(json), ruleChainId, null);
sendToRuleEngine(tenantId, tbMsg, packCallback);
sendToRuleEngine(tenantId, deviceId, sessionInfo, json, metaData, SessionMsgType.POST_TELEMETRY_REQUEST, packCallback);
}
}
}
@ -392,10 +390,7 @@ public class DefaultTransportService implements TransportService {
metaData.putValue("deviceName", sessionInfo.getDeviceName());
metaData.putValue("deviceType", sessionInfo.getDeviceType());
metaData.putValue("notifyDevice", "false");
RuleChainId ruleChainId = resolveRuleChainId(sessionInfo);
TbMsg tbMsg = TbMsg.newMsg(ServiceQueue.MAIN, SessionMsgType.POST_ATTRIBUTES_REQUEST.name(),
deviceId, metaData, gson.toJson(json), ruleChainId, null);
sendToRuleEngine(tenantId, tbMsg, new TransportTbQueueCallback(callback));
sendToRuleEngine(tenantId, deviceId, sessionInfo, json, metaData, SessionMsgType.POST_ATTRIBUTES_REQUEST, new TransportTbQueueCallback(callback));
}
}
@ -476,10 +471,8 @@ public class DefaultTransportService implements TransportService {
metaData.putValue("requestId", Integer.toString(msg.getRequestId()));
metaData.putValue("serviceId", serviceInfoProvider.getServiceId());
metaData.putValue("sessionId", sessionId.toString());
RuleChainId ruleChainId = resolveRuleChainId(sessionInfo);
TbMsg tbMsg = TbMsg.newMsg(ServiceQueue.MAIN, SessionMsgType.TO_SERVER_RPC_REQUEST.name(),
deviceId, metaData, gson.toJson(json), ruleChainId, null);
sendToRuleEngine(tenantId, tbMsg, new TransportTbQueueCallback(callback));
sendToRuleEngine(tenantId, deviceId, sessionInfo, json, metaData,
SessionMsgType.TO_SERVER_RPC_REQUEST, new TransportTbQueueCallback(callback));
String requestId = sessionId + "-" + msg.getRequestId();
toServerRpcPendingMap.put(requestId, new RpcRequestMetadata(sessionId, msg.getRequestId()));
scheduler.schedule(() -> processTimeout(requestId), clientSideRpcTimeout, TimeUnit.MILLISECONDS);
@ -731,17 +724,25 @@ public class DefaultTransportService implements TransportService {
ruleEngineMsgProducer.send(tpi, new TbProtoQueueMsg<>(tbMsg.getId(), msg), wrappedCallback);
}
private RuleChainId resolveRuleChainId(TransportProtos.SessionInfoProto sessionInfo) {
protected void sendToRuleEngine(TenantId tenantId, DeviceId deviceId, TransportProtos.SessionInfoProto sessionInfo, JsonObject json,
TbMsgMetaData metaData, SessionMsgType sessionMsgType, TbQueueCallback callback) {
DeviceProfileId deviceProfileId = new DeviceProfileId(new UUID(sessionInfo.getDeviceProfileIdMSB(), sessionInfo.getDeviceProfileIdLSB()));
DeviceProfile deviceProfile = deviceProfileCache.get(deviceProfileId);
RuleChainId ruleChainId;
String queueName;
if (deviceProfile == null) {
log.warn("[{}] Device profile is null!", deviceProfileId);
ruleChainId = null;
queueName = ServiceQueue.MAIN;
} else {
ruleChainId = deviceProfile.getDefaultRuleChainId();
String defaultQueueName = deviceProfile.getDefaultQueueName();
queueName = defaultQueueName != null ? defaultQueueName : ServiceQueue.MAIN;
}
return ruleChainId;
TbMsg tbMsg = TbMsg.newMsg(queueName, sessionMsgType.name(), deviceId, metaData, gson.toJson(json), ruleChainId, null);
sendToRuleEngine(tenantId, tbMsg, callback);
}
private class TransportTbQueueCallback implements TbQueueCallback {

1
dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java

@ -174,6 +174,7 @@ public class ModelConstants {
public static final String DEVICE_PROFILE_DESCRIPTION_PROPERTY = "description";
public static final String DEVICE_PROFILE_IS_DEFAULT_PROPERTY = "is_default";
public static final String DEVICE_PROFILE_DEFAULT_RULE_CHAIN_ID_PROPERTY = "default_rule_chain_id";
public static final String DEVICE_PROFILE_DEFAULT_QUEUE_NAME_PROPERTY = "default_queue_name";
public static final String DEVICE_PROFILE_PROVISION_DEVICE_KEY = "provision_device_key";
/**

5
dao/src/main/java/org/thingsboard/server/dao/model/sql/DeviceProfileEntity.java

@ -79,6 +79,9 @@ public final class DeviceProfileEntity extends BaseSqlEntity<DeviceProfile> impl
@Column(name = ModelConstants.DEVICE_PROFILE_DEFAULT_RULE_CHAIN_ID_PROPERTY, columnDefinition = "uuid")
private UUID defaultRuleChainId;
@Column(name = ModelConstants.DEVICE_PROFILE_DEFAULT_QUEUE_NAME_PROPERTY)
private String defaultQueueName;
@Type(type = "jsonb")
@Column(name = ModelConstants.DEVICE_PROFILE_PROFILE_DATA_PROPERTY, columnDefinition = "jsonb")
private JsonNode profileData;
@ -108,6 +111,7 @@ public final class DeviceProfileEntity extends BaseSqlEntity<DeviceProfile> impl
if (deviceProfile.getDefaultRuleChainId() != null) {
this.defaultRuleChainId = deviceProfile.getDefaultRuleChainId().getId();
}
this.defaultQueueName = deviceProfile.getDefaultQueueName();
this.provisionDeviceKey = deviceProfile.getProvisionDeviceKey();
}
@ -142,6 +146,7 @@ public final class DeviceProfileEntity extends BaseSqlEntity<DeviceProfile> impl
if (defaultRuleChainId != null) {
deviceProfile.setDefaultRuleChainId(new RuleChainId(defaultRuleChainId));
}
deviceProfile.setDefaultQueueName(defaultQueueName);
deviceProfile.setProvisionDeviceKey(provisionDeviceKey);
return deviceProfile;
}

33
dao/src/main/java/org/thingsboard/server/dao/rule/BaseRuleChainService.java

@ -20,11 +20,11 @@ import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
import org.apache.commons.lang3.StringUtils;
import org.hibernate.exception.ConstraintViolationException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
import org.springframework.util.CollectionUtils;
import org.thingsboard.server.common.data.BaseData;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.Tenant;
@ -53,10 +53,12 @@ import org.thingsboard.server.dao.tenant.TenantDao;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ExecutionException;
import java.util.stream.Collectors;
@ -136,6 +138,10 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
return null;
}
if (CollectionUtils.isNotEmpty(ruleChainMetaData.getConnections())) {
validateCircles(ruleChainMetaData.getConnections());
}
List<RuleNode> nodes = ruleChainMetaData.getNodes();
List<RuleNode> toAddOrUpdate = new ArrayList<>();
List<RuleNode> toDelete = new ArrayList<>();
@ -218,6 +224,31 @@ public class BaseRuleChainService extends AbstractEntityService implements RuleC
return loadRuleChainMetaData(tenantId, ruleChainMetaData.getRuleChainId());
}
private void validateCircles(List<NodeConnectionInfo> connectionInfos) {
Map<Integer, Set<Integer>> connectionsMap = new HashMap<>();
for (NodeConnectionInfo nodeConnection : connectionInfos) {
if (nodeConnection.getFromIndex() == nodeConnection.getToIndex()) {
throw new DataValidationException("Can't create the relation to yourself.");
}
connectionsMap
.computeIfAbsent(nodeConnection.getFromIndex(), from -> new HashSet<>())
.add(nodeConnection.getToIndex());
}
connectionsMap.keySet().forEach(key -> validateCircles(key, connectionsMap.get(key), connectionsMap));
}
private void validateCircles(int from, Set<Integer> toList, Map<Integer, Set<Integer>> connectionsMap) {
if (toList == null) {
return;
}
for (Integer to : toList) {
if (from == to) {
throw new DataValidationException("Can't create circling relations in rule chain.");
}
validateCircles(from, connectionsMap.get(to), connectionsMap);
}
}
@Override
public RuleChainMetaData loadRuleChainMetaData(TenantId tenantId, RuleChainId ruleChainId) {
Validator.validateId(ruleChainId, "Incorrect rule chain id.");

1
dao/src/main/resources/sql/schema-entities-hsql.sql

@ -170,6 +170,7 @@ CREATE TABLE IF NOT EXISTS device_profile (
is_default boolean,
tenant_id uuid,
default_rule_chain_id uuid,
default_queue_name varchar(255),
provision_device_key varchar,
CONSTRAINT device_profile_name_unq_key UNIQUE (tenant_id, name),
CONSTRAINT device_provision_key_unq_key UNIQUE (provision_device_key),

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

@ -188,6 +188,7 @@ CREATE TABLE IF NOT EXISTS device_profile (
is_default boolean,
tenant_id uuid,
default_rule_chain_id uuid,
default_queue_name varchar(255),
provision_device_key varchar,
CONSTRAINT device_profile_name_unq_key UNIQUE (tenant_id, name),
CONSTRAINT device_provision_key_unq_key UNIQUE (provision_device_key),

90
dao/src/test/java/org/thingsboard/server/dao/service/BaseRuleChainServiceTest.java

@ -317,6 +317,16 @@ public abstract class BaseRuleChainServiceTest extends AbstractServiceTest {
ruleChainService.deleteRuleChainById(tenantId, savedRuleChainMetaData.getRuleChainId());
}
@Test(expected = DataValidationException.class)
public void testUpdateRuleChainMetaDataWithCirclingRelation() throws Exception {
ruleChainService.saveRuleChainMetaData(tenantId, createRuleChainMetadataWithCirclingRelation());
}
@Test(expected = DataValidationException.class)
public void testUpdateRuleChainMetaDataWithCirclingRelation2() throws Exception {
ruleChainService.saveRuleChainMetaData(tenantId, createRuleChainMetadataWithCirclingRelation2());
}
private RuleChainMetaData createRuleChainMetadata() throws Exception {
RuleChain ruleChain = new RuleChain();
ruleChain.setName("My RuleChain");
@ -357,5 +367,85 @@ public abstract class BaseRuleChainServiceTest extends AbstractServiceTest {
return ruleChainService.saveRuleChainMetaData(tenantId, ruleChainMetaData);
}
private RuleChainMetaData createRuleChainMetadataWithCirclingRelation() throws Exception {
RuleChain ruleChain = new RuleChain();
ruleChain.setName("My RuleChain");
ruleChain.setTenantId(tenantId);
RuleChain savedRuleChain = ruleChainService.saveRuleChain(ruleChain);
RuleChainMetaData ruleChainMetaData = new RuleChainMetaData();
ruleChainMetaData.setRuleChainId(savedRuleChain.getId());
ObjectMapper mapper = new ObjectMapper();
RuleNode ruleNode1 = new RuleNode();
ruleNode1.setName("name1");
ruleNode1.setType("type1");
ruleNode1.setConfiguration(mapper.readTree("\"key1\": \"val1\""));
RuleNode ruleNode2 = new RuleNode();
ruleNode2.setName("name2");
ruleNode2.setType("type2");
ruleNode2.setConfiguration(mapper.readTree("\"key2\": \"val2\""));
RuleNode ruleNode3 = new RuleNode();
ruleNode3.setName("name3");
ruleNode3.setType("type3");
ruleNode3.setConfiguration(mapper.readTree("\"key3\": \"val3\""));
List<RuleNode> ruleNodes = new ArrayList<>();
ruleNodes.add(ruleNode1);
ruleNodes.add(ruleNode2);
ruleNodes.add(ruleNode3);
ruleChainMetaData.setFirstNodeIndex(0);
ruleChainMetaData.setNodes(ruleNodes);
ruleChainMetaData.addConnectionInfo(0,1,"success");
ruleChainMetaData.addConnectionInfo(0,2,"fail");
ruleChainMetaData.addConnectionInfo(1,2,"success");
ruleChainMetaData.addConnectionInfo(2,2,"success");
return ruleChainMetaData;
}
private RuleChainMetaData createRuleChainMetadataWithCirclingRelation2() throws Exception {
RuleChain ruleChain = new RuleChain();
ruleChain.setName("My RuleChain");
ruleChain.setTenantId(tenantId);
RuleChain savedRuleChain = ruleChainService.saveRuleChain(ruleChain);
RuleChainMetaData ruleChainMetaData = new RuleChainMetaData();
ruleChainMetaData.setRuleChainId(savedRuleChain.getId());
ObjectMapper mapper = new ObjectMapper();
RuleNode ruleNode1 = new RuleNode();
ruleNode1.setName("name1");
ruleNode1.setType("type1");
ruleNode1.setConfiguration(mapper.readTree("\"key1\": \"val1\""));
RuleNode ruleNode2 = new RuleNode();
ruleNode2.setName("name2");
ruleNode2.setType("type2");
ruleNode2.setConfiguration(mapper.readTree("\"key2\": \"val2\""));
RuleNode ruleNode3 = new RuleNode();
ruleNode3.setName("name3");
ruleNode3.setType("type3");
ruleNode3.setConfiguration(mapper.readTree("\"key3\": \"val3\""));
List<RuleNode> ruleNodes = new ArrayList<>();
ruleNodes.add(ruleNode1);
ruleNodes.add(ruleNode2);
ruleNodes.add(ruleNode3);
ruleChainMetaData.setFirstNodeIndex(0);
ruleChainMetaData.setNodes(ruleNodes);
ruleChainMetaData.addConnectionInfo(0,1,"success");
ruleChainMetaData.addConnectionInfo(0,2,"fail");
ruleChainMetaData.addConnectionInfo(1,2,"success");
ruleChainMetaData.addConnectionInfo(2,0,"success");
return ruleChainMetaData;
}
}

4
dao/src/test/resources/sql/system-test.sql

@ -1,2 +1,6 @@
TRUNCATE TABLE device_credentials;
TRUNCATE TABLE device;
TRUNCATE TABLE device_profile;
TRUNCATE TABLE rule_node_state;
TRUNCATE TABLE rule_node;
TRUNCATE TABLE rule_chain;

52
netty-mqtt/src/main/java/org/thingsboard/mqtt/MqttClientImpl.java

@ -19,18 +19,40 @@ import com.google.common.collect.HashMultimap;
import com.google.common.collect.ImmutableSet;
import io.netty.bootstrap.Bootstrap;
import io.netty.buffer.ByteBuf;
import io.netty.channel.*;
import io.netty.channel.Channel;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelFutureListener;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.EventLoopGroup;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.handler.codec.mqtt.*;
import io.netty.handler.codec.mqtt.MqttDecoder;
import io.netty.handler.codec.mqtt.MqttEncoder;
import io.netty.handler.codec.mqtt.MqttFixedHeader;
import io.netty.handler.codec.mqtt.MqttMessage;
import io.netty.handler.codec.mqtt.MqttMessageIdVariableHeader;
import io.netty.handler.codec.mqtt.MqttMessageType;
import io.netty.handler.codec.mqtt.MqttPublishMessage;
import io.netty.handler.codec.mqtt.MqttPublishVariableHeader;
import io.netty.handler.codec.mqtt.MqttQoS;
import io.netty.handler.codec.mqtt.MqttSubscribeMessage;
import io.netty.handler.codec.mqtt.MqttSubscribePayload;
import io.netty.handler.codec.mqtt.MqttTopicSubscription;
import io.netty.handler.codec.mqtt.MqttUnsubscribeMessage;
import io.netty.handler.codec.mqtt.MqttUnsubscribePayload;
import io.netty.handler.ssl.SslContext;
import io.netty.handler.timeout.IdleStateHandler;
import io.netty.util.collection.IntObjectHashMap;
import io.netty.util.concurrent.DefaultPromise;
import io.netty.util.concurrent.Future;
import io.netty.util.concurrent.Promise;
import java.util.*;
import java.util.Collections;
import java.util.HashSet;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
@ -41,11 +63,11 @@ import java.util.concurrent.atomic.AtomicInteger;
final class MqttClientImpl implements MqttClient {
private final Set<String> serverSubscriptions = new HashSet<>();
private final IntObjectHashMap<MqttPendingUnsubscription> pendingServerUnsubscribes = new IntObjectHashMap<>();
private final IntObjectHashMap<MqttIncomingQos2Publish> qos2PendingIncomingPublishes = new IntObjectHashMap<>();
private final IntObjectHashMap<MqttPendingPublish> pendingPublishes = new IntObjectHashMap<>();
private final ConcurrentMap<Integer, MqttPendingUnsubscription> pendingServerUnsubscribes = new ConcurrentHashMap<>();
private final ConcurrentMap<Integer, MqttIncomingQos2Publish> qos2PendingIncomingPublishes = new ConcurrentHashMap<>();
private final ConcurrentMap<Integer, MqttPendingPublish> pendingPublishes = new ConcurrentHashMap<>();
private final HashMultimap<String, MqttSubscription> subscriptions = HashMultimap.create();
private final IntObjectHashMap<MqttPendingSubscription> pendingSubscriptions = new IntObjectHashMap<>();
private final ConcurrentMap<Integer, MqttPendingSubscription> pendingSubscriptions = new ConcurrentHashMap<>();
private final Set<String> pendingSubscribeTopics = new HashSet<>();
private final HashMultimap<MqttHandler, MqttSubscription> handlerToSubscribtion = HashMultimap.create();
private final AtomicInteger nextMessageId = new AtomicInteger(1);
@ -340,7 +362,7 @@ final class MqttClientImpl implements MqttClient {
MqttPublishVariableHeader variableHeader = new MqttPublishVariableHeader(topic, getNewMessageId().messageId());
MqttPublishMessage message = new MqttPublishMessage(fixedHeader, variableHeader, payload);
MqttPendingPublish pendingPublish = new MqttPendingPublish(variableHeader.packetId(), future, payload.retain(), message, qos);
this.pendingPublishes.put(pendingPublish.getMessageId(), pendingPublish);
ChannelFuture channelFuture = this.sendAndFlushPacket(message);
if (channelFuture != null) {
@ -351,10 +373,12 @@ final class MqttClientImpl implements MqttClient {
}
}
if (pendingPublish.isSent() && pendingPublish.getQos() == MqttQoS.AT_MOST_ONCE) {
this.pendingPublishes.remove(pendingPublish.getMessageId());
pendingPublish.getFuture().setSuccess(null); //We don't get an ACK for QOS 0
} else if (pendingPublish.isSent()) {
this.pendingPublishes.put(pendingPublish.getMessageId(), pendingPublish);
pendingPublish.startPublishRetransmissionTimer(this.eventLoop.next(), this::sendAndFlushPacket);
} else {
this.pendingPublishes.remove(pendingPublish.getMessageId());
}
return future;
}
@ -466,7 +490,7 @@ final class MqttClientImpl implements MqttClient {
}
}
IntObjectHashMap<MqttPendingSubscription> getPendingSubscriptions() {
ConcurrentMap<Integer, MqttPendingSubscription> getPendingSubscriptions() {
return pendingSubscriptions;
}
@ -486,15 +510,15 @@ final class MqttClientImpl implements MqttClient {
return serverSubscriptions;
}
IntObjectHashMap<MqttPendingUnsubscription> getPendingServerUnsubscribes() {
ConcurrentMap<Integer, MqttPendingUnsubscription> getPendingServerUnsubscribes() {
return pendingServerUnsubscribes;
}
IntObjectHashMap<MqttPendingPublish> getPendingPublishes() {
ConcurrentMap<Integer, MqttPendingPublish> getPendingPublishes() {
return pendingPublishes;
}
IntObjectHashMap<MqttIncomingQos2Publish> getQos2PendingIncomingPublishes() {
ConcurrentMap<Integer, MqttIncomingQos2Publish> getQos2PendingIncomingPublishes() {
return qos2PendingIncomingPublishes;
}

2
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/TbContext.java

@ -146,7 +146,7 @@ public interface TbContext {
TbMsg assetCreatedMsg(Asset asset, RuleNodeId ruleNodeId);
// TODO: Does this changes the message?
TbMsg alarmCreatedMsg(Alarm alarm, RuleNodeId ruleNodeId);
TbMsg alarmActionMsg(Alarm alarm, RuleNodeId ruleNodeId, String action);
/*
*

13
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java

@ -58,13 +58,11 @@ public abstract class TbAbstractAlarmNode<C extends TbAbstractAlarmNodeConfigura
if (alarmResult.alarm == null) {
ctx.tellNext(msg, "False");
} else if (alarmResult.isCreated) {
ctx.enqueue(ctx.alarmCreatedMsg(alarmResult.alarm, ctx.getSelfId()),
() -> ctx.tellNext(toAlarmMsg(ctx, alarmResult, msg), "Created"),
throwable -> ctx.tellFailure(toAlarmMsg(ctx, alarmResult, msg), throwable));
tellNext(ctx, msg, alarmResult, DataConstants.ENTITY_CREATED, "Created");
} else if (alarmResult.isUpdated) {
ctx.tellNext(toAlarmMsg(ctx, alarmResult, msg), "Updated");
tellNext(ctx, msg, alarmResult, DataConstants.ENTITY_UPDATED, "Updated");
} else if (alarmResult.isCleared) {
ctx.tellNext(toAlarmMsg(ctx, alarmResult, msg), "Cleared");
tellNext(ctx, msg, alarmResult, DataConstants.ALARM_CLEAR, "Cleared");
} else {
ctx.tellSuccess(msg);
}
@ -109,4 +107,9 @@ public abstract class TbAbstractAlarmNode<C extends TbAbstractAlarmNodeConfigura
}
}
private void tellNext(TbContext ctx, TbMsg msg, TbAlarmResult alarmResult, String entityAction, String alarmAction) {
ctx.enqueue(ctx.alarmActionMsg(alarmResult.alarm, ctx.getSelfId(), entityAction),
() -> ctx.tellNext(toAlarmMsg(ctx, alarmResult, msg), alarmAction),
throwable -> ctx.tellFailure(toAlarmMsg(ctx, alarmResult, msg), throwable));
}
}

9
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractRelationActionNode.java

@ -35,6 +35,7 @@ import org.thingsboard.server.common.data.DashboardInfo;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.EntityView;
import org.thingsboard.server.common.data.User;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory;
@ -49,6 +50,7 @@ import org.thingsboard.server.dao.customer.CustomerService;
import org.thingsboard.server.dao.dashboard.DashboardService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.entityview.EntityViewService;
import org.thingsboard.server.dao.user.UserService;
import java.util.List;
import java.util.Optional;
@ -245,6 +247,13 @@ public abstract class TbAbstractRelationActionNode<C extends TbAbstractRelationA
}
}
break;
case USER:
UserService userService = ctx.getUserService();
User user = userService.findUserByEmail(ctx.getTenantId(), entitykey.getEntityName());
if(user != null){
targetEntity.setEntityId(user.getId());
}
break;
default:
return targetEntity;
}

2
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

7
rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbAlarmNodeTest.java

@ -35,6 +35,7 @@ import org.thingsboard.rule.engine.api.ScriptEngine;
import org.thingsboard.rule.engine.api.TbContext;
import org.thingsboard.rule.engine.api.TbNodeConfiguration;
import org.thingsboard.rule.engine.api.TbNodeException;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.alarm.Alarm;
import org.thingsboard.server.common.data.id.AlarmId;
import org.thingsboard.server.common.data.id.DeviceId;
@ -249,6 +250,8 @@ public class TbAlarmNodeTest {
node.onMsg(ctx, msg);
verify(ctx).enqueue(any(), successCaptor.capture(), failureCaptor.capture());
successCaptor.getValue().run();
verify(ctx).tellNext(any(), eq("Updated"));
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
@ -298,6 +301,8 @@ public class TbAlarmNodeTest {
node.onMsg(ctx, msg);
verify(ctx).enqueue(any(), successCaptor.capture(), failureCaptor.capture());
successCaptor.getValue().run();
verify(ctx).tellNext(any(), eq("Cleared"));
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);
@ -346,6 +351,8 @@ public class TbAlarmNodeTest {
node.onMsg(ctx, msg);
verify(ctx).enqueue(any(), successCaptor.capture(), failureCaptor.capture());
successCaptor.getValue().run();
verify(ctx).tellNext(any(), eq("Cleared"));
ArgumentCaptor<TbMsg> msgCaptor = ArgumentCaptor.forClass(TbMsg.class);

4
ui-ngx/src/app/modules/home/components/profile/add-device-profile-dialog.component.html

@ -45,6 +45,10 @@
labelText="device-profile.default-rule-chain"
formControlName="defaultRuleChainId">
</tb-rule-chain-autocomplete>
<tb-queue-type-list
[queueType]="serviceType"
formControlName="defaultQueueName">
</tb-queue-type-list>
<mat-form-field fxHide class="mat-block">
<mat-label translate>device-profile.type</mat-label>
<mat-select formControlName="type" required>

4
ui-ngx/src/app/modules/home/components/profile/add-device-profile-dialog.component.ts

@ -48,6 +48,7 @@ import { MatHorizontalStepper } from '@angular/material/stepper';
import { RuleChainId } from '@shared/models/id/rule-chain-id';
import { StepperSelectionEvent } from '@angular/cdk/stepper';
import { deepTrim } from '@core/utils';
import {ServiceType} from "@shared/models/queue.models";
export interface AddDeviceProfileDialogData {
deviceProfileName: string;
@ -89,6 +90,8 @@ export class AddDeviceProfileDialogComponent extends
provisionConfigFormGroup: FormGroup;
serviceType = ServiceType.TB_RULE_ENGINE;
constructor(protected store: Store<AppState>,
protected router: Router,
@Inject(MAT_DIALOG_DATA) public data: AddDeviceProfileDialogData,
@ -104,6 +107,7 @@ export class AddDeviceProfileDialogComponent extends
name: [data.deviceProfileName, [Validators.required]],
type: [DeviceProfileType.DEFAULT, [Validators.required]],
defaultRuleChainId: [null, []],
defaultQueueName: ['', []],
description: ['', []]
}
);

4
ui-ngx/src/app/modules/home/components/profile/device-profile.component.html

@ -53,6 +53,10 @@
labelText="device-profile.default-rule-chain"
formControlName="defaultRuleChainId">
</tb-rule-chain-autocomplete>
<tb-queue-type-list
[queueType]="serviceType"
formControlName="defaultQueueName">
</tb-queue-type-list>
<mat-form-field fxHide class="mat-block">
<mat-label translate>device-profile.type</mat-label>
<mat-select formControlName="type" required>

5
ui-ngx/src/app/modules/home/components/profile/device-profile.component.ts

@ -38,6 +38,7 @@ import {
} from '@shared/models/device.models';
import { EntityType } from '@shared/models/entity-type.models';
import { RuleChainId } from '@shared/models/id/rule-chain-id';
import {ServiceType} from "@shared/models/queue.models";
@Component({
selector: 'tb-device-profile',
@ -63,6 +64,8 @@ export class DeviceProfileComponent extends EntityComponent<DeviceProfile> {
displayTransportConfiguration: boolean;
serviceType = ServiceType.TB_RULE_ENGINE;
constructor(protected store: Store<AppState>,
protected translate: TranslateService,
@Optional() @Inject('entity') protected entityValue: DeviceProfile,
@ -101,6 +104,7 @@ export class DeviceProfileComponent extends EntityComponent<DeviceProfile> {
provisionConfiguration: [deviceProvisionConfiguration, Validators.required]
}),
defaultRuleChainId: [entity && entity.defaultRuleChainId ? entity.defaultRuleChainId.id : null, []],
defaultQueueName: [entity ? entity.defaultQueueName : '', []],
description: [entity ? entity.description : '', []],
}
);
@ -174,6 +178,7 @@ export class DeviceProfileComponent extends EntityComponent<DeviceProfile> {
provisionConfiguration: deviceProvisionConfiguration
}});
this.entityForm.patchValue({defaultRuleChainId: entity.defaultRuleChainId ? entity.defaultRuleChainId.id : null});
this.entityForm.patchValue({defaultQueueName: entity.defaultQueueName});
this.entityForm.patchValue({description: entity.description});
}

4
ui-ngx/src/app/modules/home/components/rule-chain/rule-chain-autocomplete.component.ts

@ -202,8 +202,8 @@ export class RuleChainAutocompleteComponent implements ControlValueAccessor, OnI
createDefaultRuleChain($event: Event, ruleChainName: string) {
$event.preventDefault();
this.ruleChainAutocomplete.closePanel();
this.ruleChainService.createDefaultRuleChain(ruleChainName).subscribe((ruleChain) => {
this.updateView(ruleChain.id.id);
this.ruleChainService.createDefaultRuleChain(ruleChainName.trim()).subscribe((ruleChain) => {
this.selectRuleChainFormGroup.get('ruleChainId').patchValue(ruleChain);
});
}
}

9
ui-ngx/src/app/modules/home/components/widget/lib/alarms-table-widget.component.ts

@ -29,7 +29,7 @@ import { PageComponent } from '@shared/components/page.component';
import { Store } from '@ngrx/store';
import { AppState } from '@core/core.state';
import { WidgetAction, WidgetContext } from '@home/models/widget-component.models';
import { DataKey, Datasource, WidgetActionDescriptor, WidgetConfig } from '@shared/models/widget.models';
import { DataKey, WidgetActionDescriptor, WidgetConfig } from '@shared/models/widget.models';
import { IWidgetSubscription } from '@core/api/widget-api.models';
import { UtilsService } from '@core/services/utils.service';
import { TranslateService } from '@ngx-translate/core';
@ -371,7 +371,8 @@ export class AlarmsTableWidgetComponent extends PageComponent implements OnInit,
this.subscription.alarmSource.dataKeys.forEach((alarmDataKey) => {
const dataKey: EntityColumn = deepClone(alarmDataKey) as EntityColumn;
dataKey.entityKey = dataKeyToEntityKey(alarmDataKey);
dataKey.title = this.utils.customTranslation(dataKey.label, dataKey.label);
dataKey.label = this.utils.customTranslation(dataKey.label, dataKey.label);
dataKey.title = dataKey.label;
dataKey.def = 'def' + this.columns.length;
const keySettings: TableWidgetDataKeySettings = dataKey.settings;
if (dataKey.type === DataKeyType.alarm && !isDefined(keySettings.columnWidth)) {
@ -394,7 +395,7 @@ export class AlarmsTableWidgetComponent extends PageComponent implements OnInit,
this.displayedColumns.push(...this.columns.map(column => column.def));
}
if (this.settings.defaultSortOrder && this.settings.defaultSortOrder.length) {
this.defaultSortOrder = this.settings.defaultSortOrder;
this.defaultSortOrder = this.utils.customTranslation(this.settings.defaultSortOrder, this.settings.defaultSortOrder);
}
this.pageLink.sortOrder = entityDataSortOrderFromString(this.defaultSortOrder, this.columns);
let sortColumn: EntityColumn;
@ -959,7 +960,7 @@ class AlarmsDatasource implements DataSource<AlarmDataInfo> {
}
}
}
alarm[dataKey.name] = value;
alarm[dataKey.label] = value;
});
return alarm;
}

2
ui-ngx/src/app/modules/home/components/widget/lib/table-widget.models.ts

@ -189,7 +189,7 @@ export function getAlarmValue(alarm: AlarmDataInfo, key: EntityColumn) {
if (alarmField) {
return getDescendantProp(alarm, alarmField.value);
} else {
return getDescendantProp(alarm, key.name);
return getDescendantProp(alarm, key.label);
}
}

22
ui-ngx/src/app/modules/home/components/widget/lib/timeseries-table-widget.component.ts

@ -524,16 +524,22 @@ class TimeseriesDatasource implements DataSource<TimeseriesRow> {
});
}
const rows: TimeseriesRow[] = [];
for (const value of Object.values(rowsMap)) {
if (this.hideEmptyLines && isDefinedAndNotNull(value[1])) {
rows.push(value);
} else {
rows.push(value);
let rows: TimeseriesRow[] = [];
if (this.hideEmptyLines) {
for (const t of Object.keys(rowsMap)) {
let hideLine = true;
for (let c = 0; (c < data.length) && hideLine; c++) {
if (rowsMap[t][c + 1]) {
hideLine = false;
}
}
if (!hideLine) {
rows.push(rowsMap[t]);
}
}
} else {
rows = Object.values(rowsMap);
}
return rows;
}

6
ui-ngx/src/app/modules/home/components/wizard/device-wizard-dialog.component.html

@ -97,6 +97,12 @@
formControlName="defaultRuleChainId">
</tb-rule-chain-autocomplete>
</div>
<div fxLayout="column" fxLayoutAlign="flex-end start">
<tb-queue-type-list
[queueType]="serviceType"
formControlName="defaultQueueName">
</tb-queue-type-list>
</div>
</div>
<mat-checkbox formControlName="gateway" style="padding-bottom: 16px;">
{{ 'device.is-gateway' | translate }}

20
ui-ngx/src/app/modules/home/components/wizard/device-wizard-dialog.component.ts

@ -25,8 +25,12 @@ import {
createDeviceProfileConfiguration,
createDeviceProfileTransportConfiguration,
DeviceProfile,
DeviceProfileType, DeviceProvisionConfiguration, DeviceProvisionType,
DeviceTransportType, deviceTransportTypeConfigurationInfoMap, deviceTransportTypeHintMap,
DeviceProfileType,
DeviceProvisionConfiguration,
DeviceProvisionType,
DeviceTransportType,
deviceTransportTypeConfigurationInfoMap,
deviceTransportTypeHintMap,
deviceTransportTypeTranslationMap
} from '@shared/models/device.models';
import { MatHorizontalStepper } from '@angular/material/stepper';
@ -43,6 +47,8 @@ import { StepperSelectionEvent } from '@angular/cdk/stepper';
import { BreakpointObserver, BreakpointState } from '@angular/cdk/layout';
import { MediaBreakpoints } from '@shared/models/constants';
import { RuleChainId } from '@shared/models/id/rule-chain-id';
import { ServiceType } from '@shared/models/queue.models';
import { deepTrim } from '@core/utils';
@Component({
selector: 'tb-device-wizard',
@ -84,6 +90,8 @@ export class DeviceWizardDialogComponent extends
labelPosition = 'end';
serviceType = ServiceType.TB_RULE_ENGINE;
private subscriptions: Subscription[] = [];
constructor(protected store: Store<AppState>,
@ -105,6 +113,7 @@ export class DeviceWizardDialogComponent extends
deviceProfileId: [null, Validators.required],
newDeviceProfileTitle: [{value: null, disabled: true}],
defaultRuleChainId: [{value: null, disabled: true}],
defaultQueueName: [{value: null, disabled: true}],
description: ['']
}
);
@ -117,6 +126,7 @@ export class DeviceWizardDialogComponent extends
this.deviceWizardFormGroup.get('newDeviceProfileTitle').setValidators(null);
this.deviceWizardFormGroup.get('newDeviceProfileTitle').disable();
this.deviceWizardFormGroup.get('defaultRuleChainId').disable();
this.deviceWizardFormGroup.get('defaultQueueName').disable();
this.deviceWizardFormGroup.updateValueAndValidity();
this.createProfile = false;
this.createTransportConfiguration = false;
@ -126,6 +136,8 @@ export class DeviceWizardDialogComponent extends
this.deviceWizardFormGroup.get('newDeviceProfileTitle').setValidators([Validators.required]);
this.deviceWizardFormGroup.get('newDeviceProfileTitle').enable();
this.deviceWizardFormGroup.get('defaultRuleChainId').enable();
this.deviceWizardFormGroup.get('defaultQueueName').enable();
this.deviceWizardFormGroup.updateValueAndValidity();
this.createProfile = true;
this.createTransportConfiguration = this.deviceWizardFormGroup.get('transportType').value &&
@ -281,7 +293,7 @@ export class DeviceWizardDialogComponent extends
if (this.deviceWizardFormGroup.get('defaultRuleChainId').value) {
deviceProfile.defaultRuleChainId = new RuleChainId(this.deviceWizardFormGroup.get('defaultRuleChainId').value);
}
return this.deviceProfileService.saveDeviceProfile(deviceProfile).pipe(
return this.deviceProfileService.saveDeviceProfile(deepTrim(deviceProfile)).pipe(
map(profile => profile.id),
tap((profileId) => {
this.deviceWizardFormGroup.patchValue({
@ -312,7 +324,7 @@ export class DeviceWizardDialogComponent extends
id: this.customerFormGroup.get('customerId').value
};
}
return this.data.entitiesTableConfig.saveEntity(device);
return this.data.entitiesTableConfig.saveEntity(deepTrim(device));
}
private saveCredentials(device: BaseData<HasId>): Observable<boolean> {

4
ui-ngx/src/app/modules/home/pages/admin/general-settings.component.html

@ -35,6 +35,10 @@
{{ 'admin.base-url-required' | translate }}
</mat-error>
</mat-form-field>
<tb-checkbox formControlName="prohibitDifferentUrl" style="display: block; padding-bottom: 16px;">
{{ 'admin.prohibit-different-url' | translate }}
</tb-checkbox>
<div class="tb-hint" translate>admin.prohibit-different-url-hint</div>
<div fxLayout="row" fxLayoutAlign="end center" style="width: 100%;" class="layout-wrap">
<button mat-button mat-raised-button color="primary" [disabled]="(isLoading$ | async) || generalSettings.invalid || !generalSettings.dirty"
type="submit">{{'action.save' | translate}}

3
ui-ngx/src/app/modules/home/pages/admin/general-settings.component.ts

@ -53,7 +53,8 @@ export class GeneralSettingsComponent extends PageComponent implements OnInit, H
buildGeneralServerSettingsForm() {
this.generalSettings = this.fb.group({
baseUrl: ['', [Validators.required]]
baseUrl: ['', [Validators.required]],
prohibitDifferentUrl: ['',[]]
});
}

18
ui-ngx/src/app/modules/home/pages/admin/oauth2-settings.component.ts

@ -142,6 +142,7 @@ export class OAuth2SettingsComponent extends PageComponent implements OnInit, Ha
tenantNamePattern = {value: null, disabled: true};
}
const basicGroup = this.fb.group({
emailAttributeKey: [mapperConfigBasic?.emailAttributeKey ? mapperConfigBasic.emailAttributeKey : 'email', Validators.required],
firstNameAttributeKey: [mapperConfigBasic?.firstNameAttributeKey ? mapperConfigBasic.firstNameAttributeKey : ''],
lastNameAttributeKey: [mapperConfigBasic?.lastNameAttributeKey ? mapperConfigBasic.lastNameAttributeKey : ''],
tenantNameStrategy: [mapperConfigBasic?.tenantNameStrategy ? mapperConfigBasic.tenantNameStrategy : TenantNameStrategy.DOMAIN],
@ -151,11 +152,6 @@ export class OAuth2SettingsComponent extends PageComponent implements OnInit, Ha
alwaysFullScreen: [isDefinedAndNotNull(mapperConfigBasic?.alwaysFullScreen) ? mapperConfigBasic.alwaysFullScreen : false]
});
if (MapperConfigType.GITHUB !== type) {
basicGroup.addControl('emailAttributeKey',
this.fb.control( mapperConfigBasic?.emailAttributeKey ? mapperConfigBasic.emailAttributeKey : 'email', Validators.required));
}
this.subscriptions.push(basicGroup.get('tenantNameStrategy').valueChanges.subscribe((domain) => {
if (domain === 'CUSTOM') {
basicGroup.get('tenantNamePattern').enable();
@ -347,7 +343,7 @@ export class OAuth2SettingsComponent extends PageComponent implements OnInit, Ha
clientRegistration.get('authorizationUri').disable();
clientRegistration.get('jwkSetUri').disable();
clientRegistration.get('userInfoUri').disable();
clientRegistration.patchValue(template, {emitEvent: false});
clientRegistration.patchValue(template);
}
}
@ -358,11 +354,15 @@ export class OAuth2SettingsComponent extends PageComponent implements OnInit, Ha
mapperConfig.addControl('custom', this.formCustomGroup(predefinedValue?.custom));
} else {
mapperConfig.removeControl('custom');
if (mapperConfig.get('basic')) {
mapperConfig.setControl('basic', this.formBasicGroup(type, predefinedValue?.basic));
} else {
if (!mapperConfig.get('basic')) {
mapperConfig.addControl('basic', this.formBasicGroup(type, predefinedValue?.basic));
}
if (type === MapperConfigType.GITHUB) {
mapperConfig.get('basic.emailAttributeKey').disable();
mapperConfig.get('basic.emailAttributeKey').patchValue(null, {emitEvent: false});
} else {
mapperConfig.get('basic.emailAttributeKey').enable();
}
}
}

1
ui-ngx/src/app/shared/components/queue/queue-type-list.component.html

@ -40,4 +40,5 @@
<mat-error *ngIf="queueFormGroup.get('queue').hasError('required')">
{{ 'queue.name_required' | translate }}
</mat-error>
<mat-hint class="tb-hint" translate [fxShow]="!disabled">device-profile.select-queue-hint</mat-hint>
</mat-form-field>

1
ui-ngx/src/app/shared/models/device.models.ts

@ -369,6 +369,7 @@ export interface DeviceProfile extends BaseData<DeviceProfileId> {
provisionType: DeviceProvisionType;
provisionDeviceKey?: string;
defaultRuleChainId?: RuleChainId;
defaultQueueName?: string;
profileData: DeviceProfileData;
}

6
ui-ngx/src/assets/locale/locale.constant-en_US.json

@ -80,6 +80,8 @@
"test-mail-sent": "Test mail was successfully sent!",
"base-url": "Base URL",
"base-url-required": "Base URL is required.",
"prohibit-different-url": "Prohibit to use hostname from the client request headers",
"prohibit-different-url-hint": "This setting should be enabled for production environments. May cause security issues when disabled",
"mail-from": "Mail From",
"mail-from-required": "Mail From is required.",
"smtp-protocol": "SMTP protocol",
@ -873,6 +875,7 @@
"profile-configuration": "Profile configuration",
"transport-configuration": "Transport configuration",
"default-rule-chain": "Default rule chain",
"select-queue-hint": "The queue name can be selected from a drop-down list or add a custom name.",
"delete-device-profile-title": "Are you sure you want to delete the device profile '{{deviceProfileName}}'?",
"delete-device-profile-text": "Be careful, after the confirmation the device profile and all related data will become unrecoverable.",
"delete-device-profiles-title": "Are you sure you want to delete { count, plural, 1 {1 device profile} other {# device profiles} }?",
@ -2371,7 +2374,8 @@
"el_GR": "Ελληνικά",
"ro_RO": "Română",
"lv_LV": "Latviešu",
"ka_GE": "ქართული"
"ka_GE": "ქართული",
"pt_BR": "Português do Brasil"
}
}
}

2048
ui-ngx/src/assets/locale/locale.constant-pl_BR.json

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