Browse Source

Merge remote-tracking branch 'origin/master'

pull/3703/head
Andrii Shvaika 6 years ago
parent
commit
e981edd83c
  1. 2
      application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
  2. 3
      application/src/main/java/org/thingsboard/server/actors/app/AppActor.java
  3. 3
      application/src/main/java/org/thingsboard/server/controller/BaseController.java
  4. 5
      application/src/main/java/org/thingsboard/server/controller/TelemetryController.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
  6. 2
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  7. 2
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java
  8. 3
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTenantRoutingInfoService.java
  9. 2
      application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java
  10. 2
      application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java
  11. 2
      common/dao-api/src/main/java/org/thingsboard/server/dao/tenant/TbTenantProfileCache.java
  12. 5
      common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java
  13. 8
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  14. 2
      dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java
  15. 16
      dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java
  16. 2
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java
  17. 18
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java
  18. 1
      dao/src/main/java/org/thingsboard/server/dao/sql/asset/AssetRepository.java
  19. 6
      dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java
  20. 2
      dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceRepository.java
  21. 5
      dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java
  22. 9
      dao/src/main/java/org/thingsboard/server/dao/tenant/DefaultTbTenantProfileCache.java
  23. 1
      docker/queue-kafka.env
  24. 6
      msa/js-executor/queue/awsSqsTemplate.js
  25. 54
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java

2
application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java

@ -73,7 +73,7 @@ import org.thingsboard.server.service.executors.ExternalCallExecutorService;
import org.thingsboard.server.service.executors.SharedEventLoopGroupService;
import org.thingsboard.server.service.mail.MailExecutorService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.queue.TbClusterService;
import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService;
import org.thingsboard.server.service.rpc.TbRuleEngineDeviceRpcService;

3
application/src/main/java/org/thingsboard/server/actors/app/AppActor.java

@ -30,7 +30,6 @@ import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.common.data.page.PageDataIterable;
import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent;
import org.thingsboard.server.common.msg.MsgType;
@ -41,7 +40,7 @@ import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg;
import org.thingsboard.server.common.msg.queue.RuleEngineException;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper;
import java.util.HashSet;

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

@ -35,7 +35,6 @@ import org.thingsboard.server.common.data.asset.AssetInfo;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.exception.ThingsboardErrorCode;
import org.thingsboard.server.common.data.exception.ThingsboardException;
import org.thingsboard.server.common.data.id.*;
import org.thingsboard.server.common.data.id.AlarmId;
import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
@ -94,7 +93,7 @@ import org.thingsboard.server.queue.provider.TbQueueProducerProvider;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.component.ComponentDiscoveryService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.queue.TbClusterService;
import org.thingsboard.server.service.security.model.SecurityUser;
import org.thingsboard.server.service.security.permission.AccessControlService;

5
application/src/main/java/org/thingsboard/server/controller/TelemetryController.java

@ -392,6 +392,11 @@ public class TelemetryController extends BaseController {
if (attributes.isEmpty()) {
return getImmediateDeferredResult("No attributes data found in request body!", HttpStatus.BAD_REQUEST);
}
for (AttributeKvEntry attributeKvEntry: attributes) {
if (attributeKvEntry.getKey().isEmpty() || attributeKvEntry.getKey().trim().length() == 0) {
return getImmediateDeferredResult("Key cannot be empty or contains only spaces", HttpStatus.BAD_REQUEST);
}
}
SecurityUser user = getCurrentUser();
return accessValidator.validateEntityAndCallback(getCurrentUser(), Operation.WRITE_ATTRIBUTES, entityIdSrc, (result, tenantId, entityId) -> {
tsSubService.saveAndNotify(tenantId, entityId, scope, attributes, new FutureCallback<Void>() {

2
application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java

@ -55,7 +55,7 @@ import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
import org.thingsboard.server.queue.discovery.PartitionService;
import org.thingsboard.server.queue.scheduler.SchedulerComponent;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.queue.TbClusterService;
import org.thingsboard.server.service.telemetry.InternalTelemetryService;

2
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -55,7 +55,7 @@ import org.thingsboard.server.queue.provider.TbCoreQueueFactory;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.queue.processing.AbstractConsumerService;
import org.thingsboard.server.service.rpc.FromDeviceRpcResponse;
import org.thingsboard.server.service.rpc.TbCoreDeviceRpcService;

2
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java

@ -45,7 +45,7 @@ import org.thingsboard.server.queue.settings.TbRuleEngineQueueConfiguration;
import org.thingsboard.server.queue.util.TbRuleEngineComponent;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.queue.processing.AbstractConsumerService;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingDecision;
import org.thingsboard.server.service.queue.processing.TbRuleEngineProcessingResult;

3
application/src/main/java/org/thingsboard/server/service/queue/DefaultTenantRoutingInfoService.java

@ -21,11 +21,10 @@ import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.dao.tenant.TenantProfileService;
import org.thingsboard.server.dao.tenant.TenantService;
import org.thingsboard.server.queue.discovery.TenantRoutingInfo;
import org.thingsboard.server.queue.discovery.TenantRoutingInfoService;
import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
@Slf4j
@Service

2
application/src/main/java/org/thingsboard/server/service/queue/processing/AbstractConsumerService.java

@ -38,7 +38,7 @@ import org.thingsboard.server.queue.discovery.PartitionChangeEvent;
import org.thingsboard.server.common.transport.util.DataDecodingEncodingService;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.queue.TbPackCallback;
import org.thingsboard.server.service.queue.TbPackProcessingContext;

2
application/src/main/java/org/thingsboard/server/service/transport/DefaultTransportApiService.java

@ -71,7 +71,7 @@ import org.thingsboard.server.dao.device.provision.ProvisionFailedException;
import org.thingsboard.server.service.apiusage.TbApiUsageStateService;
import org.thingsboard.server.service.executors.DbCallbackExecutorService;
import org.thingsboard.server.service.profile.TbDeviceProfileCache;
import org.thingsboard.server.service.profile.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.service.queue.TbClusterService;
import org.thingsboard.server.service.state.DeviceStateService;

2
application/src/main/java/org/thingsboard/server/service/profile/TbTenantProfileCache.java → common/dao-api/src/main/java/org/thingsboard/server/dao/tenant/TbTenantProfileCache.java

@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.profile;
package org.thingsboard.server.dao.tenant;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.TenantId;

5
common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java

@ -87,8 +87,9 @@ public class TbAwsSqsProducerTemplate<T extends TbQueueMsg> implements TbQueuePr
sendMsgRequest.withQueueUrl(getQueueUrl(tpi.getFullTopicName()));
sendMsgRequest.withMessageBody(gson.toJson(new DefaultTbQueueMsg(msg)));
sendMsgRequest.withMessageGroupId(tpi.getTopic());
sendMsgRequest.withMessageDeduplicationId(UUID.randomUUID().toString());
String sqsMsgId = UUID.randomUUID().toString();
sendMsgRequest.withMessageGroupId(sqsMsgId);
sendMsgRequest.withMessageDeduplicationId(sqsMsgId);
ListenableFuture<SendMessageResult> future = producerExecutor.submit(() -> sqsClient.sendMessage(sendMsgRequest));

8
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -122,7 +122,13 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
log.trace("[{}] Processing msg: {}", sessionId, msg);
try {
if (msg instanceof MqttMessage) {
processMqttMsg(ctx, (MqttMessage) msg);
MqttMessage message = (MqttMessage) msg;
if (message.decoderResult().isSuccess()) {
processMqttMsg(ctx, message);
} else {
log.error("[{}] Message processing failed: {}", sessionId, message.decoderResult().cause().getMessage());
ctx.close();
}
} else {
ctx.close();
}

2
dao/src/main/java/org/thingsboard/server/dao/asset/AssetDao.java

@ -166,4 +166,6 @@ public interface AssetDao extends Dao<Asset> {
*/
ListenableFuture<List<EntitySubtype>> findTenantAssetTypesAsync(UUID tenantId);
Long countAssetsByTenantId(TenantId tenantId);
}

16
dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java

@ -26,6 +26,7 @@ import org.springframework.cache.Cache;
import org.springframework.cache.CacheManager;
import org.springframework.cache.annotation.CacheEvict;
import org.springframework.cache.annotation.Cacheable;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
import org.springframework.util.StringUtils;
import org.thingsboard.server.common.data.Customer;
@ -44,12 +45,14 @@ import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.dao.customer.CustomerDao;
import org.thingsboard.server.dao.entity.AbstractEntityService;
import org.thingsboard.server.dao.entityview.EntityViewService;
import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.dao.service.PaginatedRemover;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TenantDao;
import java.util.ArrayList;
@ -90,6 +93,10 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
@Autowired
private CacheManager cacheManager;
@Autowired
@Lazy
private TbTenantProfileCache tenantProfileCache;
@Override
public AssetInfo findAssetInfoById(TenantId tenantId, AssetId assetId) {
log.trace("Executing findAssetInfoById [{}]", assetId);
@ -320,6 +327,15 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ
@Override
protected void validateCreate(TenantId tenantId, Asset asset) {
DefaultTenantProfileConfiguration profileConfiguration =
(DefaultTenantProfileConfiguration)tenantProfileCache.get(tenantId).getProfileData().getConfiguration();
long maxAssets = profileConfiguration.getMaxAssets();
if (maxAssets > 0) {
long currentAssetsCount = assetDao.countAssetsByTenantId(tenantId);
if (maxAssets >= currentAssetsCount) {
throw new DataValidationException("Can't create assets more then " + maxAssets);
}
}
}
@Override

2
dao/src/main/java/org/thingsboard/server/dao/device/DeviceDao.java

@ -203,6 +203,8 @@ public interface DeviceDao extends Dao<Device> {
*/
ListenableFuture<Device> findDeviceByTenantIdAndIdAsync(TenantId tenantId, UUID id);
Long countDevicesByTenantId(TenantId tenantId);
Long countDevicesByDeviceProfileId(TenantId tenantId, UUID deviceProfileId);
/**

18
dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java

@ -27,6 +27,7 @@ import org.springframework.cache.Cache;
import org.springframework.cache.CacheManager;
import org.springframework.cache.annotation.CacheEvict;
import org.springframework.cache.annotation.Cacheable;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.util.CollectionUtils;
@ -57,6 +58,7 @@ import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import org.thingsboard.server.common.data.security.DeviceCredentials;
import org.thingsboard.server.common.data.security.DeviceCredentialsType;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.dao.customer.CustomerDao;
import org.thingsboard.server.dao.device.provision.ProvisionFailedException;
import org.thingsboard.server.dao.device.provision.ProvisionRequest;
@ -67,6 +69,7 @@ import org.thingsboard.server.dao.event.EventService;
import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.dao.service.PaginatedRemover;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TenantDao;
import org.thingsboard.server.dao.util.mapping.JacksonUtil;
@ -120,6 +123,10 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe
@Autowired
private EventService eventService;
@Autowired
@Lazy
private TbTenantProfileCache tenantProfileCache;
@Override
public DeviceInfo findDeviceInfoById(TenantId tenantId, DeviceId deviceId) {
log.trace("Executing findDeviceInfoById [{}]", deviceId);
@ -520,6 +527,15 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe
@Override
protected void validateCreate(TenantId tenantId, Device device) {
DefaultTenantProfileConfiguration profileConfiguration =
(DefaultTenantProfileConfiguration)tenantProfileCache.get(tenantId).getProfileData().getConfiguration();
long maxDevices = profileConfiguration.getMaxDevices();
if (maxDevices > 0) {
long currentDevicesCount = deviceDao.countDevicesByTenantId(tenantId);
if (maxDevices >= currentDevicesCount) {
throw new DataValidationException("Can't create devices more then " + maxDevices);
}
}
}
@Override
@ -532,7 +548,7 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe
@Override
protected void validateDataImpl(TenantId tenantId, Device device) {
if (StringUtils.isEmpty(device.getName())) {
if (StringUtils.isEmpty(device.getName()) || device.getName().trim().length() == 0) {
throw new DataValidationException("Device name should be specified!");
}
if (device.getTenantId() == null) {

1
dao/src/main/java/org/thingsboard/server/dao/sql/asset/AssetRepository.java

@ -122,4 +122,5 @@ public interface AssetRepository extends PagingAndSortingRepository<AssetEntity,
@Query("SELECT DISTINCT a.type FROM AssetEntity a WHERE a.tenantId = :tenantId")
List<String> findTenantAssetTypes(@Param("tenantId") UUID tenantId);
Long countByTenantId(UUID tenantId);
}

6
dao/src/main/java/org/thingsboard/server/dao/sql/asset/JpaAssetDao.java

@ -176,4 +176,10 @@ public class JpaAssetDao extends JpaAbstractSearchTextDao<AssetEntity, Asset> im
}
return list;
}
@Override
public Long countAssetsByTenantId(TenantId tenantId) {
return assetRepository.countByTenantId(tenantId.getId());
}
}

2
dao/src/main/java/org/thingsboard/server/dao/sql/device/DeviceRepository.java

@ -168,4 +168,6 @@ public interface DeviceRepository extends PagingAndSortingRepository<DeviceEntit
DeviceEntity findByTenantIdAndId(UUID tenantId, UUID id);
Long countByDeviceProfileId(UUID deviceProfileId);
Long countByTenantId(UUID tenantId);
}

5
dao/src/main/java/org/thingsboard/server/dao/sql/device/JpaDeviceDao.java

@ -219,6 +219,11 @@ public class JpaDeviceDao extends JpaAbstractSearchTextDao<DeviceEntity, Device>
return deviceRepository.countByDeviceProfileId(deviceProfileId);
}
@Override
public Long countDevicesByTenantId(TenantId tenantId) {
return deviceRepository.countByTenantId(tenantId.getId());
}
private List<EntitySubtype> convertTenantDeviceTypesToDto(UUID tenantId, List<String> types) {
List<EntitySubtype> list = Collections.emptyList();
if (types != null && !types.isEmpty()) {

9
application/src/main/java/org/thingsboard/server/service/profile/DefaultTbTenantProfileCache.java → dao/src/main/java/org/thingsboard/server/dao/tenant/DefaultTbTenantProfileCache.java

@ -13,20 +13,15 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.service.profile;
package org.thingsboard.server.dao.tenant;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.DeviceProfileId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.id.TenantProfileId;
import org.thingsboard.server.dao.device.DeviceProfileService;
import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.tenant.TbTenantProfileCache;
import org.thingsboard.server.dao.tenant.TenantProfileService;
import org.thingsboard.server.dao.tenant.TenantService;

1
docker/queue-kafka.env

@ -1,3 +1,2 @@
TB_QUEUE_TYPE=kafka
TB_KAFKA_SERVERS=kafka:9092
TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES=retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600;partitions:100

6
msa/js-executor/queue/awsSqsTemplate.js

@ -52,11 +52,13 @@ function AwsSqsProducer() {
queueUrls.set(responseTopic, responseQueueUrl);
}
let msgId = uuid();
let params = {
MessageBody: msgBody,
QueueUrl: responseQueueUrl,
MessageGroupId: 'js_eval',
MessageDeduplicationId: uuid()
MessageGroupId: msgId,
MessageDeduplicationId: msgId
};
return new Promise((resolve, reject) => {

54
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/profile/AlarmState.java

@ -16,6 +16,7 @@
package org.thingsboard.rule.engine.profile;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.ObjectNode;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.rule.engine.action.TbAlarmResult;
@ -29,6 +30,7 @@ import org.thingsboard.server.common.data.alarm.AlarmStatus;
import org.thingsboard.server.common.data.device.profile.DeviceProfileAlarm;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.query.EntityKeyType;
import org.thingsboard.server.common.data.query.KeyFilter;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.queue.ServiceQueue;
@ -53,6 +55,7 @@ class AlarmState {
private volatile boolean initialFetchDone;
private volatile TbMsgMetaData lastMsgMetaData;
private volatile String lastMsgQueueName;
private volatile DataSnapshot dataSnapshot;
AlarmState(ProfileState deviceProfile, EntityId originator, DeviceProfileAlarm alarmDefinition, PersistedAlarmState alarmState) {
this.deviceProfile = deviceProfile;
@ -74,7 +77,7 @@ class AlarmState {
public <T> boolean createOrClearAlarms(TbContext ctx, T data, SnapshotUpdate update, BiFunction<AlarmRuleState, T, AlarmEvalResult> evalFunction) {
boolean stateUpdate = false;
AlarmSeverity resultSeverity = null;
AlarmRuleState resultState = null;
log.debug("[{}] processing update: {}", alarmDefinition.getId(), data);
for (AlarmRuleState state : createRulesSortedBySeverityDesc) {
if (!validateUpdate(update, state)) {
@ -84,15 +87,15 @@ class AlarmState {
AlarmEvalResult evalResult = evalFunction.apply(state, data);
stateUpdate |= state.checkUpdate();
if (AlarmEvalResult.TRUE.equals(evalResult)) {
resultSeverity = state.getSeverity();
resultState = state;
break;
} else if (AlarmEvalResult.FALSE.equals(evalResult)) {
state.clear();
stateUpdate |= state.checkUpdate();
}
}
if (resultSeverity != null) {
TbAlarmResult result = calculateAlarmResult(ctx, resultSeverity);
if (resultState != null) {
TbAlarmResult result = calculateAlarmResult(ctx, resultState);
if (result != null) {
pushMsg(ctx, result);
}
@ -187,7 +190,8 @@ class AlarmState {
}
}
private TbAlarmResult calculateAlarmResult(TbContext ctx, AlarmSeverity severity) {
private <T> TbAlarmResult calculateAlarmResult(TbContext ctx, AlarmRuleState ruleState) {
AlarmSeverity severity = ruleState.getSeverity();
if (currentAlarm != null) {
// TODO: In some extremely rare cases, we might miss the event of alarm clear (If one use in-mem queue and restarted the server) or (if one manipulated the rule chain).
// Maybe we should fetch alarm every time?
@ -213,7 +217,7 @@ class AlarmState {
currentAlarm.setSeverity(severity);
currentAlarm.setStartTs(System.currentTimeMillis());
currentAlarm.setEndTs(currentAlarm.getStartTs());
currentAlarm.setDetails(JacksonUtil.OBJECT_MAPPER.createObjectNode());
currentAlarm.setDetails(createDetails(ruleState));
currentAlarm.setOriginator(originator);
currentAlarm.setTenantId(ctx.getTenantId());
currentAlarm.setPropagate(alarmDefinition.isPropagate());
@ -226,6 +230,44 @@ class AlarmState {
}
}
private <T> JsonNode createDetails(AlarmRuleState ruleState) {
ObjectNode details = JacksonUtil.OBJECT_MAPPER.createObjectNode();
String alarmDetails = ruleState.getAlarmRule().getAlarmDetails();
if (alarmDetails != null) {
for (KeyFilter keyFilter : ruleState.getAlarmRule().getCondition().getCondition()) {
EntityKeyValue entityKeyValue = dataSnapshot.getValue(keyFilter.getKey());
alarmDetails = alarmDetails.replaceAll(String.format("\\$\\{%s}", keyFilter.getKey().getKey()), getValueAsString(entityKeyValue));
}
details.put("data", alarmDetails);
}
return details;
}
private static String getValueAsString(EntityKeyValue entityKeyValue) {
Object result = null;
switch (entityKeyValue.getDataType()) {
case STRING:
result = entityKeyValue.getStrValue();
break;
case JSON:
result = entityKeyValue.getJsonValue();
break;
case LONG:
result = entityKeyValue.getLngValue();
break;
case DOUBLE:
result = entityKeyValue.getDblValue();
break;
case BOOLEAN:
result = entityKeyValue.getBoolValue();
break;
}
return String.valueOf(result);
}
public boolean processAlarmClear(TbContext ctx, Alarm alarmNf) {
boolean updated = false;
if (currentAlarm != null && currentAlarm.getId().equals(alarmNf.getId())) {

Loading…
Cancel
Save