1012 changed files with 43750 additions and 17892 deletions
File diff suppressed because one or more lines are too long
File diff suppressed because one or more lines are too long
@ -0,0 +1,41 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.actors.calculatedField; |
|||
|
|||
import lombok.Builder; |
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.alarm.Alarm; |
|||
import org.thingsboard.server.common.data.audit.ActionType; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.msg.MsgType; |
|||
import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; |
|||
import org.thingsboard.server.common.msg.queue.TbCallback; |
|||
|
|||
@Data |
|||
@Builder |
|||
public class CalculatedFieldAlarmActionMsg implements ToCalculatedFieldSystemMsg { |
|||
|
|||
private final TenantId tenantId; |
|||
private final Alarm alarm; |
|||
private final ActionType action; |
|||
private final TbCallback callback; |
|||
|
|||
@Override |
|||
public MsgType getMsgType() { |
|||
return MsgType.CF_ALARM_ACTION_MSG; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,57 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.actors.calculatedField; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import lombok.Builder; |
|||
import lombok.Data; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.audit.ActionType; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.msg.MsgType; |
|||
import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; |
|||
import org.thingsboard.server.common.msg.queue.TbCallback; |
|||
import org.thingsboard.server.common.util.ProtoUtils; |
|||
import org.thingsboard.server.gen.transport.TransportProtos.EntityActionEventProto; |
|||
|
|||
@Data |
|||
@Builder |
|||
public class CalculatedFieldEntityActionEventMsg implements ToCalculatedFieldSystemMsg { |
|||
|
|||
private final TenantId tenantId; |
|||
private final EntityId entityId; |
|||
private final JsonNode entity; |
|||
private final ActionType action; |
|||
private final TbCallback callback; |
|||
|
|||
public static CalculatedFieldEntityActionEventMsg fromProto(EntityActionEventProto proto, |
|||
TbCallback callback) { |
|||
return CalculatedFieldEntityActionEventMsg.builder() |
|||
.tenantId((TenantId) ProtoUtils.fromProto(proto.getTenantId())) |
|||
.entityId(ProtoUtils.fromProto(proto.getEntityId())) |
|||
.entity(JacksonUtil.toJsonNode(proto.getEntity())) |
|||
.action(ActionType.valueOf(proto.getAction())) |
|||
.callback(callback) |
|||
.build(); |
|||
} |
|||
|
|||
@Override |
|||
public MsgType getMsgType() { |
|||
return MsgType.CF_ENTITY_ACTION_EVENT_MSG; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,52 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.actors.calculatedField; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.audit.ActionType; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.msg.MsgType; |
|||
import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; |
|||
import org.thingsboard.server.common.msg.queue.TbCallback; |
|||
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
|||
|
|||
@Data |
|||
public class CalculatedFieldRelationActionMsg implements ToCalculatedFieldSystemMsg { |
|||
|
|||
private final TenantId tenantId; |
|||
private final EntityId relatedEntityId; |
|||
private final ActionType action; |
|||
private final CalculatedFieldCtx calculatedField; |
|||
private final TbCallback callback; |
|||
|
|||
public CalculatedFieldRelationActionMsg(TenantId tenantId, |
|||
EntityId relatedEntityId, ActionType action, |
|||
CalculatedFieldCtx calculatedField, |
|||
TbCallback callback) { |
|||
this.tenantId = tenantId; |
|||
this.relatedEntityId = relatedEntityId; |
|||
this.action = action; |
|||
this.calculatedField = calculatedField; |
|||
this.callback = callback; |
|||
} |
|||
|
|||
@Override |
|||
public MsgType getMsgType() { |
|||
return MsgType.CF_RELATION_ACTION_MSG; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,82 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf; |
|||
|
|||
import lombok.Builder; |
|||
import lombok.Data; |
|||
import lombok.RequiredArgsConstructor; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.rule.engine.action.TbAlarmResult; |
|||
import org.thingsboard.server.common.data.DataConstants; |
|||
import org.thingsboard.server.common.data.id.CalculatedFieldId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.msg.TbMsgType; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|||
|
|||
import java.util.List; |
|||
|
|||
@Data |
|||
@Builder |
|||
@RequiredArgsConstructor |
|||
public class AlarmCalculatedFieldResult implements CalculatedFieldResult { |
|||
|
|||
private final TbAlarmResult alarmResult; |
|||
|
|||
@Override |
|||
public TbMsg toTbMsg(EntityId entityId, List<CalculatedFieldId> cfIds) { |
|||
TbMsgType msgType; |
|||
TbMsgMetaData metaData = new TbMsgMetaData(); |
|||
if (alarmResult.isCreated()) { |
|||
msgType = TbMsgType.ALARM_CREATED; |
|||
metaData.putValue(DataConstants.IS_NEW_ALARM, Boolean.TRUE.toString()); |
|||
} else if (alarmResult.isUpdated()) { |
|||
msgType = TbMsgType.ALARM_UPDATED; |
|||
metaData.putValue(DataConstants.IS_EXISTING_ALARM, Boolean.TRUE.toString()); |
|||
} else if (alarmResult.isSeverityUpdated()) { |
|||
msgType = TbMsgType.ALARM_SEVERITY_UPDATED; |
|||
metaData.putValue(DataConstants.IS_EXISTING_ALARM, Boolean.TRUE.toString()); |
|||
metaData.putValue(DataConstants.IS_SEVERITY_UPDATED_ALARM, Boolean.TRUE.toString()); |
|||
} else { |
|||
msgType = TbMsgType.ALARM_CLEAR; |
|||
metaData.putValue(DataConstants.IS_CLEARED_ALARM, Boolean.TRUE.toString()); |
|||
} |
|||
if (alarmResult.getConditionRepeats() != null) { |
|||
metaData.putValue(DataConstants.ALARM_CONDITION_REPEATS, String.valueOf(alarmResult.getConditionRepeats())); |
|||
} |
|||
if (alarmResult.getConditionDuration() != null) { |
|||
metaData.putValue(DataConstants.ALARM_CONDITION_DURATION, String.valueOf(alarmResult.getConditionDuration())); |
|||
} |
|||
|
|||
return TbMsg.newMsg() |
|||
.type(msgType) |
|||
.originator(entityId) |
|||
.data(JacksonUtil.toString(alarmResult.getAlarm())) |
|||
.metaData(metaData) |
|||
.build(); |
|||
} |
|||
|
|||
@Override |
|||
public String stringValue() { |
|||
return alarmResult != null ? JacksonUtil.toString(alarmResult) : null; |
|||
} |
|||
|
|||
@Override |
|||
public boolean isEmpty() { |
|||
return alarmResult == null; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,76 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf; |
|||
|
|||
import lombok.RequiredArgsConstructor; |
|||
import org.springframework.stereotype.Service; |
|||
import org.thingsboard.server.common.data.Customer; |
|||
import org.thingsboard.server.common.data.DeviceInfo; |
|||
import org.thingsboard.server.common.data.DeviceInfoFilter; |
|||
import org.thingsboard.server.common.data.EntityType; |
|||
import org.thingsboard.server.common.data.asset.Asset; |
|||
import org.thingsboard.server.common.data.id.AssetId; |
|||
import org.thingsboard.server.common.data.id.CustomerId; |
|||
import org.thingsboard.server.common.data.id.DeviceId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.page.PageDataIterable; |
|||
import org.thingsboard.server.dao.asset.AssetService; |
|||
import org.thingsboard.server.dao.customer.CustomerService; |
|||
import org.thingsboard.server.dao.device.DeviceService; |
|||
|
|||
import java.util.HashSet; |
|||
import java.util.Set; |
|||
|
|||
@Service |
|||
@RequiredArgsConstructor |
|||
public class OwnerService { |
|||
|
|||
private final DeviceService deviceService; |
|||
private final AssetService assetService; |
|||
private final CustomerService customerService; |
|||
|
|||
public EntityId getOwner(TenantId tenantId, EntityId entityId) { |
|||
return switch (entityId.getEntityType()) { |
|||
case DEVICE -> deviceService.findDeviceById(tenantId, (DeviceId) entityId).getOwnerId(); |
|||
case ASSET -> assetService.findAssetById(tenantId, (AssetId) entityId).getOwnerId(); |
|||
case CUSTOMER -> tenantId; |
|||
default -> throw new UnsupportedOperationException(); |
|||
}; |
|||
} |
|||
|
|||
public Set<EntityId> getOwnedEntities(TenantId tenantId, EntityId ownerId) { |
|||
Set<EntityId> ownedEntities = new HashSet<>(); |
|||
if (EntityType.CUSTOMER.equals(ownerId.getEntityType())) { |
|||
PageDataIterable<DeviceInfo> deviceIdInfos = new PageDataIterable<>(pageLink -> deviceService.findDeviceInfosByFilter(DeviceInfoFilter.builder().tenantId(tenantId).customerId((CustomerId) ownerId).build(), pageLink), 1000); |
|||
deviceIdInfos.forEach(deviceInfo -> ownedEntities.add(deviceInfo.getId())); |
|||
|
|||
PageDataIterable<Asset> assets = new PageDataIterable<>(pageLink -> assetService.findAssetsByTenantIdAndCustomerId(tenantId, (CustomerId) ownerId, pageLink), 1000); |
|||
assets.forEach(asset -> ownedEntities.add(asset.getId())); |
|||
} else if (EntityType.TENANT.equals(ownerId.getEntityType())) { |
|||
PageDataIterable<DeviceInfo> deviceIdInfos = new PageDataIterable<>(pageLink -> deviceService.findDeviceInfosByFilter(DeviceInfoFilter.builder().tenantId((TenantId) ownerId).customerId(new CustomerId(CustomerId.NULL_UUID)).build(), pageLink), 1000); |
|||
deviceIdInfos.forEach(deviceInfo -> ownedEntities.add(deviceInfo.getId())); |
|||
|
|||
PageDataIterable<Asset> assets = new PageDataIterable<>(pageLink -> assetService.findAssetsByTenantIdAndCustomerId((TenantId) ownerId, new CustomerId(CustomerId.NULL_UUID), pageLink), 1000); |
|||
assets.forEach(asset -> ownedEntities.add(asset.getId())); |
|||
|
|||
PageDataIterable<Customer> customers = new PageDataIterable<>(pageLink -> customerService.findCustomersByTenantId((TenantId) ownerId, pageLink), 1000); |
|||
customers.forEach(customer -> ownedEntities.add(customer.getId())); |
|||
} |
|||
return ownedEntities; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,49 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf; |
|||
|
|||
import lombok.Builder; |
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.id.CalculatedFieldId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.util.CollectionsUtil; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
|
|||
import java.util.List; |
|||
|
|||
@Data |
|||
@Builder |
|||
public final class PropagationCalculatedFieldResult implements CalculatedFieldResult { |
|||
|
|||
private final List<EntityId> propagationEntityIds; |
|||
private final TelemetryCalculatedFieldResult result; |
|||
|
|||
@Override |
|||
public TbMsg toTbMsg(EntityId entityId, List<CalculatedFieldId> cfIds) { |
|||
return result.toTbMsg(entityId, cfIds); |
|||
} |
|||
|
|||
@Override |
|||
public String stringValue() { |
|||
return result.stringValue(); |
|||
} |
|||
|
|||
@Override |
|||
public boolean isEmpty() { |
|||
return CollectionsUtil.isEmpty(propagationEntityIds) || result.isEmpty(); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,76 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import lombok.Builder; |
|||
import lombok.Data; |
|||
import org.thingsboard.server.common.data.AttributeScope; |
|||
import org.thingsboard.server.common.data.cf.configuration.OutputType; |
|||
import org.thingsboard.server.common.data.id.CalculatedFieldId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.msg.TbMsgType; |
|||
import org.thingsboard.server.common.msg.TbMsg; |
|||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|||
|
|||
import java.util.List; |
|||
import java.util.Map; |
|||
|
|||
import static org.thingsboard.server.common.data.DataConstants.SCOPE; |
|||
|
|||
@Data |
|||
@Builder |
|||
public final class TelemetryCalculatedFieldResult implements CalculatedFieldResult { |
|||
|
|||
private final OutputType type; |
|||
private final AttributeScope scope; |
|||
private final JsonNode result; |
|||
|
|||
public static final TelemetryCalculatedFieldResult EMPTY = TelemetryCalculatedFieldResult.builder().result(null).build(); |
|||
|
|||
@Override |
|||
public TbMsg toTbMsg(EntityId entityId, List<CalculatedFieldId> cfIds) { |
|||
TbMsgType msgType = switch (type) { |
|||
case ATTRIBUTES -> TbMsgType.POST_ATTRIBUTES_REQUEST; |
|||
case TIME_SERIES -> TbMsgType.POST_TELEMETRY_REQUEST; |
|||
}; |
|||
TbMsgMetaData metaData = switch (type) { |
|||
case ATTRIBUTES -> new TbMsgMetaData(Map.of(SCOPE, scope.name())); |
|||
case TIME_SERIES -> TbMsgMetaData.EMPTY; |
|||
}; |
|||
return TbMsg.newMsg() |
|||
.type(msgType) |
|||
.originator(entityId) |
|||
.previousCalculatedFieldIds(cfIds) |
|||
.data(stringValue()) |
|||
.metaData(metaData) |
|||
.build(); |
|||
} |
|||
|
|||
@Override |
|||
public String stringValue() { |
|||
return result == null ? null : result.toString(); |
|||
} |
|||
|
|||
@Override |
|||
public boolean isEmpty() { |
|||
return result == null || result.isMissingNode() || result.isNull() || |
|||
(result.isObject() && result.isEmpty()) || |
|||
(result.isArray() && result.isEmpty()) || |
|||
(result.isTextual() && result.asText().isEmpty()); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,241 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.aggregation; |
|||
|
|||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import lombok.Getter; |
|||
import lombok.Setter; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.actors.TbActorRef; |
|||
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|||
import org.thingsboard.server.common.data.cf.configuration.Output; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunctionInput; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggInput; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInput; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.service.cf.CalculatedFieldResult; |
|||
import org.thingsboard.server.service.cf.TelemetryCalculatedFieldResult; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
|||
import org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState; |
|||
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
|||
import org.thingsboard.server.service.cf.ctx.state.aggregation.function.AggEntry; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.HashMap; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.Map.Entry; |
|||
import java.util.concurrent.ScheduledFuture; |
|||
|
|||
import static java.util.concurrent.TimeUnit.SECONDS; |
|||
|
|||
@Slf4j |
|||
@Getter |
|||
public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculatedFieldState { |
|||
|
|||
@Setter |
|||
private long lastArgsRefreshTs = -1; |
|||
@Setter |
|||
private long lastMetricsEvalTs = -1; |
|||
@Setter |
|||
private long lastRelatedEntitiesRefreshTs = -1; |
|||
private long deduplicationIntervalMs = -1; |
|||
private Map<String, AggMetric> metrics; |
|||
|
|||
private ScheduledFuture<?> reevaluationFuture; |
|||
|
|||
public RelatedEntitiesAggregationCalculatedFieldState(EntityId entityId) { |
|||
super(entityId); |
|||
} |
|||
|
|||
@Override |
|||
public void setCtx(CalculatedFieldCtx ctx, TbActorRef actorCtx) { |
|||
super.setCtx(ctx, actorCtx); |
|||
var configuration = (RelatedEntitiesAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); |
|||
metrics = configuration.getMetrics(); |
|||
deduplicationIntervalMs = SECONDS.toMillis(configuration.getDeduplicationIntervalInSec()); |
|||
} |
|||
|
|||
@Override |
|||
public void close() { |
|||
super.close(); |
|||
if (reevaluationFuture != null) { |
|||
reevaluationFuture.cancel(true); |
|||
reevaluationFuture = null; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void reset() { // must reset everything dependent on arguments
|
|||
super.reset(); |
|||
lastArgsRefreshTs = -1; |
|||
lastMetricsEvalTs = -1; |
|||
lastRelatedEntitiesRefreshTs = -1; |
|||
metrics = null; |
|||
} |
|||
|
|||
public void updateLastRelatedEntitiesRefreshTs() { |
|||
lastRelatedEntitiesRefreshTs = System.currentTimeMillis(); |
|||
} |
|||
|
|||
@Override |
|||
public CalculatedFieldType getType() { |
|||
return CalculatedFieldType.RELATED_ENTITIES_AGGREGATION; |
|||
} |
|||
|
|||
@Override |
|||
public Map<String, ArgumentEntry> update(Map<String, ArgumentEntry> argumentValues, CalculatedFieldCtx ctx) { |
|||
lastArgsRefreshTs = System.currentTimeMillis(); |
|||
return super.update(argumentValues, ctx); |
|||
} |
|||
|
|||
public List<EntityId> checkRelatedEntities(List<EntityId> relatedEntities) { |
|||
Map<EntityId, Map<String, ArgumentEntry>> entityInputs = prepareInputs(); |
|||
findOutdatedEntities(entityInputs, relatedEntities).forEach(this::cleanupEntityData); |
|||
updateLastRelatedEntitiesRefreshTs(); |
|||
return findMissingEntities(entityInputs, relatedEntities); |
|||
} |
|||
|
|||
private List<EntityId> findMissingEntities(Map<EntityId, Map<String, ArgumentEntry>> entityInputs, List<EntityId> relatedEntities) { |
|||
List<EntityId> missing = new ArrayList<>(); |
|||
relatedEntities.forEach(entityId -> { |
|||
if (!entityInputs.containsKey(entityId)) { |
|||
missing.add(entityId); |
|||
log.warn("[{}] Missing related entity inputs for {}", ctx.getCfId(), entityId); |
|||
} |
|||
}); |
|||
return missing; |
|||
} |
|||
|
|||
private List<EntityId> findOutdatedEntities(Map<EntityId, Map<String, ArgumentEntry>> entityInputs, List<EntityId> relatedEntities) { |
|||
List<EntityId> outdated = new ArrayList<>(); |
|||
entityInputs.keySet().forEach(entityId -> { |
|||
if (!relatedEntities.contains(entityId)) { |
|||
outdated.add(entityId); |
|||
log.warn("[{}] CF state keeps outdated related entity {}", ctx.getCfId(), entityId); |
|||
} |
|||
}); |
|||
return outdated; |
|||
} |
|||
|
|||
public Map<String, ArgumentEntry> updateEntityData(Map<String, ArgumentEntry> fetchedArgs) { |
|||
lastMetricsEvalTs = -1; |
|||
return update(fetchedArgs, ctx); |
|||
} |
|||
|
|||
public void cleanupEntityData(EntityId relatedEntityId) { |
|||
arguments.values().forEach(argEntry -> { |
|||
RelatedEntitiesArgumentEntry aggEntry = (RelatedEntitiesArgumentEntry) argEntry; |
|||
aggEntry.getEntityInputs().remove(relatedEntityId); |
|||
}); |
|||
lastMetricsEvalTs = -1; |
|||
lastArgsRefreshTs = System.currentTimeMillis(); |
|||
} |
|||
|
|||
public void scheduleReevaluation() { |
|||
ScheduledFuture<?> future = ctx.scheduleReevaluation(deduplicationIntervalMs, actorCtx); |
|||
if (future != null) { |
|||
reevaluationFuture = future; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) throws Exception { |
|||
boolean cfUpdated = updatedArgs != null && updatedArgs.isEmpty(); |
|||
if (shouldRecalculate() || cfUpdated) { |
|||
Output output = ctx.getOutput(); |
|||
ObjectNode aggResult = aggregateMetrics(output); |
|||
lastMetricsEvalTs = System.currentTimeMillis(); |
|||
scheduleReevaluation(); |
|||
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder() |
|||
.type(output.getType()) |
|||
.scope(output.getScope()) |
|||
.result(toSimpleResult(ctx.isUseLatestTs(), aggResult)) |
|||
.build()); |
|||
} else { |
|||
return Futures.immediateFuture(TelemetryCalculatedFieldResult.EMPTY); |
|||
} |
|||
} |
|||
|
|||
private boolean shouldRecalculate() { |
|||
boolean intervalPassed = lastMetricsEvalTs <= System.currentTimeMillis() - deduplicationIntervalMs; |
|||
boolean argsUpdatedDuringInterval = lastArgsRefreshTs > lastMetricsEvalTs; |
|||
return intervalPassed && argsUpdatedDuringInterval; |
|||
} |
|||
|
|||
private Map<EntityId, Map<String, ArgumentEntry>> prepareInputs() { |
|||
Map<EntityId, Map<String, ArgumentEntry>> inputs = new HashMap<>(); |
|||
for (Map.Entry<String, ArgumentEntry> argEntry : arguments.entrySet()) { |
|||
String key = argEntry.getKey(); |
|||
RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry = (RelatedEntitiesArgumentEntry) argEntry.getValue(); |
|||
relatedEntitiesArgumentEntry.getEntityInputs().forEach((entityId, argumentEntry) -> { |
|||
inputs.computeIfAbsent(entityId, k -> new HashMap<>()).put(key, argumentEntry); |
|||
}); |
|||
} |
|||
return inputs; |
|||
} |
|||
|
|||
private ObjectNode aggregateMetrics(Output output) throws Exception { |
|||
ObjectNode aggResult = JacksonUtil.newObjectNode(); |
|||
Map<EntityId, Map<String, ArgumentEntry>> inputs = prepareInputs(); |
|||
for (Entry<String, AggMetric> entry : metrics.entrySet()) { |
|||
String metricKey = entry.getKey(); |
|||
AggMetric metric = entry.getValue(); |
|||
|
|||
AggEntry aggMetricEntry = AggEntry.createAggFunction(metric.getFunction()); |
|||
aggregateMetric(metric, aggMetricEntry, inputs); |
|||
aggMetricEntry.result(output.getDecimalsByDefault()).ifPresent(result -> { |
|||
aggResult.set(metricKey, JacksonUtil.valueToTree(result)); |
|||
}); |
|||
} |
|||
return aggResult; |
|||
} |
|||
|
|||
private void aggregateMetric(AggMetric metric, AggEntry aggEntry, Map<EntityId, Map<String, ArgumentEntry>> inputs) throws Exception { |
|||
for (Map<String, ArgumentEntry> entityInputs : inputs.values()) { |
|||
if (applyAggregation(metric.getFilter(), entityInputs)) { |
|||
Object arg = resolveAggregationInput(metric.getInput(), entityInputs); |
|||
if (arg != null) { |
|||
aggEntry.update(arg); |
|||
} |
|||
} |
|||
} |
|||
} |
|||
|
|||
private boolean applyAggregation(String filter, Map<String, ArgumentEntry> entityInputs) throws Exception { |
|||
if (filter == null || filter.isEmpty()) { |
|||
return true; |
|||
} else { |
|||
Object filterResult = ctx.evaluateTbelExpression(filter, entityInputs, getLatestTimestamp()).get(); |
|||
return filterResult instanceof Boolean booleanResult && booleanResult; |
|||
} |
|||
} |
|||
|
|||
private Object resolveAggregationInput(AggInput aggInput, Map<String, ArgumentEntry> entityInputs) throws Exception { |
|||
if (aggInput instanceof AggFunctionInput functionInput) { |
|||
return ctx.evaluateTbelExpression(functionInput.getFunction(), entityInputs, getLatestTimestamp()).get(); |
|||
} else { |
|||
String inputKey = ((AggKeyInput) aggInput).getKey(); |
|||
return entityInputs.get(inputKey).getValue(); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,86 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.aggregation; |
|||
|
|||
import lombok.AllArgsConstructor; |
|||
import lombok.Data; |
|||
import org.thingsboard.script.api.tbel.TbelCfArg; |
|||
import org.thingsboard.script.api.tbel.TbelCfRelatedEntitiesArgumentValue; |
|||
import org.thingsboard.script.api.tbel.TbelCfSingleValueArg; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType; |
|||
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; |
|||
|
|||
import java.util.Map; |
|||
import java.util.stream.Collectors; |
|||
|
|||
@Data |
|||
@AllArgsConstructor |
|||
public class RelatedEntitiesArgumentEntry implements ArgumentEntry { |
|||
|
|||
private final Map<EntityId, ArgumentEntry> entityInputs; |
|||
|
|||
private boolean forceResetPrevious; |
|||
|
|||
@Override |
|||
public ArgumentEntryType getType() { |
|||
return ArgumentEntryType.RELATED_ENTITIES; |
|||
} |
|||
|
|||
@Override |
|||
public Object getValue() { |
|||
return entityInputs; |
|||
} |
|||
|
|||
@Override |
|||
public boolean updateEntry(ArgumentEntry entry) { |
|||
if (entry instanceof RelatedEntitiesArgumentEntry relatedEntitiesArgumentEntry) { |
|||
entityInputs.putAll(relatedEntitiesArgumentEntry.entityInputs); |
|||
return true; |
|||
} else if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { |
|||
if (entry.isForceResetPrevious()) { |
|||
entityInputs.put(singleValueArgumentEntry.getEntityId(), singleValueArgumentEntry); |
|||
return true; |
|||
} |
|||
ArgumentEntry argumentEntry = entityInputs.get(singleValueArgumentEntry.getEntityId()); |
|||
if (argumentEntry != null) { |
|||
argumentEntry.updateEntry(singleValueArgumentEntry); |
|||
} else { |
|||
entityInputs.put(singleValueArgumentEntry.getEntityId(), singleValueArgumentEntry); |
|||
} |
|||
return true; |
|||
} else { |
|||
throw new IllegalArgumentException("Unsupported argument entry type for aggregation argument entry: " + entry.getType()); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public boolean isEmpty() { |
|||
return entityInputs.isEmpty(); |
|||
} |
|||
|
|||
@Override |
|||
public TbelCfArg toTbelCfArg() { |
|||
var inputs = entityInputs.entrySet().stream() |
|||
.collect(Collectors.toMap( |
|||
e -> e.getKey().getId(), |
|||
e -> (TbelCfSingleValueArg) e.getValue().toTbelCfArg() |
|||
)); |
|||
return new TbelCfRelatedEntitiesArgumentValue(inputs); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,58 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.aggregation.function; |
|||
|
|||
import com.fasterxml.jackson.annotation.JsonIgnore; |
|||
import com.fasterxml.jackson.annotation.JsonSubTypes; |
|||
import com.fasterxml.jackson.annotation.JsonTypeInfo; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; |
|||
|
|||
import java.util.Optional; |
|||
|
|||
@JsonTypeInfo( |
|||
use = JsonTypeInfo.Id.NAME, |
|||
include = JsonTypeInfo.As.PROPERTY, |
|||
property = "type" |
|||
) |
|||
@JsonSubTypes({ |
|||
@JsonSubTypes.Type(value = AvgAggEntry.class, name = "AVG"), |
|||
@JsonSubTypes.Type(value = CountAggEntry.class, name = "COUNT"), |
|||
@JsonSubTypes.Type(value = CountUniqueAggEntry.class, name = "COUNT_UNIQUE"), |
|||
@JsonSubTypes.Type(value = MaxAggEntry.class, name = "MAX"), |
|||
@JsonSubTypes.Type(value = MinAggEntry.class, name = "MIN"), |
|||
@JsonSubTypes.Type(value = SumAggEntry.class, name = "SUM") |
|||
}) |
|||
public interface AggEntry { |
|||
|
|||
@JsonIgnore |
|||
AggFunction getType(); |
|||
|
|||
void update(Object value); |
|||
|
|||
Optional<Object> result(Integer precision); |
|||
|
|||
static AggEntry createAggFunction(AggFunction function) { |
|||
return switch (function) { |
|||
case MIN -> new MinAggEntry(); |
|||
case MAX -> new MaxAggEntry(); |
|||
case SUM -> new SumAggEntry(); |
|||
case AVG -> new AvgAggEntry(); |
|||
case COUNT -> new CountAggEntry(); |
|||
case COUNT_UNIQUE -> new CountUniqueAggEntry(); |
|||
}; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,47 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.aggregation.function; |
|||
|
|||
import org.thingsboard.script.api.tbel.TbUtils; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; |
|||
|
|||
import java.math.BigDecimal; |
|||
import java.math.RoundingMode; |
|||
|
|||
public class AvgAggEntry extends BaseAggEntry { |
|||
|
|||
private BigDecimal sum = BigDecimal.ZERO; |
|||
private long count = 0L; |
|||
|
|||
@Override |
|||
protected void doUpdate(double value) { |
|||
if (value != 0.0) { |
|||
sum = sum.add(BigDecimal.valueOf(value)); |
|||
} |
|||
this.count++; |
|||
} |
|||
|
|||
@Override |
|||
protected Object prepareResult(Integer precision) { |
|||
double result = sum.divide(BigDecimal.valueOf(count), RoundingMode.HALF_UP).doubleValue(); |
|||
return TbUtils.roundResult(result, precision); |
|||
} |
|||
|
|||
@Override |
|||
public AggFunction getType() { |
|||
return AggFunction.AVG; |
|||
} |
|||
} |
|||
@ -0,0 +1,55 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.aggregation.function; |
|||
|
|||
import java.util.Optional; |
|||
|
|||
public abstract class BaseAggEntry implements AggEntry { |
|||
|
|||
private boolean hasResult = false; |
|||
|
|||
@Override |
|||
public void update(Object value) { |
|||
doUpdate(extractDoubleValue(value)); |
|||
hasResult = true; |
|||
} |
|||
|
|||
@Override |
|||
public Optional<Object> result(Integer precision) { |
|||
if (hasResult) { |
|||
hasResult = false; |
|||
return Optional.of(prepareResult(precision)); |
|||
} else { |
|||
return Optional.empty(); |
|||
} |
|||
} |
|||
|
|||
protected abstract void doUpdate(double value); |
|||
|
|||
protected abstract Object prepareResult(Integer precision); |
|||
|
|||
protected double extractDoubleValue(Object value) { |
|||
try { |
|||
if (value instanceof Number number) { |
|||
return number.doubleValue(); |
|||
} |
|||
return Double.parseDouble(value.toString()); |
|||
} catch (Exception e) { |
|||
throw new NumberFormatException("Cannot parse value " + value.toString()); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,41 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.aggregation.function; |
|||
|
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; |
|||
|
|||
import java.util.Optional; |
|||
|
|||
public class CountAggEntry implements AggEntry { |
|||
|
|||
private long count = 0L; |
|||
|
|||
@Override |
|||
public void update(Object value) { |
|||
count++; |
|||
} |
|||
|
|||
@Override |
|||
public Optional<Object> result(Integer precision) { |
|||
return Optional.of(count); |
|||
} |
|||
|
|||
@Override |
|||
public AggFunction getType() { |
|||
return AggFunction.COUNT; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,43 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.aggregation.function; |
|||
|
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; |
|||
|
|||
import java.util.Optional; |
|||
import java.util.Set; |
|||
|
|||
public class CountUniqueAggEntry implements AggEntry { |
|||
|
|||
private Set<String> items; |
|||
|
|||
@Override |
|||
public void update(Object value) { |
|||
if (value != null) { |
|||
items.add(String.valueOf(value)); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public Optional<Object> result(Integer precision) { |
|||
return Optional.of(items.size()); |
|||
} |
|||
|
|||
@Override |
|||
public AggFunction getType() { |
|||
return AggFunction.COUNT_UNIQUE; |
|||
} |
|||
} |
|||
@ -0,0 +1,41 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.aggregation.function; |
|||
|
|||
import org.thingsboard.script.api.tbel.TbUtils; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; |
|||
|
|||
public class MaxAggEntry extends BaseAggEntry { |
|||
|
|||
private double max = Double.MIN_VALUE; |
|||
|
|||
@Override |
|||
protected void doUpdate(double value) { |
|||
if (value > max) { |
|||
max = value; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
protected Object prepareResult(Integer precision) { |
|||
return TbUtils.roundResult(max, precision); |
|||
} |
|||
|
|||
@Override |
|||
public AggFunction getType() { |
|||
return AggFunction.MAX; |
|||
} |
|||
} |
|||
@ -0,0 +1,41 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.aggregation.function; |
|||
|
|||
import org.thingsboard.script.api.tbel.TbUtils; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; |
|||
|
|||
public class MinAggEntry extends BaseAggEntry { |
|||
|
|||
private double min = Double.MAX_VALUE; |
|||
|
|||
@Override |
|||
protected void doUpdate(double value) { |
|||
if (value < min) { |
|||
min = value; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
protected Object prepareResult(Integer precision) { |
|||
return TbUtils.roundResult(min, precision); |
|||
} |
|||
|
|||
@Override |
|||
public AggFunction getType() { |
|||
return AggFunction.MIN; |
|||
} |
|||
} |
|||
@ -0,0 +1,43 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.aggregation.function; |
|||
|
|||
import org.thingsboard.script.api.tbel.TbUtils; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; |
|||
|
|||
import java.math.BigDecimal; |
|||
|
|||
public class SumAggEntry extends BaseAggEntry { |
|||
|
|||
private BigDecimal sum = BigDecimal.ZERO; |
|||
|
|||
@Override |
|||
protected void doUpdate(double value) { |
|||
if (value != 0.0) { |
|||
sum = sum.add(BigDecimal.valueOf(value)); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
protected Object prepareResult(Integer precision) { |
|||
return TbUtils.roundResult(sum.doubleValue(), precision); |
|||
} |
|||
|
|||
@Override |
|||
public AggFunction getType() { |
|||
return AggFunction.SUM; |
|||
} |
|||
} |
|||
@ -0,0 +1,552 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.alarm; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import lombok.EqualsAndHashCode; |
|||
import lombok.Getter; |
|||
import lombok.Setter; |
|||
import lombok.SneakyThrows; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.common.util.KvUtil; |
|||
import org.thingsboard.rule.engine.action.TbAlarmResult; |
|||
import org.thingsboard.server.actors.TbActorRef; |
|||
import org.thingsboard.server.common.data.StringUtils; |
|||
import org.thingsboard.server.common.data.alarm.Alarm; |
|||
import org.thingsboard.server.common.data.alarm.AlarmApiCallResult; |
|||
import org.thingsboard.server.common.data.alarm.AlarmCreateOrUpdateActiveRequest; |
|||
import org.thingsboard.server.common.data.alarm.AlarmSeverity; |
|||
import org.thingsboard.server.common.data.alarm.AlarmUpdateRequest; |
|||
import org.thingsboard.server.common.data.alarm.rule.AlarmRule; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.AlarmCondition; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.AlarmConditionType; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.AlarmConditionValue; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.expression.AlarmConditionExpression; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.expression.AlarmConditionFilter; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.expression.ComplexOperation; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.expression.SimpleAlarmConditionExpression; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.expression.TbelAlarmConditionExpression; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.expression.predicate.BooleanFilterPredicate; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.expression.predicate.ComplexFilterPredicate; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.expression.predicate.KeyFilterPredicate; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.expression.predicate.NumericFilterPredicate; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.expression.predicate.StringFilterPredicate; |
|||
import org.thingsboard.server.common.data.audit.ActionType; |
|||
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|||
import org.thingsboard.server.common.data.cf.configuration.AlarmCalculatedFieldConfiguration; |
|||
import org.thingsboard.server.common.data.id.DashboardId; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.kv.KvEntry; |
|||
import org.thingsboard.server.service.cf.AlarmCalculatedFieldResult; |
|||
import org.thingsboard.server.service.cf.CalculatedFieldResult; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
|||
import org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState; |
|||
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
|||
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; |
|||
|
|||
import java.util.Comparator; |
|||
import java.util.Map; |
|||
import java.util.TreeMap; |
|||
import java.util.concurrent.ScheduledFuture; |
|||
import java.util.concurrent.atomic.AtomicBoolean; |
|||
import java.util.function.Function; |
|||
|
|||
import static org.thingsboard.server.common.data.StringUtils.equalsAny; |
|||
import static org.thingsboard.server.common.data.StringUtils.splitByCommaWithoutQuotes; |
|||
import static org.thingsboard.server.service.cf.ctx.state.alarm.AlarmEvalResult.Status.FALSE; |
|||
import static org.thingsboard.server.service.cf.ctx.state.alarm.AlarmEvalResult.Status.NOT_YET_TRUE; |
|||
import static org.thingsboard.server.service.cf.ctx.state.alarm.AlarmEvalResult.Status.TRUE; |
|||
|
|||
@EqualsAndHashCode(callSuper = true) |
|||
@Slf4j |
|||
public class AlarmCalculatedFieldState extends BaseCalculatedFieldState { |
|||
|
|||
private AlarmCalculatedFieldConfiguration configuration; |
|||
private String alarmType; |
|||
|
|||
@Getter |
|||
private final Map<AlarmSeverity, AlarmRuleState> createRuleStates = new TreeMap<>(Comparator.comparing(Enum::ordinal)); |
|||
@Getter |
|||
@Setter |
|||
private AlarmRuleState clearRuleState; |
|||
|
|||
@Getter |
|||
private Alarm currentAlarm; |
|||
private boolean initialFetchDone; |
|||
|
|||
// TODO: deprecate device profile node, describe the differences and improvements
|
|||
|
|||
public AlarmCalculatedFieldState(EntityId entityId) { |
|||
super(entityId); |
|||
} |
|||
|
|||
@Override |
|||
public void setCtx(CalculatedFieldCtx ctx, TbActorRef actorCtx) { |
|||
super.setCtx(ctx, actorCtx); |
|||
this.configuration = getConfiguration(ctx); |
|||
this.alarmType = ctx.getCalculatedField().getName(); |
|||
|
|||
Map<AlarmSeverity, AlarmRule> createRules = configuration.getCreateRules(); |
|||
createRules.forEach((severity, rule) -> { |
|||
AlarmRuleState ruleState = createRuleStates.get(severity); |
|||
if (ruleState != null) { |
|||
ruleState.setAlarmRule(rule); |
|||
} |
|||
}); |
|||
AlarmRule clearRule = configuration.getClearRule(); |
|||
if (clearRule != null && clearRuleState != null) { |
|||
clearRuleState.setAlarmRule(clearRule); |
|||
} |
|||
|
|||
if (currentAlarm != null && !currentAlarm.getType().equals(alarmType)) { |
|||
currentAlarm = null; |
|||
initialFetchDone = false; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public void init() { |
|||
super.init(); |
|||
AtomicBoolean reevalNeeded = new AtomicBoolean(false); |
|||
Map<AlarmSeverity, AlarmRule> createRules = configuration.getCreateRules(); |
|||
for (AlarmSeverity severity : AlarmSeverity.values()) { |
|||
AlarmRule rule = createRules.get(severity); |
|||
if (rule != null) { |
|||
createRuleStates.compute(severity, (__, ruleState) -> { |
|||
return initRuleState(severity, rule, ruleState, reevalNeeded); |
|||
}); |
|||
} else { |
|||
AlarmRuleState state = createRuleStates.remove(severity); |
|||
if (state != null) { |
|||
clearState(state); |
|||
} |
|||
} |
|||
} |
|||
|
|||
AlarmRule clearRule = configuration.getClearRule(); |
|||
if (clearRule != null) { |
|||
clearRuleState = initRuleState(null, clearRule, clearRuleState, reevalNeeded); |
|||
} else { |
|||
if (clearRuleState != null) { |
|||
clearState(clearRuleState); |
|||
clearRuleState = null; |
|||
} |
|||
} |
|||
log.debug("Initialized create rule states {} and clear rule state {} for {}", createRuleStates, clearRuleState, configuration); |
|||
|
|||
if (reevalNeeded.get()) { |
|||
initCurrentAlarm(ctx); |
|||
createOrClearAlarms(state -> { |
|||
if (state.getCondition().getType() == AlarmConditionType.DURATION) { |
|||
AlarmEvalResult evalResult = state.reeval(System.currentTimeMillis(), ctx); |
|||
if (evalResult.getStatus() == TRUE || evalResult.getStatus() == NOT_YET_TRUE) { |
|||
ScheduledFuture<?> future = ctx.scheduleReevaluation(evalResult.getLeftDuration(), actorCtx); |
|||
if (future != null) { |
|||
state.setDurationCheckFuture(future); |
|||
} |
|||
} |
|||
} |
|||
return AlarmEvalResult.NOT_YET_TRUE; |
|||
}, ctx); |
|||
} |
|||
} |
|||
|
|||
private AlarmRuleState initRuleState(AlarmSeverity severity, AlarmRule rule, AlarmRuleState ruleState, AtomicBoolean reevalNeeded) { |
|||
if (ruleState == null) { |
|||
ruleState = new AlarmRuleState(severity, rule, this); |
|||
} else { |
|||
// when restored
|
|||
ruleState.setAlarmRule(rule); |
|||
ruleState.setActive(null); |
|||
AlarmCondition condition = rule.getCondition(); |
|||
if (condition.hasSchedule() || (condition.getType() == AlarmConditionType.DURATION && !ruleState.isEmpty())) { |
|||
reevalNeeded.set(true); |
|||
} |
|||
} |
|||
return ruleState; |
|||
} |
|||
|
|||
@Override |
|||
public void reset() { |
|||
super.reset(); |
|||
configuration = null; |
|||
} |
|||
|
|||
@Override |
|||
public void close() { |
|||
super.close(); |
|||
for (AlarmRuleState state : createRuleStates.values()) { |
|||
clearState(state); |
|||
} |
|||
clearState(clearRuleState); |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) { |
|||
initCurrentAlarm(ctx); |
|||
TbAlarmResult result = createOrClearAlarms(state -> { |
|||
if (updatedArgs != null) { |
|||
boolean newEvent = !updatedArgs.isEmpty(); |
|||
AlarmEvalResult evalResult = state.eval(newEvent, ctx); |
|||
if (evalResult.getStatus() == NOT_YET_TRUE && evalResult.getLeftDuration() > 0) { |
|||
long leftDuration = evalResult.getLeftDuration(); |
|||
ScheduledFuture<?> future = ctx.scheduleReevaluation(leftDuration, actorCtx); |
|||
if (future != null) { |
|||
state.setDurationCheckFuture(future); |
|||
} |
|||
} |
|||
return evalResult; |
|||
} else { |
|||
return state.reeval(System.currentTimeMillis(), ctx); |
|||
} |
|||
}, ctx); |
|||
return Futures.immediateFuture(AlarmCalculatedFieldResult.builder() |
|||
.alarmResult(result) |
|||
.build()); |
|||
} |
|||
|
|||
public void processAlarmAction(Alarm alarm, ActionType action) { |
|||
switch (action) { |
|||
case ALARM_ACK -> processAlarmAck(alarm); |
|||
case ALARM_CLEAR -> processAlarmClear(alarm); |
|||
case ALARM_DELETE -> processAlarmDelete(alarm); |
|||
} |
|||
} |
|||
|
|||
private void processAlarmClear(Alarm alarm) { |
|||
currentAlarm = null; |
|||
createRuleStates.values().forEach(this::clearState); |
|||
clearState(clearRuleState); |
|||
} |
|||
|
|||
private void processAlarmAck(Alarm alarm) { |
|||
currentAlarm.setAcknowledged(alarm.isAcknowledged()); |
|||
currentAlarm.setAckTs(alarm.getAckTs()); |
|||
} |
|||
|
|||
private void processAlarmDelete(Alarm alarm) { |
|||
processAlarmClear(alarm); |
|||
} |
|||
|
|||
private TbAlarmResult createOrClearAlarms(Function<AlarmRuleState, AlarmEvalResult> evalFunction, |
|||
CalculatedFieldCtx ctx) { |
|||
TbAlarmResult result = null; |
|||
AlarmRuleState resultState = null; |
|||
AlarmRuleState.StateInfo resultStateInfo = null; |
|||
|
|||
for (AlarmRuleState state : createRuleStates.values()) { |
|||
AlarmEvalResult evalResult = evalFunction.apply(state); |
|||
log.debug("Evaluated create rule {} with args {}. Result: {}", state, arguments, evalResult); |
|||
if (evalResult.getStatus() == TRUE) { |
|||
resultState = state; |
|||
break; |
|||
} else if (evalResult.getStatus() == FALSE) { |
|||
clearState(state); |
|||
} |
|||
} |
|||
|
|||
if (resultState != null) { |
|||
result = calculateAlarmResult(resultState, ctx); |
|||
resultStateInfo = resultState.getStateInfo(); |
|||
log.debug("Alarm result for state {}: {}", resultState, result); |
|||
clearState(clearRuleState); |
|||
} else if (currentAlarm != null && clearRuleState != null) { |
|||
AlarmEvalResult evalResult = evalFunction.apply(clearRuleState); |
|||
log.debug("Evaluated clear rule {} with args {}. Result: {}", clearRuleState, arguments, evalResult); |
|||
if (evalResult.getStatus() == TRUE) { |
|||
resultStateInfo = clearRuleState.getStateInfo(); |
|||
clearState(clearRuleState); |
|||
for (AlarmRuleState state : createRuleStates.values()) { |
|||
clearState(state); |
|||
} |
|||
AlarmApiCallResult clearResult = ctx.getAlarmService().clearAlarm( |
|||
ctx.getTenantId(), currentAlarm.getId(), System.currentTimeMillis(), createDetails(clearRuleState), false |
|||
); |
|||
if (clearResult.isCleared()) { |
|||
result = TbAlarmResult.builder() |
|||
.isCleared(true) |
|||
.alarm(clearResult.getAlarm()) |
|||
.build(); |
|||
resultState = clearRuleState; |
|||
} |
|||
currentAlarm = null; |
|||
} else if (evalResult.getStatus() == FALSE) { |
|||
clearState(clearRuleState); |
|||
} |
|||
} |
|||
if (result != null && resultState != null) { |
|||
result.setConditionRepeats(resultStateInfo.eventCount()); |
|||
result.setConditionDuration(resultStateInfo.duration()); |
|||
} |
|||
return result; |
|||
} |
|||
|
|||
private void clearState(AlarmRuleState state) { |
|||
if (state != null) { |
|||
log.debug("Clearing rule state {}", state); |
|||
state.clear(); |
|||
} |
|||
} |
|||
|
|||
private void initCurrentAlarm(CalculatedFieldCtx ctx) { |
|||
if (!initialFetchDone) { |
|||
Alarm alarm = ctx.getAlarmService().findLatestActiveByOriginatorAndType(ctx.getTenantId(), entityId, alarmType); |
|||
if (alarm != null && !alarm.getStatus().isCleared()) { |
|||
currentAlarm = alarm; |
|||
} |
|||
initialFetchDone = true; |
|||
} |
|||
} |
|||
|
|||
private TbAlarmResult calculateAlarmResult(AlarmRuleState ruleState, CalculatedFieldCtx ctx) { |
|||
AlarmSeverity severity = ruleState.getSeverity(); |
|||
if (currentAlarm != null) { |
|||
currentAlarm.setEndTs(System.currentTimeMillis()); |
|||
AlarmSeverity oldSeverity = currentAlarm.getSeverity(); |
|||
// Skip update if severity is decreased.
|
|||
if (severity.ordinal() <= oldSeverity.ordinal()) { |
|||
currentAlarm.setDetails(createDetails(ruleState)); |
|||
currentAlarm.setSeverity(severity); |
|||
AlarmApiCallResult result = ctx.getAlarmService().updateAlarm(AlarmUpdateRequest.fromAlarm(currentAlarm)); |
|||
currentAlarm = result.getAlarm(); |
|||
return TbAlarmResult.fromAlarmResult(result); |
|||
} else { |
|||
return null; |
|||
} |
|||
} else { |
|||
var newAlarm = new Alarm(); |
|||
newAlarm.setType(alarmType); |
|||
newAlarm.setAcknowledged(false); |
|||
newAlarm.setCleared(false); |
|||
newAlarm.setSeverity(severity); |
|||
long startTs = latestTimestamp; |
|||
long currentTime = System.currentTimeMillis(); |
|||
if (startTs == 0L || startTs > currentTime) { |
|||
startTs = currentTime; |
|||
} |
|||
newAlarm.setStartTs(startTs); |
|||
newAlarm.setEndTs(startTs); |
|||
newAlarm.setDetails(createDetails(ruleState)); |
|||
newAlarm.setOriginator(entityId); |
|||
newAlarm.setTenantId(ctx.getTenantId()); |
|||
newAlarm.setPropagate(configuration.isPropagate()); |
|||
newAlarm.setPropagateToOwner(configuration.isPropagateToOwner()); |
|||
newAlarm.setPropagateToTenant(configuration.isPropagateToTenant()); |
|||
if (configuration.getPropagateRelationTypes() != null) { |
|||
newAlarm.setPropagateRelationTypes(configuration.getPropagateRelationTypes()); |
|||
} |
|||
AlarmApiCallResult result = ctx.getAlarmService().createAlarm(AlarmCreateOrUpdateActiveRequest.fromAlarm(newAlarm)); |
|||
currentAlarm = result.getAlarm(); |
|||
return TbAlarmResult.fromAlarmResult(result); |
|||
} |
|||
} |
|||
|
|||
private JsonNode createDetails(AlarmRuleState ruleState) { |
|||
JsonNode alarmDetails; |
|||
String alarmDetailsStr = ruleState.getAlarmRule().getAlarmDetails(); |
|||
DashboardId dashboardId = ruleState.getAlarmRule().getDashboardId(); |
|||
|
|||
if (StringUtils.isNotEmpty(alarmDetailsStr) || dashboardId != null) { |
|||
ObjectNode newDetails = JacksonUtil.newObjectNode(); |
|||
if (StringUtils.isNotEmpty(alarmDetailsStr)) { |
|||
for (Map.Entry<String, ArgumentEntry> entry : arguments.entrySet()) { |
|||
String key = entry.getKey(); |
|||
ArgumentEntry value = entry.getValue(); |
|||
alarmDetailsStr = alarmDetailsStr.replaceAll(String.format("\\$\\{%s}", key), String.valueOf(value.getValue())); |
|||
} |
|||
newDetails.put("data", alarmDetailsStr); |
|||
} |
|||
if (dashboardId != null) { |
|||
newDetails.put("dashboardId", dashboardId.getId().toString()); |
|||
} |
|||
alarmDetails = newDetails; |
|||
} else if (currentAlarm != null) { |
|||
alarmDetails = currentAlarm.getDetails(); |
|||
} else { |
|||
alarmDetails = JacksonUtil.newObjectNode(); |
|||
} |
|||
|
|||
return alarmDetails; |
|||
} |
|||
|
|||
@SneakyThrows |
|||
public boolean eval(AlarmConditionExpression expression, CalculatedFieldCtx ctx) { |
|||
if (expression instanceof TbelAlarmConditionExpression tbelExpression) { |
|||
Object result = ctx.evaluateTbelExpression(tbelExpression.getExpression(), this).get(); |
|||
if (result instanceof Boolean booleanResult) { |
|||
return booleanResult; |
|||
} else { |
|||
throw new IllegalStateException("Condition expression returned non-boolean value: '" + result + "'"); |
|||
} |
|||
} else { |
|||
SimpleAlarmConditionExpression simpleExpression = (SimpleAlarmConditionExpression) expression; |
|||
ComplexOperation operation = simpleExpression.getOperation(); |
|||
if (operation == null) { |
|||
operation = ComplexOperation.AND; |
|||
} |
|||
return switch (operation) { |
|||
case AND -> simpleExpression.getFilters().stream() |
|||
.allMatch(filter -> eval(getArgument(filter.getArgument()), filter)); |
|||
case OR -> simpleExpression.getFilters().stream() |
|||
.anyMatch(filter -> eval(getArgument(filter.getArgument()), filter)); |
|||
}; |
|||
} |
|||
} |
|||
|
|||
private boolean eval(SingleValueArgumentEntry argument, AlarmConditionFilter filter) { |
|||
ComplexOperation operation = filter.getOperation(); |
|||
if (operation == null) { |
|||
operation = ComplexOperation.AND; |
|||
} |
|||
return switch (operation) { |
|||
case AND -> filter.getPredicates().stream() |
|||
.allMatch(predicate -> eval(argument, predicate)); |
|||
case OR -> filter.getPredicates().stream() |
|||
.anyMatch(predicate -> eval(argument, predicate)); |
|||
}; |
|||
} |
|||
|
|||
private boolean eval(SingleValueArgumentEntry argument, KeyFilterPredicate predicate) { |
|||
return switch (predicate.getType()) { |
|||
case STRING -> evalStrPredicate(argument, (StringFilterPredicate) predicate); |
|||
case NUMERIC -> evalNumPredicate(argument, (NumericFilterPredicate) predicate); |
|||
case BOOLEAN -> evalBooleanPredicate(argument, (BooleanFilterPredicate) predicate); |
|||
case COMPLEX -> evalComplexPredicate(argument, (ComplexFilterPredicate) predicate); |
|||
}; |
|||
} |
|||
|
|||
private boolean evalComplexPredicate(SingleValueArgumentEntry argument, ComplexFilterPredicate complexPredicate) { |
|||
return switch (complexPredicate.getOperation()) { |
|||
case OR -> { |
|||
for (KeyFilterPredicate predicate : complexPredicate.getPredicates()) { |
|||
if (eval(argument, predicate)) { |
|||
yield true; |
|||
} |
|||
} |
|||
yield false; |
|||
} |
|||
case AND -> { |
|||
for (KeyFilterPredicate predicate : complexPredicate.getPredicates()) { |
|||
if (!eval(argument, predicate)) { |
|||
yield false; |
|||
} |
|||
} |
|||
yield true; |
|||
} |
|||
}; |
|||
} |
|||
|
|||
private boolean evalBooleanPredicate(SingleValueArgumentEntry argument, BooleanFilterPredicate predicate) { |
|||
Boolean value = KvUtil.getBoolValue(argument.getKvEntryValue()); |
|||
if (value == null) { |
|||
return false; |
|||
} |
|||
Boolean predicateValue = resolveValue(predicate.getValue(), KvUtil::getBoolValue); |
|||
if (predicateValue == null) { |
|||
return false; |
|||
} |
|||
return switch (predicate.getOperation()) { |
|||
case EQUAL -> value.equals(predicateValue); |
|||
case NOT_EQUAL -> !value.equals(predicateValue); |
|||
}; |
|||
} |
|||
|
|||
private boolean evalNumPredicate(SingleValueArgumentEntry argument, NumericFilterPredicate predicate) { |
|||
Double value = KvUtil.getDoubleValue(argument.getKvEntryValue()); |
|||
if (value == null) { |
|||
return false; |
|||
} |
|||
Double predicateValue = resolveValue(predicate.getValue(), KvUtil::getDoubleValue); |
|||
if (predicateValue == null) { |
|||
return false; |
|||
} |
|||
return switch (predicate.getOperation()) { |
|||
case NOT_EQUAL -> !value.equals(predicateValue); |
|||
case EQUAL -> value.equals(predicateValue); |
|||
case GREATER -> value > predicateValue; |
|||
case GREATER_OR_EQUAL -> value >= predicateValue; |
|||
case LESS -> value < predicateValue; |
|||
case LESS_OR_EQUAL -> value <= predicateValue; |
|||
}; |
|||
} |
|||
|
|||
private boolean evalStrPredicate(SingleValueArgumentEntry argument, StringFilterPredicate predicate) { |
|||
String value = KvUtil.getStringValue(argument.getKvEntryValue()); |
|||
if (value == null) { |
|||
return false; |
|||
} |
|||
String predicateValue = resolveValue(predicate.getValue(), KvUtil::getStringValue); |
|||
if (predicateValue == null) { |
|||
return false; |
|||
} |
|||
if (predicate.isIgnoreCase()) { |
|||
value = value.toLowerCase(); |
|||
predicateValue = predicateValue.toLowerCase(); |
|||
} |
|||
return switch (predicate.getOperation()) { |
|||
case CONTAINS -> value.contains(predicateValue); |
|||
case EQUAL -> value.equals(predicateValue); |
|||
case STARTS_WITH -> value.startsWith(predicateValue); |
|||
case ENDS_WITH -> value.endsWith(predicateValue); |
|||
case NOT_EQUAL -> !value.equals(predicateValue); |
|||
case NOT_CONTAINS -> !value.contains(predicateValue); |
|||
case IN -> equalsAny(value, splitByCommaWithoutQuotes(predicateValue)); |
|||
case NOT_IN -> !equalsAny(value, splitByCommaWithoutQuotes(predicateValue)); |
|||
}; |
|||
} |
|||
|
|||
protected <T> T resolveValue(AlarmConditionValue<T> conditionValue, Function<KvEntry, T> mapper) { |
|||
T value = conditionValue.getStaticValue(); |
|||
if (value == null) { |
|||
String argument = conditionValue.getDynamicValueArgument(); |
|||
SingleValueArgumentEntry entry = getArgument(argument); |
|||
value = mapper.apply(entry.getKvEntryValue()); |
|||
if (value == null) { |
|||
throw new IllegalArgumentException("No proper value found for argument " + argument); |
|||
} |
|||
} |
|||
return value; |
|||
} |
|||
|
|||
protected SingleValueArgumentEntry getArgument(String key) { |
|||
SingleValueArgumentEntry entry = (SingleValueArgumentEntry) arguments.get(key); |
|||
if (entry == null) { |
|||
throw new IllegalArgumentException("Argument '" + key + "' is missing"); |
|||
} |
|||
return entry; |
|||
} |
|||
|
|||
private AlarmCalculatedFieldConfiguration getConfiguration(CalculatedFieldCtx ctx) { |
|||
return (AlarmCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); |
|||
} |
|||
|
|||
@Override |
|||
protected void validateNewEntry(String key, ArgumentEntry newEntry) { |
|||
if (!(newEntry instanceof SingleValueArgumentEntry)) { |
|||
throw new IllegalArgumentException("Only single value arguments supported"); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public CalculatedFieldType getType() { |
|||
return CalculatedFieldType.ALARM; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,45 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.alarm; |
|||
|
|||
import lombok.Data; |
|||
import lombok.RequiredArgsConstructor; |
|||
|
|||
@Data |
|||
@RequiredArgsConstructor |
|||
public class AlarmEvalResult { |
|||
|
|||
public static final AlarmEvalResult TRUE = new AlarmEvalResult(Status.TRUE); |
|||
public static final AlarmEvalResult FALSE = new AlarmEvalResult(Status.FALSE); |
|||
public static final AlarmEvalResult NOT_YET_TRUE = new AlarmEvalResult(Status.NOT_YET_TRUE); |
|||
|
|||
private final Status status; |
|||
private final long leftDuration; |
|||
private final long leftEvents; |
|||
|
|||
public AlarmEvalResult(Status status) { |
|||
this(status, 0, 0); |
|||
} |
|||
|
|||
public static AlarmEvalResult notYetTrue(long leftEvents, long leftDuration) { |
|||
return new AlarmEvalResult(Status.NOT_YET_TRUE, leftDuration, leftEvents); |
|||
} |
|||
|
|||
public enum Status { |
|||
FALSE, NOT_YET_TRUE, TRUE; |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,344 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.alarm; |
|||
|
|||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|||
import lombok.Data; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.common.util.KvUtil; |
|||
import org.thingsboard.server.common.data.alarm.AlarmSeverity; |
|||
import org.thingsboard.server.common.data.alarm.rule.AlarmRule; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.AlarmCondition; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.AlarmConditionType; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.AlarmConditionValue; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.DurationAlarmCondition; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.RepeatingAlarmCondition; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.expression.AlarmConditionExpression; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.AlarmSchedule; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.AlarmScheduleType; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.AnyTimeSchedule; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.CustomTimeSchedule; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.CustomTimeScheduleItem; |
|||
import org.thingsboard.server.common.data.alarm.rule.condition.schedule.SpecificTimeSchedule; |
|||
import org.thingsboard.server.common.msg.tools.SchedulerUtils; |
|||
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
|||
|
|||
import java.time.Instant; |
|||
import java.time.ZoneId; |
|||
import java.time.ZonedDateTime; |
|||
import java.util.Optional; |
|||
import java.util.concurrent.ScheduledFuture; |
|||
|
|||
@Data |
|||
@Slf4j |
|||
public class AlarmRuleState { |
|||
|
|||
private final AlarmSeverity severity; |
|||
private AlarmRule alarmRule; |
|||
private AlarmCalculatedFieldState state; |
|||
|
|||
private AlarmCondition condition; |
|||
|
|||
private long eventCount; |
|||
private long firstEventTs; // when duration condition started
|
|||
private long lastEventTs; |
|||
private transient long duration; |
|||
private ScheduledFuture<?> durationCheckFuture; |
|||
private Boolean active; |
|||
|
|||
public AlarmRuleState(AlarmSeverity severity, AlarmRule alarmRule, AlarmCalculatedFieldState state) { |
|||
this.severity = severity; |
|||
if (alarmRule != null) { |
|||
setAlarmRule(alarmRule); |
|||
} |
|||
this.state = state; |
|||
} |
|||
|
|||
public AlarmEvalResult eval(boolean newEvent, CalculatedFieldCtx ctx) { // on event or config change
|
|||
long ts = newEvent ? state.getLatestTimestamp() : System.currentTimeMillis(); |
|||
active = isActive(ts); |
|||
if (!active) { |
|||
return AlarmEvalResult.FALSE; |
|||
} |
|||
return doEval(newEvent, ctx); |
|||
} |
|||
|
|||
public AlarmEvalResult reeval(long ts, CalculatedFieldCtx ctx) { // on scheduled duration check or periodic re-eval for rules with schedule
|
|||
boolean active = isActive(ts); |
|||
switch (condition.getType()) { |
|||
case SIMPLE, REPEATING -> { |
|||
if (this.active == null || active != this.active) { |
|||
this.active = active; |
|||
if (active) { |
|||
return doEval(false, ctx); |
|||
} |
|||
} |
|||
if (active) { |
|||
return AlarmEvalResult.NOT_YET_TRUE; |
|||
} else { |
|||
return AlarmEvalResult.FALSE; |
|||
} |
|||
} |
|||
case DURATION -> { |
|||
if (!active) { |
|||
return AlarmEvalResult.FALSE; |
|||
} |
|||
long requiredDuration = getRequiredDurationInMs(); |
|||
if (requiredDuration > 0 && lastEventTs > 0 && ts > lastEventTs) { |
|||
duration = ts - firstEventTs; |
|||
long leftDuration = requiredDuration - duration; |
|||
if (leftDuration <= 0) { |
|||
return AlarmEvalResult.TRUE; |
|||
} else { |
|||
return AlarmEvalResult.notYetTrue(0, leftDuration); |
|||
} |
|||
} |
|||
} |
|||
} |
|||
return AlarmEvalResult.FALSE; |
|||
} |
|||
|
|||
public AlarmEvalResult doEval(boolean newEvent, CalculatedFieldCtx ctx) { |
|||
return switch (condition.getType()) { |
|||
case SIMPLE -> evalSimple(ctx); |
|||
case DURATION -> evalDuration(ctx); |
|||
case REPEATING -> evalRepeating(newEvent, ctx); |
|||
}; |
|||
} |
|||
|
|||
private AlarmEvalResult evalSimple(CalculatedFieldCtx ctx) { |
|||
return eval(condition.getExpression(), ctx) ? AlarmEvalResult.TRUE : AlarmEvalResult.FALSE; |
|||
} |
|||
|
|||
private AlarmEvalResult evalRepeating(boolean newEvent, CalculatedFieldCtx ctx) { |
|||
if (eval(condition.getExpression(), ctx)) { |
|||
if (newEvent) { |
|||
eventCount++; |
|||
} |
|||
long requiredRepeats = getIntValue(((RepeatingAlarmCondition) condition).getCount()); |
|||
if (requiredRepeats > 0) { |
|||
long leftRepeats = requiredRepeats - eventCount; |
|||
return leftRepeats <= 0 ? AlarmEvalResult.TRUE : AlarmEvalResult.notYetTrue(leftRepeats, 0); |
|||
} else { |
|||
return AlarmEvalResult.NOT_YET_TRUE; |
|||
} |
|||
} else { |
|||
return AlarmEvalResult.FALSE; |
|||
} |
|||
} |
|||
|
|||
private AlarmEvalResult evalDuration(CalculatedFieldCtx ctx) { |
|||
if (eval(condition.getExpression(), ctx)) { |
|||
long eventTs = state.getLatestTimestamp(); |
|||
if (lastEventTs > 0) { |
|||
if (eventTs > lastEventTs) { |
|||
if (firstEventTs == 0) { |
|||
firstEventTs = lastEventTs; |
|||
} |
|||
lastEventTs = eventTs; |
|||
} |
|||
} else { |
|||
firstEventTs = eventTs; |
|||
lastEventTs = eventTs; |
|||
} |
|||
duration = lastEventTs - firstEventTs; |
|||
long requiredDuration = getRequiredDurationInMs(); |
|||
if (requiredDuration > 0) { |
|||
long leftDuration = requiredDuration - duration; |
|||
if (leftDuration <= 0) { |
|||
return AlarmEvalResult.TRUE; |
|||
} else { |
|||
return AlarmEvalResult.notYetTrue(0, leftDuration); |
|||
} |
|||
} else { |
|||
return AlarmEvalResult.NOT_YET_TRUE; |
|||
} |
|||
} else { |
|||
return AlarmEvalResult.FALSE; |
|||
} |
|||
} |
|||
|
|||
private boolean isActive(long eventTs) { |
|||
if (condition.getSchedule() == null) { |
|||
return true; |
|||
} |
|||
AlarmSchedule schedule = state.resolveValue(condition.getSchedule(), entry -> Optional.ofNullable(KvUtil.getStringValue(entry)) |
|||
.map(this::parseSchedule).orElse(null)); |
|||
boolean active = switch (schedule.getType()) { |
|||
case ANY_TIME -> true; |
|||
case SPECIFIC_TIME -> isActiveSpecific((SpecificTimeSchedule) schedule, eventTs); |
|||
case CUSTOM -> isActiveCustom((CustomTimeSchedule) schedule, eventTs); |
|||
}; |
|||
log.trace("Alarm rule active = {} for schedule {}", active, schedule); |
|||
return active; |
|||
} |
|||
|
|||
private boolean isActiveSpecific(SpecificTimeSchedule schedule, long eventTs) { |
|||
ZoneId zoneId = SchedulerUtils.getZoneId(schedule.getTimezone()); |
|||
ZonedDateTime zdt = ZonedDateTime.ofInstant(Instant.ofEpochMilli(eventTs), zoneId); |
|||
if (schedule.getDaysOfWeek().size() != 7) { |
|||
int dayOfWeek = zdt.getDayOfWeek().getValue(); |
|||
if (!schedule.getDaysOfWeek().contains(dayOfWeek)) { |
|||
return false; |
|||
} |
|||
} |
|||
long endsOn = schedule.getEndsOn(); |
|||
if (endsOn == 0) { |
|||
// 24 hours in milliseconds
|
|||
endsOn = 86400000; |
|||
} |
|||
|
|||
return isActive(eventTs, zoneId, zdt, schedule.getStartsOn(), endsOn); |
|||
} |
|||
|
|||
private boolean isActiveCustom(CustomTimeSchedule schedule, long eventTs) { |
|||
ZoneId zoneId = SchedulerUtils.getZoneId(schedule.getTimezone()); |
|||
ZonedDateTime zdt = ZonedDateTime.ofInstant(Instant.ofEpochMilli(eventTs), zoneId); |
|||
int dayOfWeek = zdt.toLocalDate().getDayOfWeek().getValue(); |
|||
for (CustomTimeScheduleItem item : schedule.getItems()) { |
|||
if (item.getDayOfWeek() == dayOfWeek) { |
|||
if (item.isEnabled()) { |
|||
long endsOn = item.getEndsOn(); |
|||
if (endsOn == 0) { |
|||
// 24 hours in milliseconds
|
|||
endsOn = 86400000; |
|||
} |
|||
return isActive(eventTs, zoneId, zdt, item.getStartsOn(), endsOn); |
|||
} else { |
|||
return false; |
|||
} |
|||
} |
|||
} |
|||
return false; |
|||
} |
|||
|
|||
private boolean isActive(long eventTs, ZoneId zoneId, ZonedDateTime zdt, long startsOn, long endsOn) { |
|||
long startOfDay = zdt.toLocalDate().atStartOfDay(zoneId).toInstant().toEpochMilli(); |
|||
long msFromStartOfDay = eventTs - startOfDay; |
|||
if (startsOn <= endsOn) { |
|||
return startsOn <= msFromStartOfDay && endsOn > msFromStartOfDay; |
|||
} else { |
|||
return startsOn < msFromStartOfDay || (0 < msFromStartOfDay && msFromStartOfDay < endsOn); |
|||
} |
|||
} |
|||
|
|||
public void clear() { |
|||
clearRepeatingConditionState(); |
|||
clearDurationConditionState(); |
|||
} |
|||
|
|||
private void clearRepeatingConditionState() { |
|||
eventCount = 0L; |
|||
} |
|||
|
|||
private void clearDurationConditionState() { |
|||
firstEventTs = 0L; |
|||
lastEventTs = 0L; |
|||
duration = 0L; |
|||
if (durationCheckFuture != null) { |
|||
durationCheckFuture.cancel(true); |
|||
durationCheckFuture = null; |
|||
} |
|||
} |
|||
|
|||
public boolean isEmpty() { |
|||
return eventCount == 0L && firstEventTs == 0L && lastEventTs == 0L && durationCheckFuture == null; |
|||
} |
|||
|
|||
private AlarmSchedule parseSchedule(String str) { |
|||
ObjectNode json = (ObjectNode) JacksonUtil.toJsonNode(str); |
|||
if (json.isEmpty()) { |
|||
return new AnyTimeSchedule(); // only if valid json, fail otherwise
|
|||
} |
|||
|
|||
if (!json.hasNonNull("type")) { |
|||
// deducting the schedule type
|
|||
AlarmScheduleType type; |
|||
if (json.hasNonNull("daysOfWeek")) { |
|||
type = AlarmScheduleType.SPECIFIC_TIME; |
|||
} else if (json.hasNonNull("items")) { |
|||
type = AlarmScheduleType.CUSTOM; |
|||
} else { |
|||
throw new IllegalArgumentException("Failed to parse alarm schedule from '" + str + "'"); |
|||
} |
|||
json.put("type", type.name()); |
|||
} |
|||
|
|||
return JacksonUtil.treeToValue(json, AlarmSchedule.class); |
|||
} |
|||
|
|||
private Integer getIntValue(AlarmConditionValue<Integer> value) { |
|||
return state.resolveValue(value, entry -> Optional.ofNullable(KvUtil.getLongValue(entry)).map(Long::intValue).orElse(null)); |
|||
} |
|||
|
|||
private long getRequiredDurationInMs() { |
|||
DurationAlarmCondition durationCondition = (DurationAlarmCondition) condition; |
|||
return durationCondition.getUnit().toMillis(state.resolveValue(durationCondition.getValue(), KvUtil::getLongValue)); |
|||
} |
|||
|
|||
private boolean eval(AlarmConditionExpression expression, CalculatedFieldCtx ctx) { |
|||
return state.eval(expression, ctx); |
|||
} |
|||
|
|||
public void setAlarmRule(AlarmRule alarmRule) { |
|||
this.alarmRule = alarmRule; |
|||
this.condition = alarmRule.getCondition(); |
|||
|
|||
// clearing state for other condition types (possibly left from a previous condition type)
|
|||
switch (condition.getType()) { |
|||
case SIMPLE -> { |
|||
clearRepeatingConditionState(); |
|||
clearDurationConditionState(); |
|||
} |
|||
case REPEATING -> { |
|||
clearDurationConditionState(); |
|||
} |
|||
case DURATION -> { |
|||
clearRepeatingConditionState(); |
|||
} |
|||
} |
|||
} |
|||
|
|||
public StateInfo getStateInfo() { |
|||
if (condition.getType() == AlarmConditionType.REPEATING) { |
|||
return new StateInfo(eventCount, null); |
|||
} else if (condition.getType() == AlarmConditionType.DURATION) { |
|||
return new StateInfo(null, duration); |
|||
} else { |
|||
return StateInfo.EMPTY; |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public String toString() { |
|||
return "AlarmRuleState{" + |
|||
"severity=" + severity + |
|||
", condition=" + condition + |
|||
", eventCount=" + eventCount + |
|||
", firstEventTs=" + firstEventTs + |
|||
", lastEventTs=" + lastEventTs + |
|||
", duration=" + duration + |
|||
", durationCheckFuture=" + durationCheckFuture + |
|||
'}'; |
|||
} |
|||
|
|||
public record StateInfo(Long eventCount, Long duration) { |
|||
static final StateInfo EMPTY = new StateInfo(null, null); |
|||
|
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,73 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.propagation; |
|||
|
|||
import lombok.Data; |
|||
import org.thingsboard.script.api.tbel.TbelCfArg; |
|||
import org.thingsboard.script.api.tbel.TbelCfPropagationArg; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.util.CollectionsUtil; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.List; |
|||
|
|||
@Data |
|||
public class PropagationArgumentEntry implements ArgumentEntry { |
|||
|
|||
private List<EntityId> propagationEntityIds; |
|||
|
|||
private boolean forceResetPrevious; |
|||
|
|||
public PropagationArgumentEntry(List<EntityId> propagationEntityIds) { |
|||
this.propagationEntityIds = new ArrayList<>(propagationEntityIds); |
|||
} |
|||
|
|||
@Override |
|||
public ArgumentEntryType getType() { |
|||
return ArgumentEntryType.PROPAGATION; |
|||
} |
|||
|
|||
@Override |
|||
public Object getValue() { |
|||
return propagationEntityIds; |
|||
} |
|||
|
|||
@Override |
|||
public boolean updateEntry(ArgumentEntry entry) { |
|||
if (!(entry instanceof PropagationArgumentEntry propagationArgumentEntry)) { |
|||
throw new IllegalArgumentException("Unsupported argument entry type for propagation argument entry: " + entry.getType()); |
|||
} |
|||
if (propagationArgumentEntry.isEmpty()) { |
|||
propagationEntityIds.clear(); |
|||
} else { |
|||
propagationEntityIds = propagationArgumentEntry.getPropagationEntityIds(); |
|||
} |
|||
return true; |
|||
} |
|||
|
|||
@Override |
|||
public boolean isEmpty() { |
|||
return CollectionsUtil.isEmpty(propagationEntityIds); |
|||
} |
|||
|
|||
@Override |
|||
public TbelCfArg toTbelCfArg() { |
|||
return new TbelCfPropagationArg(propagationEntityIds); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,107 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.cf.ctx.state.propagation; |
|||
|
|||
import com.fasterxml.jackson.databind.node.ObjectNode; |
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import com.google.common.util.concurrent.MoreExecutors; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.actors.TbActorRef; |
|||
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|||
import org.thingsboard.server.common.data.cf.configuration.Output; |
|||
import org.thingsboard.server.common.data.cf.configuration.OutputType; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.service.cf.CalculatedFieldResult; |
|||
import org.thingsboard.server.service.cf.PropagationCalculatedFieldResult; |
|||
import org.thingsboard.server.service.cf.TelemetryCalculatedFieldResult; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
|||
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
|||
import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState; |
|||
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.Map; |
|||
|
|||
import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; |
|||
|
|||
public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState { |
|||
|
|||
public PropagationCalculatedFieldState(EntityId entityId) { |
|||
super(entityId); |
|||
} |
|||
|
|||
@Override |
|||
public void setCtx(CalculatedFieldCtx ctx, TbActorRef actorCtx) { |
|||
this.ctx = ctx; |
|||
this.actorCtx = actorCtx; |
|||
this.requiredArguments = new ArrayList<>(ctx.getArgNames()); |
|||
requiredArguments.add(PROPAGATION_CONFIG_ARGUMENT); |
|||
this.readinessStatus = checkReadiness(requiredArguments, arguments); |
|||
if (ctx.isApplyExpressionForResolvedArguments()) { |
|||
this.tbelExpression = ctx.getTbelExpressions().get(ctx.getExpression()); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public CalculatedFieldType getType() { |
|||
return CalculatedFieldType.PROPAGATION; |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) { |
|||
ArgumentEntry argumentEntry = arguments.get(PROPAGATION_CONFIG_ARGUMENT); |
|||
if (!(argumentEntry instanceof PropagationArgumentEntry propagationArgumentEntry) || propagationArgumentEntry.isEmpty()) { |
|||
return Futures.immediateFuture(PropagationCalculatedFieldResult.builder().build()); |
|||
} |
|||
if (ctx.isApplyExpressionForResolvedArguments()) { |
|||
return Futures.transform(super.performCalculation(updatedArgs, ctx), telemetryCfResult -> |
|||
PropagationCalculatedFieldResult.builder() |
|||
.propagationEntityIds(propagationArgumentEntry.getPropagationEntityIds()) |
|||
.result((TelemetryCalculatedFieldResult) telemetryCfResult) |
|||
.build(), |
|||
MoreExecutors.directExecutor()); |
|||
} |
|||
return Futures.immediateFuture(PropagationCalculatedFieldResult.builder() |
|||
.propagationEntityIds(propagationArgumentEntry.getPropagationEntityIds()) |
|||
.result(toTelemetryResult(ctx)) |
|||
.build()); |
|||
} |
|||
|
|||
private TelemetryCalculatedFieldResult toTelemetryResult(CalculatedFieldCtx ctx) { |
|||
Output output = ctx.getOutput(); |
|||
TelemetryCalculatedFieldResult.TelemetryCalculatedFieldResultBuilder telemetryCfBuilder = |
|||
TelemetryCalculatedFieldResult.builder() |
|||
.type(output.getType()) |
|||
.scope(output.getScope()); |
|||
ObjectNode valuesNode = JacksonUtil.newObjectNode(); |
|||
arguments.forEach((outputKey, argumentEntry) -> { |
|||
if (argumentEntry instanceof PropagationArgumentEntry) { |
|||
return; |
|||
} |
|||
if (argumentEntry instanceof SingleValueArgumentEntry singleArgumentEntry) { |
|||
JacksonUtil.addKvEntry(valuesNode, singleArgumentEntry.getKvEntryValue(), outputKey); |
|||
return; |
|||
} |
|||
throw new IllegalArgumentException("Unsupported argument type: " + argumentEntry.getType() + " detected for argument: " + outputKey + ". " + |
|||
"Only Latest telemetry or Attribute arguments supported for 'Arguments Only' propagation mode!"); |
|||
}); |
|||
ObjectNode result = toSimpleResult(output.getType() == OutputType.TIME_SERIES, valuesNode); |
|||
telemetryCfBuilder.result(result); |
|||
return telemetryCfBuilder.build(); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,130 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.edge.rpc.processor.ai; |
|||
|
|||
import com.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.data.util.Pair; |
|||
import org.springframework.stereotype.Component; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.EdgeUtils; |
|||
import org.thingsboard.server.common.data.ai.AiModel; |
|||
import org.thingsboard.server.common.data.edge.Edge; |
|||
import org.thingsboard.server.common.data.edge.EdgeEvent; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventActionType; |
|||
import org.thingsboard.server.common.data.edge.EdgeEventType; |
|||
import org.thingsboard.server.common.data.id.AiModelId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.msg.TbMsgType; |
|||
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|||
import org.thingsboard.server.dao.exception.DataValidationException; |
|||
import org.thingsboard.server.gen.edge.v1.AiModelUpdateMsg; |
|||
import org.thingsboard.server.gen.edge.v1.DownlinkMsg; |
|||
import org.thingsboard.server.gen.edge.v1.EdgeVersion; |
|||
import org.thingsboard.server.gen.edge.v1.UpdateMsgType; |
|||
import org.thingsboard.server.queue.util.TbCoreComponent; |
|||
import org.thingsboard.server.service.edge.EdgeMsgConstructorUtils; |
|||
|
|||
import java.util.Optional; |
|||
import java.util.UUID; |
|||
|
|||
@Slf4j |
|||
@Component |
|||
@TbCoreComponent |
|||
public class AiModelEdgeProcessor extends BaseAiModelProcessor implements AiModelProcessor { |
|||
|
|||
@Override |
|||
public ListenableFuture<Void> processAiModelMsgFromEdge(TenantId tenantId, Edge edge, AiModelUpdateMsg aiModelUpdateMsg) { |
|||
AiModelId aiModelId = new AiModelId(new UUID(aiModelUpdateMsg.getIdMSB(), aiModelUpdateMsg.getIdLSB())); |
|||
try { |
|||
edgeSynchronizationManager.getEdgeId().set(edge.getId()); |
|||
|
|||
switch (aiModelUpdateMsg.getMsgType()) { |
|||
case ENTITY_CREATED_RPC_MESSAGE: |
|||
case ENTITY_UPDATED_RPC_MESSAGE: |
|||
processAiModel(tenantId, aiModelId, aiModelUpdateMsg, edge); |
|||
return Futures.immediateFuture(null); |
|||
case UNRECOGNIZED: |
|||
default: |
|||
return handleUnsupportedMsgType(aiModelUpdateMsg.getMsgType()); |
|||
} |
|||
} catch (DataValidationException e) { |
|||
return Futures.immediateFailedFuture(e); |
|||
} finally { |
|||
edgeSynchronizationManager.getEdgeId().remove(); |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public DownlinkMsg convertEdgeEventToDownlink(EdgeEvent edgeEvent, EdgeVersion edgeVersion) { |
|||
AiModelId aiModelId = new AiModelId(edgeEvent.getEntityId()); |
|||
switch (edgeEvent.getAction()) { |
|||
case ADDED, UPDATED -> { |
|||
Optional<AiModel> aiModel = edgeCtx.getAiModelService().findAiModelById(edgeEvent.getTenantId(), aiModelId); |
|||
if (aiModel.isPresent()) { |
|||
UpdateMsgType msgType = getUpdateMsgType(edgeEvent.getAction()); |
|||
AiModelUpdateMsg aiModelUpdateMsg = EdgeMsgConstructorUtils.constructAiModelUpdatedMsg(msgType, aiModel.get()); |
|||
return DownlinkMsg.newBuilder() |
|||
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|||
.addAiModelUpdateMsg(aiModelUpdateMsg) |
|||
.build(); |
|||
} |
|||
} |
|||
case DELETED -> { |
|||
AiModelUpdateMsg aiModelUpdateMsg = EdgeMsgConstructorUtils.constructAiModelDeleteMsg(aiModelId); |
|||
return DownlinkMsg.newBuilder() |
|||
.setDownlinkMsgId(EdgeUtils.nextPositiveInt()) |
|||
.addAiModelUpdateMsg(aiModelUpdateMsg) |
|||
.build(); |
|||
} |
|||
} |
|||
return null; |
|||
} |
|||
|
|||
@Override |
|||
public EdgeEventType getEdgeEventType() { |
|||
return EdgeEventType.AI_MODEL; |
|||
} |
|||
|
|||
private void processAiModel(TenantId tenantId, AiModelId aiModelId, AiModelUpdateMsg aiModelUpdateMsg, Edge edge) { |
|||
Pair<Boolean, Boolean> resultPair = super.saveOrUpdateAiModel(tenantId, aiModelId, aiModelUpdateMsg); |
|||
Boolean wasCreated = resultPair.getFirst(); |
|||
if (wasCreated) { |
|||
pushAiModelCreatedEventToRuleEngine(tenantId, edge, aiModelId); |
|||
} |
|||
Boolean nameWasUpdated = resultPair.getSecond(); |
|||
if (nameWasUpdated) { |
|||
saveEdgeEvent(tenantId, edge.getId(), EdgeEventType.AI_MODEL, EdgeEventActionType.UPDATED, aiModelId, null); |
|||
} |
|||
} |
|||
|
|||
private void pushAiModelCreatedEventToRuleEngine(TenantId tenantId, Edge edge, AiModelId aiModelId) { |
|||
try { |
|||
Optional<AiModel> aiModel = edgeCtx.getAiModelService().findAiModelById(tenantId, aiModelId); |
|||
if (aiModel.isPresent()) { |
|||
String aiModelAsString = JacksonUtil.toString(aiModel.get()); |
|||
TbMsgMetaData msgMetaData = getEdgeActionTbMsgMetaData(edge, edge.getCustomerId()); |
|||
pushEntityEventToRuleEngine(tenantId, aiModelId, edge.getCustomerId(), TbMsgType.ENTITY_CREATED, aiModelAsString, msgMetaData); |
|||
} else { |
|||
log.warn("[{}][{}] Failed to find aiModel", tenantId, aiModelId); |
|||
} |
|||
} catch (Exception e) { |
|||
log.warn("[{}][{}] Failed to push aiModel action to rule engine: {}", tenantId, aiModelId, TbMsgType.ENTITY_CREATED.name(), e); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,28 @@ |
|||
/** |
|||
* Copyright © 2016-2025 The Thingsboard Authors |
|||
* |
|||
* Licensed under the Apache License, Version 2.0 (the "License"); |
|||
* you may not use this file except in compliance with the License. |
|||
* You may obtain a copy of the License at |
|||
* |
|||
* http://www.apache.org/licenses/LICENSE-2.0
|
|||
* |
|||
* Unless required by applicable law or agreed to in writing, software |
|||
* distributed under the License is distributed on an "AS IS" BASIS, |
|||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
|||
* See the License for the specific language governing permissions and |
|||
* limitations under the License. |
|||
*/ |
|||
package org.thingsboard.server.service.edge.rpc.processor.ai; |
|||
|
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import org.thingsboard.server.common.data.edge.Edge; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.gen.edge.v1.AiModelUpdateMsg; |
|||
import org.thingsboard.server.service.edge.rpc.processor.EdgeProcessor; |
|||
|
|||
public interface AiModelProcessor extends EdgeProcessor { |
|||
|
|||
ListenableFuture<Void> processAiModelMsgFromEdge(TenantId tenantId, Edge edge, AiModelUpdateMsg aiModelUpdateMsg); |
|||
|
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue