938 changed files with 51406 additions and 16734 deletions
@ -1,43 +1,150 @@ |
|||
# ThingsBoard |
|||
[](https://builds.thingsboard.io/viewType.html?buildTypeId=ThingsBoard_Build&guest=1) |
|||
 |
|||
|
|||
<div align="center"> |
|||
|
|||
# Open-source IoT platform for data collection, processing, visualization, and device management. |
|||
|
|||
</div> |
|||
<br> |
|||
<div align="center"> |
|||
|
|||
💡 [Get started](https://thingsboard.io/docs/getting-started-guides/helloworld/) • 🌐 [Website](https://thingsboard.io/) • 📚 [Documentation](https://thingsboard.io/docs/) • 📔 [Blog](https://thingsboard.io/blog/) • ▶️ [Live demo](https://demo.thingsboard.io/signup) • 🔗 [LinkedIn](https://www.linkedin.com/company/thingsboard/posts/?feedView=all) |
|||
|
|||
</div> |
|||
|
|||
## 🚀 Installation options |
|||
|
|||
* Install ThingsBoard [On-premise](https://thingsboard.io/docs/user-guide/install/installation-options/?ceInstallType=onPremise) |
|||
* Try [ThingsBoard Cloud](https://thingsboard.io/installations/) |
|||
* or [Use our Live demo](https://demo.thingsboard.io/signup) |
|||
|
|||
## 💡 Getting started with ThingsBoard |
|||
|
|||
Check out our [Getting Started guide](https://thingsboard.io/docs/getting-started-guides/helloworld/) or [watch the video](https://www.youtube.com/watch?v=80L0ubQLXsc) to learn the basics of ThingsBoard and create your first dashboard! You will learn to: |
|||
|
|||
* Connect devices to ThingsBoard |
|||
* Push data from devices to ThingsBoard |
|||
* Build real-time dashboards |
|||
* Create a Customer and assign the dashboard with them. |
|||
* Define thresholds and trigger alarms |
|||
* Set up notifications via email, SMS, mobile apps, or integrate with third-party services. |
|||
|
|||
## ✨ Features |
|||
|
|||
<table> |
|||
<tr> |
|||
<td width="50%" valign="top"> |
|||
<br> |
|||
<div align="center"> |
|||
<img src="https://github.com/user-attachments/assets/255cca4f-b111-44e8-99ea-0af55f8e3681" alt="Provision and manage devices and assets" width="378" /> |
|||
<h3>Provision and manage <br> devices and assets</h3> |
|||
</div> |
|||
<div align="center"> |
|||
<p>Provision, monitor and control your IoT entities in secure way using rich server-side APIs. Define relations between your devices, assets, customers or any other entities.</p> |
|||
</div> |
|||
<br> |
|||
<div align="center"> |
|||
<a href="https://thingsboard.io/docs/user-guide/entities-and-relations/">Read more ➜</a> |
|||
</div> |
|||
<br> |
|||
</td> |
|||
<td width="50%" valign="top"> |
|||
<br> |
|||
<div align="center"> |
|||
<img src="https://github.com/user-attachments/assets/24b41d10-150a-42dd-ab1a-32ac9b5978c1" alt="Collect and visualize your data" width="378" /> |
|||
<h3>Collect and visualize <br> your data</h3> |
|||
</div> |
|||
<div align="center"> |
|||
<p>Collect and store telemetry data in scalable and fault-tolerant way. Visualize your data with built-in or custom widgets and flexible dashboards. Share dashboards with your customers.</p> |
|||
</div> |
|||
<br> |
|||
<div align="center"> |
|||
<a href="https://thingsboard.io/iot-data-visualization/">Read more ➜</a> |
|||
</div> |
|||
<br> |
|||
</td> |
|||
</tr> |
|||
<tr> |
|||
<td width="50%" valign="top"> |
|||
<br> |
|||
<div align="center"> |
|||
<img src="https://github.com/user-attachments/assets/6f2a6dd2-7b33-4d17-8b92-d1f995adda2c" alt="SCADA Dashboards" width="378" /> |
|||
<h3>SCADA Dashboards</h3> |
|||
</div> |
|||
<div align="center"> |
|||
<p>Monitor and control your industrial processes in real time with SCADA. Use SCADA symbols on dashboards to create and manage any workflow, offering full flexibility to design and oversee operations according to your requirements.</p> |
|||
</div> |
|||
<br> |
|||
<div align="center"> |
|||
<a href="https://thingsboard.io/use-cases/scada/">Read more ➜</a> |
|||
</div> |
|||
<br> |
|||
</td> |
|||
<td width="50%" valign="top"> |
|||
<br> |
|||
<div align="center"> |
|||
<img src="https://github.com/user-attachments/assets/c23dcc9b-aeba-40ef-9973-49b953fc1257" alt="Process and React" width="378" /> |
|||
<h3>Process and React</h3> |
|||
</div> |
|||
<div align="center"> |
|||
<p>Define data processing rule chains. Transform and normalize your device data. Raise alarms on incoming telemetry events, attribute updates, device inactivity and user actions.<br></p> |
|||
</div> |
|||
<br> |
|||
<br> |
|||
<div align="center"> |
|||
<a href="https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/">Read more ➜</a> |
|||
</div> |
|||
<br> |
|||
</td> |
|||
</tr> |
|||
</table> |
|||
|
|||
## ⚙️ Powerful IoT Rule Engine |
|||
|
|||
ThingsBoard allows you to create complex [Rule Chains](https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/) to process data from your devices and match your application specific use cases. |
|||
|
|||
[](https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/) |
|||
|
|||
<div align="center"> |
|||
|
|||
[**Read more about Rule Engine ➜**](https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/) |
|||
|
|||
</div> |
|||
|
|||
## 📦 Real-Time IoT Dashboards |
|||
|
|||
ThingsBoard is a scalable, user-friendly, and device-agnostic IoT platform that speeds up time-to-market with powerful built-in solution templates. It enables data collection and analysis from any devices, saving resources on routine tasks and letting you focus on your solution’s unique aspects. See more our Use Cases [here](https://thingsboard.io/iot-use-cases/). |
|||
|
|||
[**Smart energy**](https://thingsboard.io/use-cases/smart-energy/) |
|||
|
|||
[](https://thingsboard.io/use-cases/smart-energy/) |
|||
|
|||
[**SCADA swimming pool**](https://thingsboard.io/use-cases/scada/) |
|||
|
|||
[](https://thingsboard.io/use-cases/scada/) |
|||
|
|||
[**Fleet tracking**](https://thingsboard.io/use-cases/fleet-tracking/) |
|||
|
|||
[](https://thingsboard.io/use-cases/fleet-tracking/) |
|||
|
|||
[**Smart farming**](https://thingsboard.io/use-cases/smart-farming/) |
|||
|
|||
[](https://thingsboard.io/use-cases/smart-farming/) |
|||
|
|||
ThingsBoard is an open-source IoT platform for data collection, processing, visualization, and device management. |
|||
|
|||
<img src="./img/logo.png?raw=true" width="100" height="100"> |
|||
|
|||
|
|||
## Documentation |
|||
|
|||
ThingsBoard documentation is hosted on [thingsboard.io](https://thingsboard.io/docs). |
|||
|
|||
## IoT use cases |
|||
|
|||
[**Smart energy**](https://thingsboard.io/smart-energy/) |
|||
[](https://thingsboard.io/smart-energy/) |
|||
|
|||
[**SCADA Swimming pool**](https://thingsboard.io/use-cases/scada/) |
|||
[](https://thingsboard.io/use-cases/scada/) |
|||
|
|||
[**Fleet tracking**](https://thingsboard.io/fleet-tracking/) |
|||
[](https://thingsboard.io/fleet-tracking/) |
|||
|
|||
[**Smart farming**](https://thingsboard.io/smart-farming/) |
|||
[](https://thingsboard.io/smart-farming/) |
|||
[**Smart metering**](https://thingsboard.io/smart-metering/) |
|||
|
|||
[**IoT Rule Engine**](https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/) |
|||
[](https://thingsboard.io/docs/user-guide/rule-engine-2-0/re-getting-started/) |
|||
[](https://thingsboard.io/smart-metering/) |
|||
|
|||
[**Smart metering**](https://thingsboard.io/smart-metering/) |
|||
[](https://thingsboard.io/smart-metering/) |
|||
<div align="center"> |
|||
|
|||
## Getting Started |
|||
[**Check more of our use cases ➜**](https://thingsboard.io/iot-use-cases/) |
|||
|
|||
Collect and Visualize your IoT data in minutes by following this [guide](https://thingsboard.io/docs/getting-started-guides/helloworld/). |
|||
</div> |
|||
|
|||
## Support |
|||
## 🫶 Support |
|||
|
|||
- [Stackoverflow](http://stackoverflow.com/questions/tagged/thingsboard) |
|||
To get support, please visit our [GitHub issues page](https://github.com/thingsboard/thingsboard/issues) |
|||
|
|||
## Licenses |
|||
## 📄 Licenses |
|||
|
|||
This project is released under [Apache 2.0 License](./LICENSE). |
|||
This project is released under [Apache 2.0 License](./LICENSE) |
|||
|
|||
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,37 @@ |
|||
/** |
|||
* 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.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 CalculatedFieldArgumentResetMsg implements ToCalculatedFieldSystemMsg { |
|||
|
|||
private final TenantId tenantId; |
|||
private final CalculatedFieldCtx ctx; |
|||
private final TbCallback callback; |
|||
|
|||
@Override |
|||
public MsgType getMsgType() { |
|||
return MsgType.CF_ARGUMENT_RESET_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,35 @@ |
|||
/** |
|||
* 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.id.TenantId; |
|||
import org.thingsboard.server.common.msg.MsgType; |
|||
import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; |
|||
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; |
|||
|
|||
@Data |
|||
public class CalculatedFieldReevaluateMsg implements ToCalculatedFieldSystemMsg { |
|||
|
|||
private final TenantId tenantId; |
|||
private final CalculatedFieldCtx ctx; |
|||
|
|||
@Override |
|||
public MsgType getMsgType() { |
|||
return MsgType.CF_REEVALUATE_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,332 @@ |
|||
/** |
|||
* 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.google.common.util.concurrent.Futures; |
|||
import com.google.common.util.concurrent.ListenableFuture; |
|||
import com.google.common.util.concurrent.ListeningExecutorService; |
|||
import com.google.common.util.concurrent.MoreExecutors; |
|||
import jakarta.annotation.PostConstruct; |
|||
import jakarta.annotation.PreDestroy; |
|||
import lombok.Data; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.common.util.ThingsBoardExecutors; |
|||
import org.thingsboard.server.common.data.cf.CalculatedField; |
|||
import org.thingsboard.server.common.data.cf.configuration.Argument; |
|||
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
|||
import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; |
|||
import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.kv.Aggregation; |
|||
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|||
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
|||
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; |
|||
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
|||
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; |
|||
import org.thingsboard.server.common.data.kv.TsKvEntry; |
|||
import org.thingsboard.server.common.data.relation.EntityRelation; |
|||
import org.thingsboard.server.common.data.relation.EntityRelationPathQuery; |
|||
import org.thingsboard.server.common.data.relation.RelationPathLevel; |
|||
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; |
|||
import org.thingsboard.server.dao.attributes.AttributesService; |
|||
import org.thingsboard.server.dao.relation.RelationService; |
|||
import org.thingsboard.server.dao.timeseries.TimeseriesService; |
|||
import org.thingsboard.server.dao.usagerecord.ApiLimitService; |
|||
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.SingleValueArgumentEntry; |
|||
|
|||
import java.util.Collections; |
|||
import java.util.HashMap; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.Optional; |
|||
import java.util.Set; |
|||
import java.util.concurrent.ExecutionException; |
|||
import java.util.function.Predicate; |
|||
import java.util.stream.Collectors; |
|||
|
|||
import static org.thingsboard.server.common.data.cf.CalculatedFieldType.PROPAGATION; |
|||
import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; |
|||
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY; |
|||
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY; |
|||
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultAttributeEntry; |
|||
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry; |
|||
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformSingleValueArgument; |
|||
|
|||
@Data |
|||
@Slf4j |
|||
public abstract class AbstractCalculatedFieldProcessingService { |
|||
|
|||
protected final AttributesService attributesService; |
|||
protected final TimeseriesService timeseriesService; |
|||
protected final ApiLimitService apiLimitService; |
|||
protected final RelationService relationService; |
|||
protected final OwnerService ownerService; |
|||
|
|||
protected ListeningExecutorService calculatedFieldCallbackExecutor; |
|||
|
|||
@PostConstruct |
|||
public void init() { |
|||
calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( |
|||
Math.max(4, Runtime.getRuntime().availableProcessors()), getExecutorNamePrefix())); |
|||
} |
|||
|
|||
@PreDestroy |
|||
public void stop() { |
|||
if (calculatedFieldCallbackExecutor != null) { |
|||
calculatedFieldCallbackExecutor.shutdownNow(); |
|||
} |
|||
} |
|||
|
|||
protected abstract String getExecutorNamePrefix(); |
|||
|
|||
protected ListenableFuture<Map<String, ArgumentEntry>> fetchArguments(CalculatedFieldCtx ctx, EntityId entityId, long ts) { |
|||
Map<String, ListenableFuture<ArgumentEntry>> argFutures = switch (ctx.getCfType()) { |
|||
case GEOFENCING -> fetchGeofencingCalculatedFieldArguments(ctx, entityId, false, ts); |
|||
case SIMPLE, SCRIPT, ALARM, PROPAGATION -> getBaseCalculatedFieldArguments(ctx, entityId, ts); |
|||
case RELATED_ENTITIES_AGGREGATION -> fetchRelatedEntitiesAggArguments(ctx, entityId, ts); |
|||
}; |
|||
if (ctx.getCfType() == PROPAGATION) { |
|||
argFutures.put(PROPAGATION_CONFIG_ARGUMENT, fetchPropagationCalculatedFieldArgument(ctx, entityId)); |
|||
} |
|||
return Futures.whenAllComplete(argFutures.values()) |
|||
.call(() -> resolveArgumentFutures(argFutures), |
|||
MoreExecutors.directExecutor()); |
|||
} |
|||
|
|||
private Map<String, ListenableFuture<ArgumentEntry>> getBaseCalculatedFieldArguments(CalculatedFieldCtx ctx, EntityId entityId, long ts) { |
|||
Map<String, ListenableFuture<ArgumentEntry>> futures = new HashMap<>(); |
|||
for (var entry : ctx.getArguments().entrySet()) { |
|||
var argEntityId = resolveEntityId(ctx.getTenantId(), entityId, entry.getValue()); |
|||
var argValueFuture = fetchArgumentValue(ctx.getTenantId(), argEntityId, entry.getValue(), ts); |
|||
futures.put(entry.getKey(), argValueFuture); |
|||
} |
|||
return futures; |
|||
} |
|||
|
|||
protected EntityId resolveEntityId(TenantId tenantId, EntityId entityId, Argument argument) { |
|||
if (argument.getRefEntityId() != null) { |
|||
return argument.getRefEntityId(); |
|||
} |
|||
if (!argument.hasOwnerSource()) { |
|||
return entityId; |
|||
} |
|||
return resolveOwnerArgument(tenantId, entityId); |
|||
} |
|||
|
|||
protected Map<String, ArgumentEntry> resolveArgumentFutures(Map<String, ListenableFuture<ArgumentEntry>> argFutures) { |
|||
return argFutures.entrySet().stream() |
|||
.collect(Collectors.toMap( |
|||
Map.Entry::getKey, // Keep the key as is
|
|||
entry -> { |
|||
try { |
|||
return entry.getValue().get(); |
|||
} catch (ExecutionException e) { |
|||
Throwable cause = e.getCause(); |
|||
throw new RuntimeException("Failed to fetch " + entry.getKey() + ": " + cause.getMessage(), cause); |
|||
} catch (InterruptedException e) { |
|||
throw new RuntimeException("Failed to fetch" + entry.getKey(), e); |
|||
} |
|||
} |
|||
)); |
|||
} |
|||
|
|||
protected ListenableFuture<ArgumentEntry> fetchPropagationCalculatedFieldArgument(CalculatedFieldCtx ctx, EntityId entityId) { |
|||
ListenableFuture<List<EntityId>> propagationEntityIds = fromDynamicSource(ctx.getTenantId(), entityId, ctx.getPropagationArgument()); |
|||
return Futures.transform(propagationEntityIds, ArgumentEntry::createPropagationArgument, MoreExecutors.directExecutor()); |
|||
} |
|||
|
|||
protected Map<String, ListenableFuture<ArgumentEntry>> fetchGeofencingCalculatedFieldArguments(CalculatedFieldCtx ctx, EntityId entityId, boolean dynamicArgumentsOnly, long startTs) { |
|||
Map<String, ListenableFuture<ArgumentEntry>> argFutures = new HashMap<>(); |
|||
Set<Map.Entry<String, Argument>> entries = ctx.getArguments().entrySet(); |
|||
if (dynamicArgumentsOnly) { |
|||
entries = entries.stream() |
|||
.filter(entry -> entry.getValue().hasRelationQuerySource()) |
|||
.collect(Collectors.toSet()); |
|||
} |
|||
for (var entry : entries) { |
|||
switch (entry.getKey()) { |
|||
case ENTITY_ID_LATITUDE_ARGUMENT_KEY, ENTITY_ID_LONGITUDE_ARGUMENT_KEY -> |
|||
argFutures.put(entry.getKey(), fetchArgumentValue(ctx.getTenantId(), entityId, entry.getValue(), startTs)); |
|||
default -> { |
|||
var resolvedEntityIdsFuture = resolveGeofencingEntityIds(ctx.getTenantId(), entityId, entry); |
|||
argFutures.put(entry.getKey(), Futures.transformAsync(resolvedEntityIdsFuture, resolvedEntityIds -> |
|||
fetchGeofencingKvEntry(ctx.getTenantId(), resolvedEntityIds, entry.getValue()), MoreExecutors.directExecutor())); |
|||
} |
|||
} |
|||
} |
|||
return argFutures; |
|||
} |
|||
|
|||
protected Map<String, ListenableFuture<ArgumentEntry>> fetchRelatedEntitiesAggArguments(CalculatedFieldCtx ctx, EntityId entityId, long ts) { |
|||
if (!(ctx.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config)) { |
|||
return Collections.emptyMap(); |
|||
} |
|||
ListenableFuture<List<EntityId>> relatedEntitiesFut = resolveRelatedEntities(ctx.getTenantId(), entityId, config.getRelation()); |
|||
|
|||
return config.getArguments().entrySet().stream() |
|||
.collect(Collectors.toMap( |
|||
Map.Entry::getKey, |
|||
entry -> Futures.transformAsync(relatedEntitiesFut, relatedEntities -> fetchRelatedEntitiesArgumentEntry(ctx.getTenantId(), relatedEntities, entry.getValue(), ts), MoreExecutors.directExecutor()) |
|||
)); |
|||
} |
|||
|
|||
protected ListenableFuture<List<EntityId>> resolveRelatedEntities(TenantId tenantId, EntityId entityId, RelationPathLevel relation) { |
|||
Predicate<EntityRelation> filter = entityRelation -> CalculatedField.isSupportedRefEntity(entityRelation.getFrom()) && CalculatedField.isSupportedRefEntity(entityRelation.getTo()); |
|||
ListenableFuture<List<EntityRelation>> relationsFut = relationService.findFilteredRelationsByPathQueryAsync(tenantId, new EntityRelationPathQuery(entityId, List.of(relation)), filter); |
|||
|
|||
return Futures.transform(relationsFut, relations -> { |
|||
if (relations == null) { |
|||
return Collections.emptyList(); |
|||
} |
|||
|
|||
return switch (relation.direction()) { |
|||
case FROM -> relations.stream() |
|||
.map(EntityRelation::getTo) |
|||
.toList(); |
|||
case TO -> relations.stream() |
|||
.map(EntityRelation::getFrom) |
|||
.findFirst() |
|||
.map(List::of) |
|||
.orElseGet(Collections::emptyList); |
|||
}; |
|||
}, calculatedFieldCallbackExecutor); |
|||
} |
|||
|
|||
private ListenableFuture<List<EntityId>> resolveGeofencingEntityIds(TenantId tenantId, EntityId entityId, Map.Entry<String, Argument> entry) { |
|||
Argument value = entry.getValue(); |
|||
if (value.getRefEntityId() != null) { |
|||
return Futures.immediateFuture(List.of(value.getRefEntityId())); |
|||
} |
|||
if (!value.hasDynamicSource()) { |
|||
return Futures.immediateFuture(List.of(entityId)); |
|||
} |
|||
return fromDynamicSource(tenantId, entityId, value); |
|||
} |
|||
|
|||
private ListenableFuture<List<EntityId>> fromDynamicSource(TenantId tenantId, EntityId entityId, Argument value) { |
|||
var refDynamicSourceConfiguration = value.getRefDynamicSourceConfiguration(); |
|||
return switch (refDynamicSourceConfiguration.getType()) { |
|||
case CURRENT_OWNER -> Futures.immediateFuture(List.of(resolveOwnerArgument(tenantId, entityId))); |
|||
case RELATION_PATH_QUERY -> { |
|||
var configuration = (RelationPathQueryDynamicSourceConfiguration) refDynamicSourceConfiguration; |
|||
Predicate<EntityRelation> filter = entityRelation -> CalculatedField.isSupportedRefEntity(entityRelation.getFrom()) && CalculatedField.isSupportedRefEntity(entityRelation.getTo()); |
|||
yield Futures.transform(relationService.findFilteredRelationsByPathQueryAsync(tenantId, configuration.toRelationPathQuery(entityId), filter), |
|||
configuration::resolveEntityIds, calculatedFieldCallbackExecutor); |
|||
} |
|||
}; |
|||
} |
|||
|
|||
private EntityId resolveOwnerArgument(TenantId tenantId, EntityId entityId) { |
|||
return ownerService.getOwner(tenantId, entityId); |
|||
} |
|||
|
|||
private ListenableFuture<ArgumentEntry> fetchGeofencingKvEntry(TenantId tenantId, List<EntityId> geofencingEntities, Argument argument) { |
|||
if (argument.getRefEntityKey().getType() != ArgumentType.ATTRIBUTE) { |
|||
throw new IllegalStateException("Unsupported argument key type: " + argument.getRefEntityKey().getType()); |
|||
} |
|||
List<ListenableFuture<Map.Entry<EntityId, AttributeKvEntry>>> kvFutures = geofencingEntities.stream() |
|||
.map(entityId -> { |
|||
var attributesFuture = attributesService.find( |
|||
tenantId, |
|||
entityId, |
|||
argument.getRefEntityKey().getScope(), |
|||
argument.getRefEntityKey().getKey() |
|||
); |
|||
return Futures.transform(attributesFuture, resultOpt -> |
|||
Map.entry(entityId, resultOpt.orElseGet(() -> createDefaultAttributeEntry(argument, System.currentTimeMillis()))), |
|||
calculatedFieldCallbackExecutor |
|||
); |
|||
}).collect(Collectors.toList()); |
|||
|
|||
ListenableFuture<List<Map.Entry<EntityId, AttributeKvEntry>>> allFutures = Futures.allAsList(kvFutures); |
|||
|
|||
return Futures.transform(allFutures, entries -> ArgumentEntry.createGeofencingValueArgument(entries.stream() |
|||
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue))), MoreExecutors.directExecutor()); |
|||
} |
|||
|
|||
public ListenableFuture<ArgumentEntry> fetchRelatedEntitiesArgumentEntry(TenantId tenantId, List<EntityId> aggEntities, Argument argument, long startTs) { |
|||
List<ListenableFuture<Map.Entry<EntityId, ArgumentEntry>>> futures = aggEntities.stream() |
|||
.map(entityId -> { |
|||
ListenableFuture<ArgumentEntry> argumentEntryFut = fetchArgumentValue(tenantId, entityId, argument, startTs); |
|||
return Futures.transform(argumentEntryFut, argumentEntry -> Map.entry(entityId, ArgumentEntry.createSingleValueArgument(entityId, argumentEntry)), MoreExecutors.directExecutor()); |
|||
}) |
|||
.toList(); |
|||
|
|||
ListenableFuture<List<Map.Entry<EntityId, ? extends ArgumentEntry>>> allFutures = Futures.allAsList(futures); |
|||
|
|||
return Futures.transform(allFutures, |
|||
entries -> ArgumentEntry.createAggArgument( |
|||
entries.stream().collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)) |
|||
), |
|||
MoreExecutors.directExecutor()); |
|||
} |
|||
|
|||
protected ListenableFuture<ArgumentEntry> fetchArgumentValue(TenantId tenantId, EntityId entityId, Argument argument, long startTs) { |
|||
return switch (argument.getRefEntityKey().getType()) { |
|||
case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument, startTs); |
|||
case ATTRIBUTE -> fetchAttribute(tenantId, entityId, argument, startTs); |
|||
case TS_LATEST -> fetchTsLatest(tenantId, entityId, argument, startTs); |
|||
}; |
|||
} |
|||
|
|||
private ListenableFuture<ArgumentEntry> fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument, long queryEndTs) { |
|||
long argTimeWindow = argument.getTimeWindow() == 0 ? queryEndTs : argument.getTimeWindow(); |
|||
long startInterval = queryEndTs - argTimeWindow; |
|||
ReadTsKvQuery query = buildTsRollingQuery(tenantId, argument, startInterval, queryEndTs); |
|||
|
|||
log.trace("[{}][{}] Fetching timeseries for query {}", tenantId, entityId, query); |
|||
ListenableFuture<List<TsKvEntry>> tsRollingFuture = timeseriesService.findAll(tenantId, entityId, List.of(query)); |
|||
return Futures.transform(tsRollingFuture, tsRolling -> { |
|||
log.debug("[{}][{}] Fetched {} timeseries for query {}", tenantId, entityId, tsRolling == null ? 0 : tsRolling.size(), query); |
|||
return ArgumentEntry.createTsRollingArgument(tsRolling, query.getLimit(), argTimeWindow); |
|||
}, calculatedFieldCallbackExecutor); |
|||
} |
|||
|
|||
private ListenableFuture<ArgumentEntry> fetchAttribute(TenantId tenantId, EntityId entityId, Argument argument, long defaultLastUpdateTs) { |
|||
log.trace("[{}][{}] Fetching attribute for key {}", tenantId, entityId, argument.getRefEntityKey()); |
|||
var attributeOptFuture = attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey()); |
|||
|
|||
return Futures.transform(attributeOptFuture, attrOpt -> { |
|||
log.debug("[{}][{}] Fetched attribute for key {}: {}", tenantId, entityId, argument.getRefEntityKey(), attrOpt); |
|||
AttributeKvEntry attributeKvEntry = attrOpt.orElseGet(() -> new BaseAttributeKvEntry(createDefaultKvEntry(argument), defaultLastUpdateTs, SingleValueArgumentEntry.DEFAULT_VERSION)); |
|||
return transformSingleValueArgument(Optional.of(attributeKvEntry)); |
|||
}, calculatedFieldCallbackExecutor); |
|||
} |
|||
|
|||
protected ListenableFuture<ArgumentEntry> fetchTsLatest(TenantId tenantId, EntityId entityId, Argument argument, long defaultTs) { |
|||
String timeseriesKey = argument.getRefEntityKey().getKey(); |
|||
log.trace("[{}][{}] Fetching latest timeseries {}", tenantId, entityId, timeseriesKey); |
|||
return transformSingleValueArgument( |
|||
Futures.transform( |
|||
timeseriesService.findLatest(tenantId, entityId, timeseriesKey), |
|||
result -> { |
|||
log.debug("[{}][{}] Fetched latest timeseries {}: {}", tenantId, entityId, timeseriesKey, result); |
|||
return result.or(() -> Optional.of(new BasicTsKvEntry(defaultTs, createDefaultKvEntry(argument), SingleValueArgumentEntry.DEFAULT_VERSION))); |
|||
}, calculatedFieldCallbackExecutor)); |
|||
} |
|||
|
|||
private ReadTsKvQuery buildTsRollingQuery(TenantId tenantId, Argument argument, long startTs, long endTs) { |
|||
long maxDataPoints = apiLimitService.getLimit( |
|||
tenantId, DefaultTenantProfileConfiguration::getMaxDataPointsPerRollingArg); |
|||
int argumentLimit = argument.getLimit(); |
|||
int limit = argumentLimit == 0 || argumentLimit > maxDataPoints ? (int) maxDataPoints : argumentLimit; |
|||
return new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, endTs, 0, limit, Aggregation.NONE); |
|||
} |
|||
|
|||
} |
|||
@ -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,111 @@ |
|||
/** |
|||
* 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.geofencing; |
|||
|
|||
import lombok.Data; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.script.api.tbel.TbelCfArg; |
|||
import org.thingsboard.script.api.tbel.TbelCfGeofencingArg; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.kv.KvEntry; |
|||
import org.thingsboard.server.common.util.ProtoUtils; |
|||
import org.thingsboard.server.gen.transport.TransportProtos; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; |
|||
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntryType; |
|||
|
|||
import java.util.Map; |
|||
import java.util.stream.Collectors; |
|||
|
|||
@Data |
|||
@Slf4j |
|||
public class GeofencingArgumentEntry implements ArgumentEntry { |
|||
|
|||
private Map<EntityId, GeofencingZoneState> zoneStates; |
|||
|
|||
private boolean forceResetPrevious; |
|||
|
|||
public GeofencingArgumentEntry() { |
|||
} |
|||
|
|||
public GeofencingArgumentEntry(EntityId entityId, TransportProtos.AttributeValueProto entry) { |
|||
this(Map.of(entityId, ProtoUtils.fromProto(entry))); |
|||
} |
|||
|
|||
public GeofencingArgumentEntry(Map<EntityId, KvEntry> entityIdkvEntryMap) { |
|||
this.zoneStates = toZones(entityIdkvEntryMap); |
|||
} |
|||
|
|||
@Override |
|||
public ArgumentEntryType getType() { |
|||
return ArgumentEntryType.GEOFENCING; |
|||
} |
|||
|
|||
@Override |
|||
public Object getValue() { |
|||
return zoneStates; |
|||
} |
|||
|
|||
@Override |
|||
public boolean updateEntry(ArgumentEntry entry) { |
|||
if (!(entry instanceof GeofencingArgumentEntry geofencingArgumentEntry)) { |
|||
throw new IllegalArgumentException("Unsupported argument entry type for geofencing argument entry: " + entry.getType()); |
|||
} |
|||
if (geofencingArgumentEntry.isEmpty()) { |
|||
zoneStates.clear(); |
|||
return true; |
|||
} |
|||
boolean updated = false; |
|||
for (var zoneEntry : geofencingArgumentEntry.getZoneStates().entrySet()) { |
|||
if (updateZone(zoneEntry)) { |
|||
updated = true; |
|||
} |
|||
} |
|||
return updated; |
|||
} |
|||
|
|||
@Override |
|||
public boolean isEmpty() { |
|||
return zoneStates == null || zoneStates.isEmpty(); |
|||
} |
|||
|
|||
@Override |
|||
public TbelCfArg toTbelCfArg() { |
|||
return new TbelCfGeofencingArg(zoneStates); |
|||
} |
|||
|
|||
private Map<EntityId, GeofencingZoneState> toZones(Map<EntityId, KvEntry> entityIdKvEntryMap) { |
|||
return entityIdKvEntryMap.entrySet().stream() |
|||
.collect(Collectors.toMap(Map.Entry::getKey, |
|||
entry -> new GeofencingZoneState(entry.getKey(), entry.getValue()))); |
|||
} |
|||
|
|||
private boolean updateZone(Map.Entry<EntityId, GeofencingZoneState> zoneEntry) { |
|||
EntityId zoneId = zoneEntry.getKey(); |
|||
GeofencingZoneState newZoneState = zoneEntry.getValue(); |
|||
|
|||
GeofencingZoneState existingZoneState = zoneStates.get(zoneId); |
|||
if (existingZoneState == null) { |
|||
zoneStates.put(zoneId, newZoneState); |
|||
return true; |
|||
} |
|||
if (newZoneState.getPerimeterDefinition() == null) { |
|||
zoneStates.remove(zoneId); |
|||
return true; |
|||
} |
|||
return existingZoneState.update(newZoneState); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,197 @@ |
|||
/** |
|||
* 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.geofencing; |
|||
|
|||
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 com.google.common.util.concurrent.MoreExecutors; |
|||
import lombok.EqualsAndHashCode; |
|||
import lombok.Getter; |
|||
import lombok.Setter; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.common.util.geo.Coordinates; |
|||
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|||
import org.thingsboard.server.common.data.cf.configuration.OutputType; |
|||
import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingCalculatedFieldConfiguration; |
|||
import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingReportStrategy; |
|||
import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingTransitionEvent; |
|||
import org.thingsboard.server.common.data.cf.configuration.geofencing.ZoneGroupConfiguration; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.relation.EntityRelation; |
|||
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.ArgumentEntryType; |
|||
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.ArrayList; |
|||
import java.util.List; |
|||
import java.util.Map; |
|||
import java.util.stream.Collectors; |
|||
|
|||
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY; |
|||
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY; |
|||
import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingPresenceStatus.INSIDE; |
|||
import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingPresenceStatus.OUTSIDE; |
|||
|
|||
@Getter |
|||
@Setter |
|||
@Slf4j |
|||
@EqualsAndHashCode(callSuper = true) |
|||
public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { |
|||
|
|||
private long lastDynamicArgumentsRefreshTs = -1; |
|||
|
|||
public GeofencingCalculatedFieldState(EntityId entityId) { |
|||
super(entityId); |
|||
} |
|||
|
|||
@Override |
|||
public CalculatedFieldType getType() { |
|||
return CalculatedFieldType.GEOFENCING; |
|||
} |
|||
|
|||
@Override |
|||
protected void validateNewEntry(String key, ArgumentEntry newEntry) { |
|||
switch (key) { |
|||
case ENTITY_ID_LATITUDE_ARGUMENT_KEY, ENTITY_ID_LONGITUDE_ARGUMENT_KEY -> { |
|||
if (!(newEntry instanceof SingleValueArgumentEntry)) { |
|||
throw new IllegalArgumentException("Unsupported argument entry type for " + key + " argument: " + newEntry.getType() + ". " + |
|||
"Only SINGLE_VALUE type is allowed."); |
|||
} |
|||
} |
|||
default -> { |
|||
if (!(newEntry instanceof GeofencingArgumentEntry)) { |
|||
throw new IllegalArgumentException("Unsupported argument entry type for " + key + " argument: " + newEntry.getType() + ". " + |
|||
"Only GEOFENCING type is allowed."); |
|||
} |
|||
} |
|||
} |
|||
} |
|||
|
|||
@Override |
|||
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) { |
|||
double latitude = (double) arguments.get(ENTITY_ID_LATITUDE_ARGUMENT_KEY).getValue(); |
|||
double longitude = (double) arguments.get(ENTITY_ID_LONGITUDE_ARGUMENT_KEY).getValue(); |
|||
Coordinates entityCoordinates = new Coordinates(latitude, longitude); |
|||
|
|||
var geofencingCfg = (GeofencingCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); |
|||
Map<String, ZoneGroupConfiguration> zoneGroups = geofencingCfg.getZoneGroups(); |
|||
|
|||
ObjectNode valuesNode = JacksonUtil.newObjectNode(); |
|||
List<ListenableFuture<Boolean>> relationFutures = new ArrayList<>(); |
|||
|
|||
getGeofencingArguments().forEach((argumentKey, argumentEntry) -> { |
|||
ZoneGroupConfiguration zoneGroupCfg = zoneGroups.get(argumentKey); |
|||
if (zoneGroupCfg == null) { |
|||
throw new RuntimeException("Zone group configuration is missing for the: " + entityId); |
|||
} |
|||
boolean createRelationsWithMatchedZones = zoneGroupCfg.isCreateRelationsWithMatchedZones(); |
|||
List<GeofencingEvalResult> zoneResults = new ArrayList<>(argumentEntry.getZoneStates().size()); |
|||
argumentEntry.getZoneStates().forEach((zoneId, zoneState) -> { |
|||
GeofencingEvalResult eval = zoneState.evaluate(entityCoordinates); |
|||
zoneResults.add(eval); |
|||
if (createRelationsWithMatchedZones) { |
|||
GeofencingTransitionEvent transitionEvent = eval.transition(); |
|||
if (transitionEvent == null) { |
|||
return; |
|||
} |
|||
EntityRelation relation = switch (zoneGroupCfg.getDirection()) { |
|||
case TO -> new EntityRelation(zoneId, entityId, zoneGroupCfg.getRelationType()); |
|||
case FROM -> new EntityRelation(entityId, zoneId, zoneGroupCfg.getRelationType()); |
|||
}; |
|||
ListenableFuture<Boolean> f = switch (transitionEvent) { |
|||
case ENTERED -> ctx.getRelationService().saveRelationAsync(ctx.getTenantId(), relation); |
|||
case LEFT -> ctx.getRelationService().deleteRelationAsync(ctx.getTenantId(), relation); |
|||
}; |
|||
relationFutures.add(f); |
|||
} |
|||
}); |
|||
updateValuesNode(argumentKey, zoneResults, zoneGroupCfg.getReportStrategy(), valuesNode); |
|||
}); |
|||
|
|||
OutputType outputType = ctx.getOutput().getType(); |
|||
var result = TelemetryCalculatedFieldResult.builder() |
|||
.type(outputType) |
|||
.scope(ctx.getOutput().getScope()) |
|||
.result(toResultNode(outputType, valuesNode)) |
|||
.build(); |
|||
if (relationFutures.isEmpty()) { |
|||
return Futures.immediateFuture(result); |
|||
} |
|||
return Futures.whenAllComplete(relationFutures).call(() -> result, MoreExecutors.directExecutor()); |
|||
} |
|||
|
|||
@Override |
|||
public void reset() { |
|||
super.reset(); |
|||
lastDynamicArgumentsRefreshTs = -1; |
|||
} |
|||
|
|||
public void updateLastDynamicArgumentsRefreshTs() { |
|||
lastDynamicArgumentsRefreshTs = System.currentTimeMillis(); |
|||
} |
|||
|
|||
private Map<String, GeofencingArgumentEntry> getGeofencingArguments() { |
|||
return arguments.entrySet() |
|||
.stream() |
|||
.filter(entry -> entry.getValue().getType().equals(ArgumentEntryType.GEOFENCING)) |
|||
.collect(Collectors.toMap(Map.Entry::getKey, entry -> (GeofencingArgumentEntry) entry.getValue())); |
|||
} |
|||
|
|||
private void updateValuesNode(String argumentKey, List<GeofencingEvalResult> zoneResults, GeofencingReportStrategy geofencingReportStrategy, ObjectNode resultNode) { |
|||
GeofencingEvalResult aggregationResult = aggregateZoneGroup(zoneResults); |
|||
final String eventKey = argumentKey + "Event"; |
|||
final String statusKey = argumentKey + "Status"; |
|||
switch (geofencingReportStrategy) { |
|||
case REPORT_TRANSITION_EVENTS_ONLY -> addTransitionEventIfExists(resultNode, aggregationResult, eventKey); |
|||
case REPORT_PRESENCE_STATUS_ONLY -> resultNode.put(statusKey, aggregationResult.status().name()); |
|||
case REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS -> { |
|||
addTransitionEventIfExists(resultNode, aggregationResult, eventKey); |
|||
resultNode.put(statusKey, aggregationResult.status().name()); |
|||
} |
|||
} |
|||
} |
|||
|
|||
private JsonNode toResultNode(OutputType outputType, ObjectNode valuesNode) { |
|||
return toSimpleResult(outputType == OutputType.TIME_SERIES, valuesNode); |
|||
} |
|||
|
|||
private GeofencingEvalResult aggregateZoneGroup(List<GeofencingEvalResult> zoneResults) { |
|||
boolean nowInside = zoneResults.stream().anyMatch(r -> INSIDE.equals(r.status())); |
|||
boolean prevInside = zoneResults.stream() |
|||
.anyMatch(r -> GeofencingTransitionEvent.LEFT.equals(r.transition()) || r.transition() == null && r.status() == INSIDE); |
|||
GeofencingTransitionEvent transition = null; |
|||
if (!prevInside && nowInside) { |
|||
transition = GeofencingTransitionEvent.ENTERED; |
|||
} else if (prevInside && !nowInside) { |
|||
transition = GeofencingTransitionEvent.LEFT; |
|||
} |
|||
return new GeofencingEvalResult(transition, nowInside ? INSIDE : OUTSIDE); |
|||
} |
|||
|
|||
private void addTransitionEventIfExists(ObjectNode resultNode, GeofencingEvalResult aggregationResult, String eventKey) { |
|||
if (aggregationResult.transition() != null) { |
|||
resultNode.put(eventKey, aggregationResult.transition().name()); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,24 @@ |
|||
/** |
|||
* 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.geofencing; |
|||
|
|||
import jakarta.annotation.Nullable; |
|||
import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingPresenceStatus; |
|||
import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingTransitionEvent; |
|||
|
|||
public record GeofencingEvalResult(@Nullable GeofencingTransitionEvent transition, |
|||
GeofencingPresenceStatus status) { |
|||
} |
|||
@ -0,0 +1,106 @@ |
|||
/** |
|||
* 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.geofencing; |
|||
|
|||
import lombok.Data; |
|||
import lombok.EqualsAndHashCode; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.common.util.geo.Coordinates; |
|||
import org.thingsboard.common.util.geo.PerimeterDefinition; |
|||
import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingPresenceStatus; |
|||
import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingTransitionEvent; |
|||
import org.thingsboard.server.common.data.id.EntityId; |
|||
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|||
import org.thingsboard.server.common.data.kv.KvEntry; |
|||
import org.thingsboard.server.common.util.ProtoUtils; |
|||
import org.thingsboard.server.gen.transport.TransportProtos.GeofencingZoneProto; |
|||
|
|||
import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingPresenceStatus.INSIDE; |
|||
import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingPresenceStatus.OUTSIDE; |
|||
|
|||
@Data |
|||
public class GeofencingZoneState { |
|||
|
|||
private final EntityId zoneId; |
|||
|
|||
private long ts; |
|||
private Long version; |
|||
private PerimeterDefinition perimeterDefinition; |
|||
|
|||
@EqualsAndHashCode.Exclude |
|||
private GeofencingPresenceStatus lastPresence; |
|||
|
|||
public GeofencingZoneState(EntityId zoneId, KvEntry entry) { |
|||
this.zoneId = zoneId; |
|||
if (!(entry instanceof AttributeKvEntry attributeKvEntry)) { |
|||
throw new IllegalArgumentException("Unsupported KvEntry type for geofencing zone state: " + entry.getClass().getSimpleName()); |
|||
} |
|||
this.ts = attributeKvEntry.getLastUpdateTs(); |
|||
this.version = attributeKvEntry.getVersion(); |
|||
this.perimeterDefinition = JacksonUtil.fromString(entry.getValueAsString(), PerimeterDefinition.class); |
|||
} |
|||
|
|||
public GeofencingZoneState(GeofencingZoneProto proto) { |
|||
this.zoneId = ProtoUtils.fromProto(proto.getZoneId()); |
|||
this.ts = proto.getTs(); |
|||
this.version = proto.getVersion(); |
|||
this.perimeterDefinition = JacksonUtil.fromString(proto.getPerimeterDefinition(), PerimeterDefinition.class); |
|||
if (proto.hasInside()) { |
|||
this.lastPresence = proto.getInside() ? INSIDE : OUTSIDE; |
|||
} |
|||
} |
|||
|
|||
public boolean update(GeofencingZoneState newZoneState) { |
|||
if (newZoneState.getTs() <= this.ts) { |
|||
return false; |
|||
} |
|||
Long newVersion = newZoneState.getVersion(); |
|||
if (newVersion == null || this.version == null || newVersion > this.version) { |
|||
this.ts = newZoneState.getTs(); |
|||
this.version = newVersion; |
|||
this.perimeterDefinition = newZoneState.getPerimeterDefinition(); |
|||
this.lastPresence = null; |
|||
return true; |
|||
} |
|||
return false; |
|||
} |
|||
|
|||
public GeofencingEvalResult evaluate(Coordinates entityCoordinates) { |
|||
boolean nowInside = perimeterDefinition.checkMatches(entityCoordinates); |
|||
|
|||
GeofencingPresenceStatus status = nowInside ? INSIDE : OUTSIDE; |
|||
|
|||
// first evaluation
|
|||
if (this.lastPresence == null) { |
|||
this.lastPresence = status; |
|||
GeofencingTransitionEvent transition = null; |
|||
if (status == GeofencingPresenceStatus.INSIDE) { |
|||
transition = GeofencingTransitionEvent.ENTERED; |
|||
} |
|||
return new GeofencingEvalResult(transition, status); |
|||
} |
|||
// State changed
|
|||
if (this.lastPresence != status) { |
|||
this.lastPresence = status; |
|||
GeofencingTransitionEvent transition = (status == GeofencingPresenceStatus.INSIDE) ? |
|||
GeofencingTransitionEvent.ENTERED : GeofencingTransitionEvent.LEFT; |
|||
return new GeofencingEvalResult(transition, status); |
|||
} |
|||
// State unchanged
|
|||
return new GeofencingEvalResult(null, status); |
|||
} |
|||
|
|||
} |
|||
@ -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(); |
|||
} |
|||
|
|||
} |
|||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue