Browse Source

geofencing cf init commit

pull/13920/head
dshvaika 1 year ago
parent
commit
57146ff94f
  1. 27
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  2. 10
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java
  3. 92
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
  4. 9
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntryType.java
  6. 6
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  7. 3
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java
  8. 85
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingArgumentEntry.java
  9. 243
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldState.java
  10. 4
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java
  11. 4
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java
  12. 4
      application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java
  13. 2
      common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldType.java
  14. 2
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Argument.java
  15. 22
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CFArgumentDynamicSourceType.java
  16. 3
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CalculatedFieldConfiguration.java
  17. 39
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CfArgumentDynamicSourceConfiguration.java
  18. 31
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/GeofencingCalculatedFieldConfiguration.java
  19. 57
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfiguration.java
  20. 40
      common/util/src/main/java/org/thingsboard/common/util/geo/CirclePerimeterDefinition.java
  21. 39
      common/util/src/main/java/org/thingsboard/common/util/geo/PerimeterDefinition.java
  22. 35
      common/util/src/main/java/org/thingsboard/common/util/geo/PolygonPerimeterDefinition.java

27
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java

@ -288,17 +288,34 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
boolean stateSizeChecked = false;
try {
if (ctx.isInitialized() && state.isReady()) {
CalculatedFieldResult calculationResult = state.performCalculation(ctx).get(systemContext.getCfCalculationResultTimeout(), TimeUnit.SECONDS);
List<CalculatedFieldResult> calculationResults = state.performCalculation(ctx).get(systemContext.getCfCalculationResultTimeout(), TimeUnit.SECONDS);
state.checkStateSize(ctxId, ctx.getMaxStateSize());
stateSizeChecked = true;
if (state.isSizeOk()) {
if (!calculationResult.isEmpty()) {
cfService.pushMsgToRuleEngine(tenantId, entityId, calculationResult, cfIdList, callback);
} else {
if (calculationResults.isEmpty()) {
callback.onSuccess();
} else {
TbCallback effectiveCallback = calculationResults.size() > 1 ?
new MultipleTbCallback(calculationResults.size(), callback) : callback;
for (CalculatedFieldResult calculationResult : calculationResults) {
if (calculationResult.isEmpty()) {
effectiveCallback.onSuccess();
} else {
cfService.pushMsgToRuleEngine(tenantId, entityId, calculationResult, cfIdList, callback);
}
}
}
if (DebugModeUtil.isDebugAllAvailable(ctx.getCalculatedField())) {
systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId, state.getArguments(), tbMsgId, tbMsgType, calculationResult.getResult().toString(), null);
if (calculationResults.isEmpty()) {
systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId,
state.getArguments(), tbMsgId, tbMsgType, null, null);
} else {
for (CalculatedFieldResult calculationResult : calculationResults) {
systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId,
state.getArguments(), tbMsgId, tbMsgType, calculationResult.getResultAsString(), null);
}
}
}
}
} else {

10
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java

@ -34,4 +34,14 @@ public final class CalculatedFieldResult {
(result.isTextual() && result.asText().isEmpty());
}
public String getResultAsString() {
if (result == null) {
return null;
}
if (result.isTextual()) {
return result.asText();
}
return result.toString();
}
}

92
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java

@ -32,12 +32,16 @@ import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
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.OutputType;
import org.thingsboard.server.common.data.cf.configuration.RelationQueryDynamicSourceConfiguration;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
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;
@ -55,6 +59,7 @@ import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
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.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsgProto;
@ -70,6 +75,7 @@ import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
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.CalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.GeofencingCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
@ -86,6 +92,10 @@ import java.util.concurrent.ExecutionException;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.DataConstants.SCOPE;
import static org.thingsboard.server.service.cf.ctx.state.GeofencingCalculatedFieldState.ENTITY_ID_LATITUDE_ARGUMENT_KEY;
import static org.thingsboard.server.service.cf.ctx.state.GeofencingCalculatedFieldState.ENTITY_ID_LONGITUDE_ARGUMENT_KEY;
import static org.thingsboard.server.service.cf.ctx.state.GeofencingCalculatedFieldState.RESTRICTED_ZONES_ARGUMENT_KEY;
import static org.thingsboard.server.service.cf.ctx.state.GeofencingCalculatedFieldState.SAVE_ZONES_ARGUMENT_KEY;
import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto;
@TbRuleEngineComponent
@ -99,6 +109,7 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP
private final TbClusterService clusterService;
private final ApiLimitService apiLimitService;
private final PartitionService partitionService;
private final RelationService relationService;
private ListeningExecutorService calculatedFieldCallbackExecutor;
@ -118,11 +129,29 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP
@Override
public ListenableFuture<CalculatedFieldState> fetchStateFromDb(CalculatedFieldCtx ctx, EntityId entityId) {
Map<String, ListenableFuture<ArgumentEntry>> argFutures = new HashMap<>();
for (var entry : ctx.getArguments().entrySet()) {
var argEntityId = entry.getValue().getRefEntityId() != null ? entry.getValue().getRefEntityId() : entityId;
var argValueFuture = fetchKvEntry(ctx.getTenantId(), argEntityId, entry.getValue());
argFutures.put(entry.getKey(), argValueFuture);
if (ctx.getCalculatedField().getType().equals(CalculatedFieldType.GEOFENCING)) {
// Ignoring any other arguments except ENTITY_ID_LATITUDE_ARGUMENT_KEY,
// ENTITY_ID_LONGITUDE_ARGUMENT_KEY, SAVE_ZONES_ARGUMENT_KEY, RESTRICTED_ZONES_ARGUMENT_KEY.
for (var entry : ctx.getArguments().entrySet()) {
switch (entry.getKey()) {
case ENTITY_ID_LATITUDE_ARGUMENT_KEY, ENTITY_ID_LONGITUDE_ARGUMENT_KEY ->
argFutures.put(entry.getKey(), fetchKvEntry(ctx.getTenantId(), resolveEntityId(entityId, entry), entry.getValue()));
case SAVE_ZONES_ARGUMENT_KEY, RESTRICTED_ZONES_ARGUMENT_KEY -> {
var resolvedEntityIdsFuture = resolveGeofencingEntityIds(ctx.getTenantId(), entityId, entry);
argFutures.put(entry.getKey(), Futures.transformAsync(resolvedEntityIdsFuture, resolvedEntityIds ->
fetchGeofencingKvEntry(ctx.getTenantId(), resolvedEntityIds, entry.getValue()), MoreExecutors.directExecutor()));
}
}
}
} else {
for (var entry : ctx.getArguments().entrySet()) {
var argEntityId = resolveEntityId(entityId, entry);
var argValueFuture = fetchKvEntry(ctx.getTenantId(), argEntityId, entry.getValue());
argFutures.put(entry.getKey(), argValueFuture);
}
}
return Futures.whenAllComplete(argFutures.values()).call(() -> {
var result = createStateByType(ctx);
result.updateState(ctx, argFutures.entrySet().stream()
@ -145,7 +174,7 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP
public Map<String, ArgumentEntry> fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map<String, Argument> arguments) {
Map<String, ListenableFuture<ArgumentEntry>> argFutures = new HashMap<>();
for (var entry : arguments.entrySet()) {
var argEntityId = entry.getValue().getRefEntityId() != null ? entry.getValue().getRefEntityId() : entityId;
var argEntityId = resolveEntityId(entityId, entry);
var argValueFuture = fetchKvEntry(tenantId, argEntityId, entry.getValue());
argFutures.put(entry.getKey(), argValueFuture);
}
@ -241,6 +270,58 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP
return builder.build();
}
private EntityId resolveEntityId(EntityId entityId, Entry<String, Argument> entry) {
return entry.getValue().getRefEntityId() != null ? entry.getValue().getRefEntityId() : entityId;
}
private ListenableFuture<List<EntityId>> resolveGeofencingEntityIds(TenantId tenantId, EntityId entityId, Entry<String, Argument> entry) {
Argument value = entry.getValue();
if (value.getRefEntityId() != null) {
return Futures.immediateFuture(List.of(value.getRefEntityId()));
}
var refDynamicSource = value.getRefDynamicSource();
if (refDynamicSource == null) {
return Futures.immediateFuture(List.of(entityId));
}
return switch (value.getRefDynamicSource()) {
case RELATION_QUERY -> {
var relationQueryDynamicSourceConfiguration = (RelationQueryDynamicSourceConfiguration) value.getRefDynamicSourceConfiguration();
yield Futures.transform(relationService.findByQuery(tenantId, relationQueryDynamicSourceConfiguration.toEntityRelationsQuery(entityId)),
relationQueryDynamicSourceConfiguration::resolveEntityIds, MoreExecutors.directExecutor());
}
};
}
private ListenableFuture<ArgumentEntry> fetchGeofencingKvEntry(TenantId tenantId, List<EntityId> geofencingEntities, Argument argument) {
// TODO: Should we handle any other case?
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(() ->
new BaseAttributeKvEntry(createDefaultKvEntry(argument), System.currentTimeMillis(), 0L))),
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))),
calculatedFieldCallbackExecutor
);
}
private ListenableFuture<ArgumentEntry> fetchKvEntry(TenantId tenantId, EntityId entityId, Argument argument) {
return switch (argument.getRefEntityKey().getType()) {
case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument);
@ -301,6 +382,7 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP
return switch (ctx.getCfType()) {
case SIMPLE -> new SimpleCalculatedFieldState(ctx.getArgNames());
case SCRIPT -> new ScriptCalculatedFieldState(ctx.getArgNames());
case GEOFENCING -> new GeofencingCalculatedFieldState(ctx.getArgNames());
};
}

9
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java

@ -19,10 +19,12 @@ import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import org.thingsboard.script.api.tbel.TbelCfArg;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import java.util.List;
import java.util.Map;
@JsonTypeInfo(
use = JsonTypeInfo.Id.NAME,
@ -31,7 +33,8 @@ import java.util.List;
)
@JsonSubTypes({
@JsonSubTypes.Type(value = SingleValueArgumentEntry.class, name = "SINGLE_VALUE"),
@JsonSubTypes.Type(value = TsRollingArgumentEntry.class, name = "TS_ROLLING")
@JsonSubTypes.Type(value = TsRollingArgumentEntry.class, name = "TS_ROLLING"),
@JsonSubTypes.Type(value = GeofencingArgumentEntry.class, name = "GEOFENCING")
})
public interface ArgumentEntry {
@ -58,4 +61,8 @@ public interface ArgumentEntry {
return new TsRollingArgumentEntry(kvEntries, limit, timeWindow);
}
static ArgumentEntry createGeofencingValueArgument(Map<EntityId, KvEntry> entityIdkvEntryMap) {
return new GeofencingArgumentEntry(entityIdkvEntryMap);
}
}

2
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntryType.java

@ -16,5 +16,5 @@
package org.thingsboard.server.service.cf.ctx.state;
public enum ArgumentEntryType {
SINGLE_VALUE, TS_ROLLING
SINGLE_VALUE, TS_ROLLING, GEOFENCING
}

6
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java

@ -59,6 +59,7 @@ public class CalculatedFieldCtx {
private final Map<String, Argument> arguments;
private final Map<ReferencedEntityKey, String> mainEntityArguments;
private final Map<EntityId, Map<ReferencedEntityKey, String>> linkedEntityArguments;
private final Map<ReferencedEntityKey, String> dynamicEntityArguments;
private final List<String> argNames;
private Output output;
private String expression;
@ -84,10 +85,13 @@ public class CalculatedFieldCtx {
this.arguments = configuration.getArguments();
this.mainEntityArguments = new HashMap<>();
this.linkedEntityArguments = new HashMap<>();
this.dynamicEntityArguments = new HashMap<>();
for (Map.Entry<String, Argument> entry : arguments.entrySet()) {
var refId = entry.getValue().getRefEntityId();
var refKey = entry.getValue().getRefEntityKey();
if (refId == null || refId.equals(calculatedField.getEntityId())) {
if (refId == null && entry.getValue().getRefDynamicSource() != null) {
dynamicEntityArguments.put(refKey, entry.getKey());
} else if (refId == null || refId.equals(calculatedField.getEntityId())) {
mainEntityArguments.put(refKey, entry.getKey());
} else {
linkedEntityArguments.computeIfAbsent(refId, key -> new HashMap<>()).put(refKey, entry.getKey());

3
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java

@ -34,6 +34,7 @@ import java.util.Map;
@JsonSubTypes({
@JsonSubTypes.Type(value = SimpleCalculatedFieldState.class, name = "SIMPLE"),
@JsonSubTypes.Type(value = ScriptCalculatedFieldState.class, name = "SCRIPT"),
@JsonSubTypes.Type(value = GeofencingCalculatedFieldState.class, name = "GEOFENCING"),
})
public interface CalculatedFieldState {
@ -48,7 +49,7 @@ public interface CalculatedFieldState {
boolean updateState(CalculatedFieldCtx ctx, Map<String, ArgumentEntry> argumentValues);
ListenableFuture<CalculatedFieldResult> performCalculation(CalculatedFieldCtx ctx);
ListenableFuture<List<CalculatedFieldResult>> performCalculation(CalculatedFieldCtx ctx);
@JsonIgnore
boolean isReady();

85
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingArgumentEntry.java

@ -0,0 +1,85 @@
/**
* 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;
import lombok.Data;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.geo.PerimeterDefinition;
import org.thingsboard.script.api.tbel.TbelCfArg;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.KvEntry;
import java.util.Map;
import java.util.Objects;
import java.util.stream.Collectors;
// TODO: implement
@Data
public class GeofencingArgumentEntry implements ArgumentEntry {
private Map<EntityId, PerimeterDefinition> geofencingIdToPerimeter;
private boolean forceResetPrevious;
public GeofencingArgumentEntry(Map<EntityId, KvEntry> entityIdKvEntryMap) {
this.geofencingIdToPerimeter = toPerimetersMap(entityIdKvEntryMap);
}
@Override
public ArgumentEntryType getType() {
return ArgumentEntryType.GEOFENCING;
}
@Override
public Object getValue() {
return geofencingIdToPerimeter;
}
@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 (Objects.equals(this.geofencingIdToPerimeter, geofencingArgumentEntry.getGeofencingIdToPerimeter())) {
return false; // No change
}
this.geofencingIdToPerimeter = geofencingArgumentEntry.getGeofencingIdToPerimeter();
return true;
}
@Override
public boolean isEmpty() {
return geofencingIdToPerimeter == null || geofencingIdToPerimeter.isEmpty();
}
@Override
public TbelCfArg toTbelCfArg() {
return null;
}
private Map<EntityId, PerimeterDefinition> toPerimetersMap(Map<EntityId, KvEntry> entityIdKvEntryMap) {
return entityIdKvEntryMap.entrySet().stream().map(entry -> {
if (entry.getValue().getJsonValue().isEmpty()) {
return null;
}
String rawPerimeterValue = entry.getValue().getJsonValue().get();
PerimeterDefinition perimeter = JacksonUtil.fromString(rawPerimeterValue, PerimeterDefinition.class);
return Map.entry(entry.getKey(), perimeter);
})
.filter(Objects::nonNull)
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
}
}

243
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldState.java

@ -0,0 +1,243 @@
/**
* 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;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import lombok.Data;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.geo.Coordinates;
import org.thingsboard.common.util.geo.PerimeterDefinition;
import org.thingsboard.rule.engine.geo.EntityGeofencingState;
import org.thingsboard.rule.engine.util.GpsGeofencingEvents;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.service.cf.CalculatedFieldResult;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
@Data
public class GeofencingCalculatedFieldState implements CalculatedFieldState {
public static final String ENTITY_ID_LATITUDE_ARGUMENT_KEY = "latitude";
public static final String ENTITY_ID_LONGITUDE_ARGUMENT_KEY = "longitude";
public static final String SAVE_ZONES_ARGUMENT_KEY = "saveZones";
public static final String RESTRICTED_ZONES_ARGUMENT_KEY = "restrictedZones";
private List<String> requiredArguments;
private Map<String, ArgumentEntry> arguments;
private long latestTimestamp = -1;
private Map<EntityId, EntityGeofencingState> saveZoneStates;
private Map<EntityId, EntityGeofencingState> restrictedZoneStates;
public GeofencingCalculatedFieldState() {
this(List.of(ENTITY_ID_LATITUDE_ARGUMENT_KEY, ENTITY_ID_LONGITUDE_ARGUMENT_KEY, SAVE_ZONES_ARGUMENT_KEY, RESTRICTED_ZONES_ARGUMENT_KEY));
}
public GeofencingCalculatedFieldState(List<String> argNames) {
this.requiredArguments = argNames;
this.arguments = new HashMap<>();
this.saveZoneStates = new HashMap<>();
this.restrictedZoneStates = new HashMap<>();
}
@Override
public CalculatedFieldType getType() {
return CalculatedFieldType.GEOFENCING;
}
@Override
public boolean updateState(CalculatedFieldCtx ctx, Map<String, ArgumentEntry> argumentValues) {
// TODO: Do I need to check argument for null?
if (arguments == null) {
arguments = new HashMap<>();
}
boolean stateUpdated = false;
for (Map.Entry<String, ArgumentEntry> entry : argumentValues.entrySet()) {
String key = entry.getKey();
ArgumentEntry newEntry = entry.getValue();
// TODO: Do I need to check argument size?
// checkArgumentSize(key, newEntry, ctx);
ArgumentEntry existingEntry = arguments.get(key);
boolean entryUpdated;
// TODO: What is force reset previos?
// if (existingEntry == null || newEntry.isForceResetPrevious()) {
// fresh start of state. No entry exists yet.
if (existingEntry == null) {
switch (key) {
case ENTITY_ID_LATITUDE_ARGUMENT_KEY:
case ENTITY_ID_LONGITUDE_ARGUMENT_KEY:
if (!(newEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry)) {
throw new IllegalArgumentException(key + " argument must be a single value argument.");
}
arguments.put(key, singleValueArgumentEntry);
entryUpdated = true;
break;
case SAVE_ZONES_ARGUMENT_KEY:
case RESTRICTED_ZONES_ARGUMENT_KEY:
if (!(newEntry instanceof GeofencingArgumentEntry geofencingArgumentEntry)) {
throw new IllegalArgumentException(key + " argument must be a geofencing argument entry.");
}
arguments.put(key, geofencingArgumentEntry);
entryUpdated = true;
break;
default:
throw new IllegalArgumentException("Unsupported argument: " + key);
}
} else {
entryUpdated = switch (key) {
case ENTITY_ID_LATITUDE_ARGUMENT_KEY,
ENTITY_ID_LONGITUDE_ARGUMENT_KEY -> existingEntry.updateEntry(newEntry);
case SAVE_ZONES_ARGUMENT_KEY,
RESTRICTED_ZONES_ARGUMENT_KEY -> {
// TODO: ensure zone cleanup working correctly.
boolean updated = existingEntry.updateEntry(newEntry);
if (updated) {
Map<EntityId, EntityGeofencingState> currentStates =
key.equals(SAVE_ZONES_ARGUMENT_KEY) ? saveZoneStates : restrictedZoneStates;
Set<EntityId> newZoneIds = ((GeofencingArgumentEntry) newEntry).getGeofencingIdToPerimeter().keySet();
currentStates.keySet().removeIf(existingZoneId -> !newZoneIds.contains(existingZoneId));
}
yield updated;
}
default -> throw new IllegalStateException("Unsupported argument: " + key);
};
}
if (entryUpdated) {
stateUpdated = true;
updateLastUpdateTimestamp(newEntry);
}
}
return stateUpdated;
}
@Override
public ListenableFuture<List<CalculatedFieldResult>> performCalculation(CalculatedFieldCtx ctx) {
List<CalculatedFieldResult> savedZonesStatesResults = updateSavedGeofencingZonesState(ctx);
List<CalculatedFieldResult> restrictedZonesStatesResults = updateRestrictedGeofencingZonesState(ctx);
List<CalculatedFieldResult> allZoneStatesResults =
new ArrayList<>(savedZonesStatesResults.size() + restrictedZonesStatesResults.size());
allZoneStatesResults.addAll(savedZonesStatesResults);
allZoneStatesResults.addAll(restrictedZonesStatesResults);
return Futures.immediateFuture(allZoneStatesResults);
}
@Override
public boolean isReady() {
return arguments.keySet().containsAll(requiredArguments) &&
arguments.values().stream().noneMatch(ArgumentEntry::isEmpty);
}
// TODO: implement
@Override
public boolean isSizeExceedsLimit() {
return false;
}
// TODO: implement
@Override
public void checkStateSize(CalculatedFieldEntityCtxId ctxId, long maxStateSize) {
}
// TODO: implement
@Override
public void checkArgumentSize(String name, ArgumentEntry entry, CalculatedFieldCtx ctx) {
}
private void updateLastUpdateTimestamp(ArgumentEntry entry) {
long newTs = this.latestTimestamp;
if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) {
newTs = singleValueArgumentEntry.getTs();
}
this.latestTimestamp = Math.max(this.latestTimestamp, newTs);
}
private List<CalculatedFieldResult> updateSavedGeofencingZonesState(CalculatedFieldCtx ctx) {
return updateGeofencingZonesState(ctx, saveZoneStates, false);
}
private List<CalculatedFieldResult> updateRestrictedGeofencingZonesState(CalculatedFieldCtx ctx) {
return updateGeofencingZonesState(ctx, restrictedZoneStates, true);
}
// TODO: Ensure all cases are covered based on rule node logic.
private List<CalculatedFieldResult> updateGeofencingZonesState(CalculatedFieldCtx ctx, Map<EntityId, EntityGeofencingState> zoneStates, boolean restricted) {
var results = new ArrayList<CalculatedFieldResult>();
long stateSwitchTime = System.currentTimeMillis();
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);
String zoneKey = restricted ? RESTRICTED_ZONES_ARGUMENT_KEY : SAVE_ZONES_ARGUMENT_KEY;
GeofencingArgumentEntry zonesEntry = (GeofencingArgumentEntry) arguments.get(zoneKey);
for (Map.Entry<EntityId, PerimeterDefinition> entry : zonesEntry.getGeofencingIdToPerimeter().entrySet()) {
EntityId zoneId = entry.getKey();
PerimeterDefinition perimeter = entry.getValue();
boolean inside = perimeter.checkMatches(entityCoordinates);
// Always present or created
EntityGeofencingState state = zoneStates.computeIfAbsent(
zoneId, id -> new EntityGeofencingState(false, 0L, false)
);
String event;
if (state.getStateSwitchTime() == 0L || state.isInside() != inside) {
// First state or transition (entered/left)
state.setInside(inside);
state.setStateSwitchTime(stateSwitchTime);
state.setStayed(false);
event = inside ? GpsGeofencingEvents.ENTERED : GpsGeofencingEvents.LEFT;
} else {
// No transition
event = inside ? GpsGeofencingEvents.INSIDE : GpsGeofencingEvents.OUTSIDE;
}
ObjectNode stateNode = JacksonUtil.newObjectNode();
stateNode.put("entityId", ctx.getEntityId().toString());
stateNode.put("zoneId", zoneId.getId().toString());
stateNode.put("restricted", restricted);
stateNode.put("event", event);
results.add(new CalculatedFieldResult(ctx.getOutput().getType(), ctx.getOutput().getScope(), stateNode));
}
return results;
}
}

4
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java

@ -53,7 +53,7 @@ public class ScriptCalculatedFieldState extends BaseCalculatedFieldState {
}
@Override
public ListenableFuture<CalculatedFieldResult> performCalculation(CalculatedFieldCtx ctx) {
public ListenableFuture<List<CalculatedFieldResult>> performCalculation(CalculatedFieldCtx ctx) {
Map<String, TbelCfArg> arguments = new LinkedHashMap<>();
List<Object> args = new ArrayList<>(ctx.getArgNames().size() + 1);
args.add(new Object()); // first element is a ctx, but we will set it later;
@ -70,7 +70,7 @@ public class ScriptCalculatedFieldState extends BaseCalculatedFieldState {
ListenableFuture<JsonNode> resultFuture = ctx.getCalculatedFieldScriptEngine().executeJsonAsync(args.toArray());
Output output = ctx.getOutput();
return Futures.transform(resultFuture,
result -> new CalculatedFieldResult(output.getType(), output.getScope(), result),
result -> List.of(new CalculatedFieldResult(output.getType(), output.getScope(), result)),
MoreExecutors.directExecutor()
);
}

4
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java

@ -52,7 +52,7 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState {
}
@Override
public ListenableFuture<CalculatedFieldResult> performCalculation(CalculatedFieldCtx ctx) {
public ListenableFuture<List<CalculatedFieldResult>> performCalculation(CalculatedFieldCtx ctx) {
var expr = ctx.getCustomExpression().get();
for (Map.Entry<String, ArgumentEntry> entry : this.arguments.entrySet()) {
@ -76,7 +76,7 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState {
Object result = formatResult(expressionResult, output.getDecimalsByDefault());
JsonNode outputResult = createResultJson(ctx.isUseLatestTs(), output.getName(), result);
return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), outputResult));
return Futures.immediateFuture(List.of(new CalculatedFieldResult(output.getType(), output.getScope(), outputResult)));
}
private Object formatResult(double expressionResult, Integer decimals) {

4
application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java

@ -32,6 +32,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.TsRollingArgumentPro
import org.thingsboard.server.gen.transport.TransportProtos.TsValueProto;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.GeofencingCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
@ -118,8 +119,11 @@ public class CalculatedFieldUtils {
CalculatedFieldState state = switch (type) {
case SIMPLE -> new SimpleCalculatedFieldState();
case SCRIPT -> new ScriptCalculatedFieldState();
case GEOFENCING -> new GeofencingCalculatedFieldState();
};
// TODO: add logic to restore geofencing state from proto
proto.getSingleValueArgumentsList().forEach(argProto ->
state.getArguments().put(argProto.getArgName(), fromSingleValueArgumentProto(argProto)));

2
common/data/src/main/java/org/thingsboard/server/common/data/cf/CalculatedFieldType.java

@ -17,6 +17,6 @@ package org.thingsboard.server.common.data.cf;
public enum CalculatedFieldType {
SIMPLE, SCRIPT
SIMPLE, SCRIPT, GEOFENCING
}

2
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/Argument.java

@ -26,6 +26,8 @@ public class Argument {
@Nullable
private EntityId refEntityId;
private CFArgumentDynamicSourceType refDynamicSource;
private CfArgumentDynamicSourceConfiguration refDynamicSourceConfiguration;
private ReferencedEntityKey refEntityKey;
private String defaultValue;

22
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CFArgumentDynamicSourceType.java

@ -0,0 +1,22 @@
/**
* 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.common.data.cf.configuration;
public enum CFArgumentDynamicSourceType {
RELATION_QUERY
}

3
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CalculatedFieldConfiguration.java

@ -35,7 +35,8 @@ import java.util.Map;
)
@JsonSubTypes({
@JsonSubTypes.Type(value = SimpleCalculatedFieldConfiguration.class, name = "SIMPLE"),
@JsonSubTypes.Type(value = ScriptCalculatedFieldConfiguration.class, name = "SCRIPT")
@JsonSubTypes.Type(value = ScriptCalculatedFieldConfiguration.class, name = "SCRIPT"),
@JsonSubTypes.Type(value = GeofencingCalculatedFieldConfiguration.class, name = "GEOFENCING")
})
@JsonIgnoreProperties(ignoreUnknown = true)
public interface CalculatedFieldConfiguration {

39
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CfArgumentDynamicSourceConfiguration.java

@ -0,0 +1,39 @@
/**
* 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.common.data.cf.configuration;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
@JsonTypeInfo(
use = JsonTypeInfo.Id.NAME,
include = JsonTypeInfo.As.PROPERTY,
property = "type"
)
@JsonSubTypes({
@JsonSubTypes.Type(value = RelationQueryDynamicSourceConfiguration.class, name = "RELATION_QUERY"),
})
@JsonIgnoreProperties(ignoreUnknown = true)
public interface CfArgumentDynamicSourceConfiguration {
@JsonIgnore
CFArgumentDynamicSourceType getType();
default void validate() {}
}

31
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/GeofencingCalculatedFieldConfiguration.java

@ -0,0 +1,31 @@
/**
* 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.common.data.cf.configuration;
import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
@Data
@EqualsAndHashCode(callSuper = true)
public class GeofencingCalculatedFieldConfiguration extends BaseCalculatedFieldConfiguration implements CalculatedFieldConfiguration {
@Override
public CalculatedFieldType getType() {
return CalculatedFieldType.GEOFENCING;
}
}

57
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfiguration.java

@ -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.common.data.cf.configuration;
import lombok.Data;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntityRelationsQuery;
import org.thingsboard.server.common.data.relation.EntitySearchDirection;
import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter;
import org.thingsboard.server.common.data.relation.RelationsSearchParameters;
import java.util.Collections;
import java.util.List;
@Data
public class RelationQueryDynamicSourceConfiguration implements CfArgumentDynamicSourceConfiguration {
private int maxLevel;
private EntitySearchDirection direction;
private String relationType;
private List<EntityType> profiles;
@Override
public CFArgumentDynamicSourceType getType() {
return CFArgumentDynamicSourceType.RELATION_QUERY;
}
public EntityRelationsQuery toEntityRelationsQuery(EntityId rootEntityId) {
var entityRelationsQuery = new EntityRelationsQuery();
entityRelationsQuery.setParameters(new RelationsSearchParameters(rootEntityId, direction, maxLevel, false));
entityRelationsQuery.setFilters(Collections.singletonList(new RelationEntityTypeFilter(relationType, profiles)));
return entityRelationsQuery;
}
public List<EntityId> resolveEntityIds(List<EntityRelation> relations) {
return switch (direction) {
case FROM -> relations.stream().map(EntityRelation::getTo).toList();
case TO -> relations.stream().map(EntityRelation::getFrom).toList();
};
}
}

40
common/util/src/main/java/org/thingsboard/common/util/geo/CirclePerimeterDefinition.java

@ -0,0 +1,40 @@
/**
* 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.common.util.geo;
import lombok.Data;
@Data
public class CirclePerimeterDefinition implements PerimeterDefinition {
private Double centerLatitude;
private Double centerLongitude;
private Double range;
private RangeUnit rangeUnit;
@Override
public PerimeterType getType() {
return PerimeterType.CIRCLE;
}
@Override
public boolean checkMatches(Coordinates entityCoordinates) {
Coordinates perimeterCoordinates = new Coordinates(centerLatitude, centerLongitude);
return range > GeoUtil.distance(entityCoordinates, perimeterCoordinates, rangeUnit);
}
}

39
common/util/src/main/java/org/thingsboard/common/util/geo/PerimeterDefinition.java

@ -0,0 +1,39 @@
/**
* 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.common.util.geo;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import java.io.Serializable;
@JsonTypeInfo(
use = JsonTypeInfo.Id.NAME,
include = JsonTypeInfo.As.PROPERTY,
property = "type")
@JsonSubTypes({
@JsonSubTypes.Type(value = PolygonPerimeterDefinition.class, name = "POLYGON"),
@JsonSubTypes.Type(value = CirclePerimeterDefinition.class, name = "CIRCLE")})
@JsonIgnoreProperties(ignoreUnknown = true)
public interface PerimeterDefinition extends Serializable {
@JsonIgnore
PerimeterType getType();
boolean checkMatches(Coordinates entityCoordinates);
}

35
common/util/src/main/java/org/thingsboard/common/util/geo/PolygonPerimeterDefinition.java

@ -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.common.util.geo;
import lombok.Data;
@Data
public class PolygonPerimeterDefinition implements PerimeterDefinition {
private String polygonsDefinition;
@Override
public PerimeterType getType() {
return PerimeterType.POLYGON;
}
@Override
public boolean checkMatches(Coordinates entityCoordinates) {
return GeoUtil.contains(polygonsDefinition, entityCoordinates);
}
}
Loading…
Cancel
Save