committed by
GitHub
279 changed files with 9832 additions and 3563 deletions
File diff suppressed because one or more lines are too long
@ -0,0 +1,252 @@ |
|||||
|
/** |
||||
|
* 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.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.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.RelationTypeGroup; |
||||
|
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.CalculatedFieldState; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; |
||||
|
|
||||
|
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.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.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry; |
||||
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createStateByType; |
||||
|
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 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(); |
||||
|
|
||||
|
public ListenableFuture<CalculatedFieldState> fetchStateFromDb(CalculatedFieldCtx ctx, EntityId entityId) { |
||||
|
Map<String, ListenableFuture<ArgumentEntry>> argFutures = switch (ctx.getCalculatedField().getType()) { |
||||
|
case GEOFENCING -> fetchGeofencingCalculatedFieldArguments(ctx, entityId, false); |
||||
|
case SIMPLE, SCRIPT -> { |
||||
|
Map<String, ListenableFuture<ArgumentEntry>> futures = new HashMap<>(); |
||||
|
for (var entry : ctx.getArguments().entrySet()) { |
||||
|
var argEntityId = resolveEntityId(entityId, entry.getValue()); |
||||
|
var argValueFuture = fetchArgumentValue(ctx.getTenantId(), argEntityId, entry.getValue(), System.currentTimeMillis()); |
||||
|
futures.put(entry.getKey(), argValueFuture); |
||||
|
} |
||||
|
yield futures; |
||||
|
} |
||||
|
}; |
||||
|
return Futures.whenAllComplete(argFutures.values()).call(() -> { |
||||
|
var result = createStateByType(ctx); |
||||
|
result.updateState(ctx, resolveArgumentFutures(argFutures)); |
||||
|
// TODO: move to state.init() method after merge with alarm rules 2.0
|
||||
|
if (ctx.hasRelationQueryDynamicArguments() && result instanceof GeofencingCalculatedFieldState geofencingCalculatedFieldState) { |
||||
|
geofencingCalculatedFieldState.setLastDynamicArgumentsRefreshTs(System.currentTimeMillis()); |
||||
|
} |
||||
|
return result; |
||||
|
}, MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
protected EntityId resolveEntityId(EntityId entityId, Argument argument) { |
||||
|
return argument.getRefEntityId() != null ? argument.getRefEntityId() : 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 Map<String, ListenableFuture<ArgumentEntry>> fetchGeofencingCalculatedFieldArguments(CalculatedFieldCtx ctx, EntityId entityId, boolean dynamicArgumentsOnly) { |
||||
|
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().hasDynamicSource()) |
||||
|
.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(), System.currentTimeMillis())); |
||||
|
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; |
||||
|
} |
||||
|
|
||||
|
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)); |
||||
|
} |
||||
|
var refDynamicSourceConfiguration = value.getRefDynamicSourceConfiguration(); |
||||
|
return switch (refDynamicSourceConfiguration.getType()) { |
||||
|
case RELATION_PATH_QUERY -> { |
||||
|
var configuration = (RelationPathQueryDynamicSourceConfiguration) refDynamicSourceConfiguration; |
||||
|
yield Futures.transform(relationService.findByRelationPathQueryAsync(tenantId, configuration.toRelationPathQuery(entityId)), |
||||
|
configuration::resolveEntityIds, calculatedFieldCallbackExecutor); |
||||
|
} |
||||
|
}; |
||||
|
} |
||||
|
|
||||
|
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(() -> |
||||
|
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))), 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, 0L)); |
||||
|
return transformSingleValueArgument(Optional.of(attributeKvEntry)); |
||||
|
}, calculatedFieldCallbackExecutor); |
||||
|
} |
||||
|
|
||||
|
protected ListenableFuture<ArgumentEntry> fetchTsLatest(TenantId tenantId, EntityId entityId, Argument argument, long startTs) { |
||||
|
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(System.currentTimeMillis(), createDefaultKvEntry(argument), 0L))); |
||||
|
}, 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,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.TbelCfTsGeofencingArg; |
||||
|
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 TbelCfTsGeofencingArg(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,188 @@ |
|||||
|
/** |
||||
|
* 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.Data; |
||||
|
import lombok.EqualsAndHashCode; |
||||
|
import lombok.NoArgsConstructor; |
||||
|
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.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; |
||||
|
|
||||
|
@Data |
||||
|
@Slf4j |
||||
|
@NoArgsConstructor |
||||
|
@EqualsAndHashCode(callSuper = true) |
||||
|
public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { |
||||
|
|
||||
|
private long lastDynamicArgumentsRefreshTs = -1; |
||||
|
|
||||
|
public GeofencingCalculatedFieldState(List<String> requiredArguments) { |
||||
|
super(requiredArguments); |
||||
|
} |
||||
|
|
||||
|
@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(EntityId entityId, 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 = new CalculatedFieldResult(ctx.getOutput().getType(), ctx.getOutput().getScope(), toResultNode(outputType, valuesNode)); |
||||
|
if (relationFutures.isEmpty()) { |
||||
|
return Futures.immediateFuture(result); |
||||
|
} |
||||
|
return Futures.whenAllComplete(relationFutures).call(() -> result, MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
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) { |
||||
|
if (OutputType.ATTRIBUTES.equals(outputType) || latestTimestamp == -1) { |
||||
|
return valuesNode; |
||||
|
} |
||||
|
ObjectNode resultNode = JacksonUtil.newObjectNode(); |
||||
|
resultNode.put("ts", latestTimestamp); |
||||
|
resultNode.set("values", valuesNode); |
||||
|
return resultNode; |
||||
|
} |
||||
|
|
||||
|
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,75 @@ |
|||||
|
/** |
||||
|
* 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.utils; |
||||
|
|
||||
|
import com.google.common.util.concurrent.Futures; |
||||
|
import com.google.common.util.concurrent.ListenableFuture; |
||||
|
import com.google.common.util.concurrent.MoreExecutors; |
||||
|
import org.apache.commons.lang3.math.NumberUtils; |
||||
|
import org.thingsboard.server.common.data.StringUtils; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.Argument; |
||||
|
import org.thingsboard.server.common.data.kv.BooleanDataEntry; |
||||
|
import org.thingsboard.server.common.data.kv.DoubleDataEntry; |
||||
|
import org.thingsboard.server.common.data.kv.KvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.StringDataEntry; |
||||
|
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.geofencing.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; |
||||
|
|
||||
|
import java.util.Optional; |
||||
|
|
||||
|
public class CalculatedFieldArgumentUtils { |
||||
|
|
||||
|
public static ListenableFuture<ArgumentEntry> transformSingleValueArgument(ListenableFuture<Optional<? extends KvEntry>> kvEntryFuture) { |
||||
|
return Futures.transform(kvEntryFuture, CalculatedFieldArgumentUtils::transformSingleValueArgument, MoreExecutors.directExecutor()); |
||||
|
} |
||||
|
|
||||
|
public static ArgumentEntry transformSingleValueArgument(Optional<? extends KvEntry> kvEntry) { |
||||
|
if (kvEntry.isPresent() && kvEntry.get().getValue() != null) { |
||||
|
return ArgumentEntry.createSingleValueArgument(kvEntry.get()); |
||||
|
} else { |
||||
|
return new SingleValueArgumentEntry(); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public static KvEntry createDefaultKvEntry(Argument argument) { |
||||
|
String key = argument.getRefEntityKey().getKey(); |
||||
|
String defaultValue = argument.getDefaultValue(); |
||||
|
if (StringUtils.isBlank(defaultValue)) { |
||||
|
return new StringDataEntry(key, null); |
||||
|
} |
||||
|
if (NumberUtils.isParsable(defaultValue)) { |
||||
|
return new DoubleDataEntry(key, Double.parseDouble(defaultValue)); |
||||
|
} |
||||
|
if ("true".equalsIgnoreCase(defaultValue) || "false".equalsIgnoreCase(defaultValue)) { |
||||
|
return new BooleanDataEntry(key, Boolean.parseBoolean(defaultValue)); |
||||
|
} |
||||
|
return new StringDataEntry(key, defaultValue); |
||||
|
} |
||||
|
|
||||
|
public static CalculatedFieldState createStateByType(CalculatedFieldCtx ctx) { |
||||
|
return switch (ctx.getCfType()) { |
||||
|
case SIMPLE -> new SimpleCalculatedFieldState(ctx.getArgNames()); |
||||
|
case SCRIPT -> new ScriptCalculatedFieldState(ctx.getArgNames()); |
||||
|
case GEOFENCING -> new GeofencingCalculatedFieldState(ctx.getArgNames()); |
||||
|
}; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,478 @@ |
|||||
|
/** |
||||
|
* 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.google.common.util.concurrent.Futures; |
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.mockito.ArgumentCaptor; |
||||
|
import org.mockito.Mock; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.cf.CalculatedField; |
||||
|
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.Output; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.OutputType; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates; |
||||
|
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.ZoneGroupConfiguration; |
||||
|
import org.thingsboard.server.common.data.id.AssetId; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.DoubleDataEntry; |
||||
|
import org.thingsboard.server.common.data.kv.JsonDataEntry; |
||||
|
import org.thingsboard.server.common.data.relation.EntityRelation; |
||||
|
import org.thingsboard.server.common.data.relation.EntitySearchDirection; |
||||
|
import org.thingsboard.server.common.data.relation.RelationPathLevel; |
||||
|
import org.thingsboard.server.dao.relation.RelationService; |
||||
|
import org.thingsboard.server.dao.usagerecord.ApiLimitService; |
||||
|
import org.thingsboard.server.service.cf.CalculatedFieldResult; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; |
||||
|
|
||||
|
import java.util.HashMap; |
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.ExecutionException; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.assertj.core.api.Assertions.assertThatThrownBy; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.ArgumentMatchers.eq; |
||||
|
import static org.mockito.Mockito.times; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
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.GeofencingReportStrategy.REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
public class GeofencingCalculatedFieldStateTest { |
||||
|
|
||||
|
private final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("8f83eeca-b5cd-4955-9241-09d1393768c6")); |
||||
|
private final DeviceId DEVICE_ID = new DeviceId(UUID.fromString("688b529d-cfbe-4430-91c5-60b4f4e5d3cf")); |
||||
|
private final AssetId ZONE_1_ID = new AssetId(UUID.fromString("c0e3031c-7df1-45e4-9590-cfd621a4d714")); |
||||
|
private final AssetId ZONE_2_ID = new AssetId(UUID.fromString("e7da6200-2096-4038-a343-ade9ea4fa3e4")); |
||||
|
|
||||
|
private final SingleValueArgumentEntry latitudeArgEntry = new SingleValueArgumentEntry(System.currentTimeMillis() - 10, new DoubleDataEntry("latitude", 50.4730), 145L); |
||||
|
private final SingleValueArgumentEntry longitudeArgEntry = new SingleValueArgumentEntry(System.currentTimeMillis() - 6, new DoubleDataEntry("longitude", 30.5050), 165L); |
||||
|
|
||||
|
private final JsonDataEntry allowedZoneDataEntry = new JsonDataEntry("zone", "[[50.472000, 30.504000], [50.472000, 30.506000], [50.474000, 30.506000], [50.474000, 30.504000]]"); |
||||
|
private final BaseAttributeKvEntry allowedZoneAttributeKvEntry = new BaseAttributeKvEntry(allowedZoneDataEntry, System.currentTimeMillis(), 0L); |
||||
|
private final GeofencingArgumentEntry geofencingAllowedZoneArgEntry = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, allowedZoneAttributeKvEntry)); |
||||
|
|
||||
|
private final JsonDataEntry restrictedZoneDataEntry = new JsonDataEntry("zone", "[[50.475000, 30.510000], [50.475000, 30.512000], [50.477000, 30.512000], [50.477000, 30.510000]]"); |
||||
|
private final BaseAttributeKvEntry restrictedZoneAttributeKvEntry = new BaseAttributeKvEntry(restrictedZoneDataEntry, System.currentTimeMillis(), 0L); |
||||
|
private final GeofencingArgumentEntry geofencingRestrictedZoneArgEntry = new GeofencingArgumentEntry(Map.of(ZONE_2_ID, restrictedZoneAttributeKvEntry)); |
||||
|
|
||||
|
|
||||
|
private GeofencingCalculatedFieldState state; |
||||
|
private CalculatedFieldCtx ctx; |
||||
|
|
||||
|
@Mock |
||||
|
private ApiLimitService apiLimitService; |
||||
|
@Mock |
||||
|
private RelationService relationService; |
||||
|
|
||||
|
@BeforeEach |
||||
|
void setUp() { |
||||
|
when(apiLimitService.getLimit(any(), any())).thenReturn(1000L); |
||||
|
ctx = new CalculatedFieldCtx(getCalculatedField(), null, apiLimitService, relationService); |
||||
|
ctx.init(); |
||||
|
state = new GeofencingCalculatedFieldState(ctx.getArgNames()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testType() { |
||||
|
assertThat(state.getType()).isEqualTo(CalculatedFieldType.GEOFENCING); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateState() { |
||||
|
state.arguments = new HashMap<>(Map.of( |
||||
|
ENTITY_ID_LATITUDE_ARGUMENT_KEY, latitudeArgEntry, |
||||
|
ENTITY_ID_LONGITUDE_ARGUMENT_KEY, longitudeArgEntry |
||||
|
)); |
||||
|
|
||||
|
Map<String, ArgumentEntry> newArgs = Map.of("allowedZones", geofencingAllowedZoneArgEntry); |
||||
|
boolean stateUpdated = state.updateState(ctx, newArgs); |
||||
|
|
||||
|
assertThat(stateUpdated).isTrue(); |
||||
|
assertThat(state.getArguments()).containsExactlyInAnyOrderEntriesOf( |
||||
|
Map.of( |
||||
|
ENTITY_ID_LATITUDE_ARGUMENT_KEY, latitudeArgEntry, |
||||
|
ENTITY_ID_LONGITUDE_ARGUMENT_KEY, longitudeArgEntry, |
||||
|
"allowedZones", geofencingAllowedZoneArgEntry |
||||
|
) |
||||
|
); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateStateWithInvalidArgumentTypeForLatitudeArgument() { |
||||
|
assertThatThrownBy(() -> state.updateState(ctx, Map.of(ENTITY_ID_LATITUDE_ARGUMENT_KEY, geofencingAllowedZoneArgEntry))) |
||||
|
.isInstanceOf(IllegalArgumentException.class) |
||||
|
.hasMessage("Unsupported argument entry type for latitude argument: GEOFENCING. Only SINGLE_VALUE type is allowed."); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateStateWithInvalidArgumentTypeForLongitudeArgument() { |
||||
|
assertThatThrownBy(() -> state.updateState(ctx, Map.of(ENTITY_ID_LONGITUDE_ARGUMENT_KEY, geofencingAllowedZoneArgEntry))) |
||||
|
.isInstanceOf(IllegalArgumentException.class) |
||||
|
.hasMessage("Unsupported argument entry type for longitude argument: GEOFENCING. Only SINGLE_VALUE type is allowed."); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateStateWithInvalidArgumentTypeForGeofencingArgument() { |
||||
|
assertThatThrownBy(() -> state.updateState(ctx, Map.of("someArgumentName", latitudeArgEntry))) |
||||
|
.isInstanceOf(IllegalArgumentException.class) |
||||
|
.hasMessage("Unsupported argument entry type for someArgumentName argument: SINGLE_VALUE. Only GEOFENCING type is allowed."); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateStateWhenUpdateExistingSingleValueArgumentEntry() { |
||||
|
state.arguments = new HashMap<>(Map.of(ENTITY_ID_LATITUDE_ARGUMENT_KEY, latitudeArgEntry)); |
||||
|
|
||||
|
SingleValueArgumentEntry newArgEntry = new SingleValueArgumentEntry(System.currentTimeMillis(), new DoubleDataEntry("latitude", 50.4760), 190L); |
||||
|
Map<String, ArgumentEntry> newArgs = Map.of("latitude", newArgEntry); |
||||
|
boolean stateUpdated = state.updateState(ctx, newArgs); |
||||
|
|
||||
|
assertThat(stateUpdated).isTrue(); |
||||
|
assertThat(state.getArguments()).isEqualTo(newArgs); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateStateWhenUpdateExistingGeofencingValueArgumentEntryWithTheSameValue() { |
||||
|
state.arguments = new HashMap<>(Map.of("allowedZones", geofencingAllowedZoneArgEntry)); |
||||
|
|
||||
|
Map<String, ArgumentEntry> newArgs = Map.of("allowedZones", geofencingAllowedZoneArgEntry); |
||||
|
|
||||
|
boolean stateUpdated = state.updateState(ctx, newArgs); |
||||
|
|
||||
|
assertThat(stateUpdated).isFalse(); |
||||
|
assertThat(state.getArguments()).isEqualTo(newArgs); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateStateWhenUpdateExistingSingleValueArgumentEntryWithValueOfAnotherType() { |
||||
|
state.arguments = new HashMap<>(Map.of(ENTITY_ID_LATITUDE_ARGUMENT_KEY, latitudeArgEntry)); |
||||
|
|
||||
|
assertThatThrownBy(() -> state.updateState(ctx, Map.of(ENTITY_ID_LATITUDE_ARGUMENT_KEY, geofencingAllowedZoneArgEntry))) |
||||
|
.isInstanceOf(IllegalArgumentException.class) |
||||
|
.hasMessage("Unsupported argument entry type for single value argument entry: GEOFENCING"); |
||||
|
} |
||||
|
|
||||
|
|
||||
|
@Test |
||||
|
void testUpdateStateWhenUpdateExistingGeofencingValueArgumentEntryWithValueOfAnotherType() { |
||||
|
state.arguments = new HashMap<>(Map.of("allowedZones", geofencingAllowedZoneArgEntry)); |
||||
|
|
||||
|
assertThatThrownBy(() -> state.updateState(ctx, Map.of("allowedZones", latitudeArgEntry))) |
||||
|
.isInstanceOf(IllegalArgumentException.class) |
||||
|
.hasMessage("Unsupported argument entry type for geofencing argument entry: SINGLE_VALUE"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testIsReadyWhenNotAllArgPresent() { |
||||
|
assertThat(state.isReady()).isFalse(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testIsReadyWhenAllArgPresent() { |
||||
|
state.arguments = new HashMap<>(Map.of( |
||||
|
ENTITY_ID_LATITUDE_ARGUMENT_KEY, latitudeArgEntry, |
||||
|
ENTITY_ID_LONGITUDE_ARGUMENT_KEY, longitudeArgEntry, |
||||
|
"allowedZones", geofencingAllowedZoneArgEntry, |
||||
|
"restrictedZones", geofencingRestrictedZoneArgEntry |
||||
|
)); |
||||
|
assertThat(state.isReady()).isTrue(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testIsReadyWhenEmptyEntryPresents() { |
||||
|
state.arguments = new HashMap<>(Map.of( |
||||
|
ENTITY_ID_LATITUDE_ARGUMENT_KEY, latitudeArgEntry, |
||||
|
ENTITY_ID_LONGITUDE_ARGUMENT_KEY, longitudeArgEntry, |
||||
|
"allowedZones", geofencingAllowedZoneArgEntry, |
||||
|
"restrictedZones", geofencingRestrictedZoneArgEntry |
||||
|
)); |
||||
|
|
||||
|
state.getArguments().put("noParkingZones", new GeofencingArgumentEntry()); |
||||
|
|
||||
|
assertThat(state.isReady()).isFalse(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testPerformCalculation() throws ExecutionException, InterruptedException { |
||||
|
state.arguments = new HashMap<>(Map.of( |
||||
|
ENTITY_ID_LATITUDE_ARGUMENT_KEY, latitudeArgEntry, |
||||
|
ENTITY_ID_LONGITUDE_ARGUMENT_KEY, longitudeArgEntry, |
||||
|
"allowedZones", geofencingAllowedZoneArgEntry, |
||||
|
"restrictedZones", geofencingRestrictedZoneArgEntry |
||||
|
)); |
||||
|
|
||||
|
Output output = ctx.getOutput(); |
||||
|
var configuration = (GeofencingCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); |
||||
|
|
||||
|
when(relationService.saveRelationAsync(any(), any())).thenReturn(Futures.immediateFuture(true)); |
||||
|
when(relationService.deleteRelationAsync(any(), any())).thenReturn(Futures.immediateFuture(true)); |
||||
|
|
||||
|
CalculatedFieldResult result = state.performCalculation(ctx.getEntityId(), ctx).get(); |
||||
|
|
||||
|
assertThat(result).isNotNull(); |
||||
|
assertThat(result.getType()).isEqualTo(output.getType()); |
||||
|
assertThat(result.getScope()).isEqualTo(output.getScope()); |
||||
|
assertThat(result.getResult()).isEqualTo( |
||||
|
JacksonUtil.newObjectNode() |
||||
|
.put("allowedZonesEvent", "ENTERED") |
||||
|
.put("allowedZonesStatus", "INSIDE") |
||||
|
.put("restrictedZonesStatus", "OUTSIDE") |
||||
|
); |
||||
|
|
||||
|
SingleValueArgumentEntry newLatitude = new SingleValueArgumentEntry(System.currentTimeMillis(), new DoubleDataEntry("latitude", 50.4760), 146L); |
||||
|
SingleValueArgumentEntry newLongitude = new SingleValueArgumentEntry(System.currentTimeMillis(), new DoubleDataEntry("longitude", 30.5110), 166L); |
||||
|
|
||||
|
// move the device to new coordinates → leaves allowed, enters restricted
|
||||
|
state.updateState(ctx, Map.of(ENTITY_ID_LATITUDE_ARGUMENT_KEY, newLatitude, ENTITY_ID_LONGITUDE_ARGUMENT_KEY, newLongitude)); |
||||
|
|
||||
|
CalculatedFieldResult result2 = state.performCalculation(ctx.getEntityId(), ctx).get(); |
||||
|
|
||||
|
assertThat(result2).isNotNull(); |
||||
|
assertThat(result2.getType()).isEqualTo(output.getType()); |
||||
|
assertThat(result2.getScope()).isEqualTo(output.getScope()); |
||||
|
assertThat(result2.getResult().get("values")).isEqualTo( |
||||
|
JacksonUtil.newObjectNode() |
||||
|
.put("allowedZonesEvent", "LEFT") |
||||
|
.put("allowedZonesStatus", "OUTSIDE") |
||||
|
.put("restrictedZonesEvent", "ENTERED") |
||||
|
.put("restrictedZonesStatus", "INSIDE") |
||||
|
); |
||||
|
|
||||
|
// Check relations are created and deleted correctly for both iterations.
|
||||
|
ArgumentCaptor<EntityRelation> saveCaptor = ArgumentCaptor.forClass(EntityRelation.class); |
||||
|
verify(relationService, times(2)).saveRelationAsync(eq(ctx.getTenantId()), saveCaptor.capture()); |
||||
|
List<EntityRelation> saveValues = saveCaptor.getAllValues(); |
||||
|
assertThat(saveValues).hasSize(2); |
||||
|
|
||||
|
EntityRelation relationFromFirstIteration = saveValues.get(0); |
||||
|
assertThat(relationFromFirstIteration.getTo()).isEqualTo(ctx.getEntityId()); |
||||
|
assertThat(relationFromFirstIteration.getFrom()).isEqualTo(ZONE_1_ID); |
||||
|
assertThat(relationFromFirstIteration.getType()).isEqualTo("CurrentZone"); |
||||
|
|
||||
|
EntityRelation relationFromSecondIteration = saveValues.get(1); |
||||
|
assertThat(relationFromSecondIteration.getTo()).isEqualTo(ctx.getEntityId()); |
||||
|
assertThat(relationFromSecondIteration.getFrom()).isEqualTo(ZONE_2_ID); |
||||
|
assertThat(relationFromSecondIteration.getType()).isEqualTo("CurrentZone"); |
||||
|
|
||||
|
ArgumentCaptor<EntityRelation> deleteCaptor = ArgumentCaptor.forClass(EntityRelation.class); |
||||
|
verify(relationService).deleteRelationAsync(eq(ctx.getTenantId()), deleteCaptor.capture()); |
||||
|
EntityRelation leftRelation = deleteCaptor.getValue(); |
||||
|
assertThat(leftRelation.getFrom()).isEqualTo(ZONE_1_ID); |
||||
|
assertThat(leftRelation.getTo()).isEqualTo(ctx.getEntityId()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testPerformCalculationWithOnlyTransitionEventsReportingStrategy() throws ExecutionException, InterruptedException { |
||||
|
state.arguments = new HashMap<>(Map.of( |
||||
|
ENTITY_ID_LATITUDE_ARGUMENT_KEY, latitudeArgEntry, |
||||
|
ENTITY_ID_LONGITUDE_ARGUMENT_KEY, longitudeArgEntry, |
||||
|
"allowedZones", geofencingAllowedZoneArgEntry, |
||||
|
"restrictedZones", geofencingRestrictedZoneArgEntry |
||||
|
)); |
||||
|
|
||||
|
Output output = ctx.getOutput(); |
||||
|
|
||||
|
var calculatedFieldConfig = getCalculatedFieldConfig(GeofencingReportStrategy.REPORT_TRANSITION_EVENTS_ONLY); |
||||
|
|
||||
|
ctx.setCalculatedField(getCalculatedField(calculatedFieldConfig)); |
||||
|
ctx.init(); |
||||
|
|
||||
|
var configuration = (GeofencingCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); |
||||
|
|
||||
|
when(relationService.saveRelationAsync(any(), any())).thenReturn(Futures.immediateFuture(true)); |
||||
|
when(relationService.deleteRelationAsync(any(), any())).thenReturn(Futures.immediateFuture(true)); |
||||
|
|
||||
|
CalculatedFieldResult result = state.performCalculation(ctx.getEntityId(), ctx).get(); |
||||
|
|
||||
|
assertThat(result).isNotNull(); |
||||
|
assertThat(result.getType()).isEqualTo(output.getType()); |
||||
|
assertThat(result.getScope()).isEqualTo(output.getScope()); |
||||
|
assertThat(result.getResult()).isEqualTo( |
||||
|
JacksonUtil.newObjectNode().put("allowedZonesEvent", "ENTERED") |
||||
|
); |
||||
|
|
||||
|
SingleValueArgumentEntry newLatitude = new SingleValueArgumentEntry(System.currentTimeMillis(), new DoubleDataEntry("latitude", 50.4760), 146L); |
||||
|
SingleValueArgumentEntry newLongitude = new SingleValueArgumentEntry(System.currentTimeMillis(), new DoubleDataEntry("longitude", 30.5110), 166L); |
||||
|
|
||||
|
// move the device to new coordinates → leaves allowed, enters restricted
|
||||
|
state.updateState(ctx, Map.of(ENTITY_ID_LATITUDE_ARGUMENT_KEY, newLatitude, ENTITY_ID_LONGITUDE_ARGUMENT_KEY, newLongitude)); |
||||
|
|
||||
|
CalculatedFieldResult result2 = state.performCalculation(ctx.getEntityId(), ctx).get(); |
||||
|
|
||||
|
assertThat(result2).isNotNull(); |
||||
|
assertThat(result2.getType()).isEqualTo(output.getType()); |
||||
|
assertThat(result2.getScope()).isEqualTo(output.getScope()); |
||||
|
assertThat(result2.getResult().get("values")).isEqualTo( |
||||
|
JacksonUtil.newObjectNode() |
||||
|
.put("allowedZonesEvent", "LEFT") |
||||
|
.put("restrictedZonesEvent", "ENTERED") |
||||
|
); |
||||
|
|
||||
|
// Check relations are created and deleted correctly for both iterations.
|
||||
|
ArgumentCaptor<EntityRelation> saveCaptor = ArgumentCaptor.forClass(EntityRelation.class); |
||||
|
verify(relationService, times(2)).saveRelationAsync(eq(ctx.getTenantId()), saveCaptor.capture()); |
||||
|
List<EntityRelation> saveValues = saveCaptor.getAllValues(); |
||||
|
assertThat(saveValues).hasSize(2); |
||||
|
|
||||
|
EntityRelation relationFromFirstIteration = saveValues.get(0); |
||||
|
assertThat(relationFromFirstIteration.getTo()).isEqualTo(ctx.getEntityId()); |
||||
|
assertThat(relationFromFirstIteration.getFrom()).isEqualTo(ZONE_1_ID); |
||||
|
assertThat(relationFromFirstIteration.getType()).isEqualTo("CurrentZone"); |
||||
|
|
||||
|
EntityRelation relationFromSecondIteration = saveValues.get(1); |
||||
|
assertThat(relationFromSecondIteration.getTo()).isEqualTo(ctx.getEntityId()); |
||||
|
assertThat(relationFromSecondIteration.getFrom()).isEqualTo(ZONE_2_ID); |
||||
|
assertThat(relationFromSecondIteration.getType()).isEqualTo("CurrentZone"); |
||||
|
|
||||
|
ArgumentCaptor<EntityRelation> deleteCaptor = ArgumentCaptor.forClass(EntityRelation.class); |
||||
|
verify(relationService).deleteRelationAsync(eq(ctx.getTenantId()), deleteCaptor.capture()); |
||||
|
EntityRelation leftRelation = deleteCaptor.getValue(); |
||||
|
assertThat(leftRelation.getFrom()).isEqualTo(ZONE_1_ID); |
||||
|
assertThat(leftRelation.getTo()).isEqualTo(ctx.getEntityId()); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testPerformCalculationWithOnlyPresenceStatusReportingStrategy() throws ExecutionException, InterruptedException { |
||||
|
state.arguments = new HashMap<>(Map.of( |
||||
|
ENTITY_ID_LATITUDE_ARGUMENT_KEY, latitudeArgEntry, |
||||
|
ENTITY_ID_LONGITUDE_ARGUMENT_KEY, longitudeArgEntry, |
||||
|
"allowedZones", geofencingAllowedZoneArgEntry, |
||||
|
"restrictedZones", geofencingRestrictedZoneArgEntry |
||||
|
)); |
||||
|
|
||||
|
Output output = ctx.getOutput(); |
||||
|
|
||||
|
var calculatedFieldConfig = getCalculatedFieldConfig(GeofencingReportStrategy.REPORT_PRESENCE_STATUS_ONLY); |
||||
|
|
||||
|
ctx.setCalculatedField(getCalculatedField(calculatedFieldConfig)); |
||||
|
ctx.init(); |
||||
|
|
||||
|
var configuration = (GeofencingCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); |
||||
|
|
||||
|
when(relationService.saveRelationAsync(any(), any())).thenReturn(Futures.immediateFuture(true)); |
||||
|
when(relationService.deleteRelationAsync(any(), any())).thenReturn(Futures.immediateFuture(true)); |
||||
|
|
||||
|
CalculatedFieldResult result = state.performCalculation(ctx.getEntityId(), ctx).get(); |
||||
|
|
||||
|
assertThat(result).isNotNull(); |
||||
|
assertThat(result.getType()).isEqualTo(output.getType()); |
||||
|
assertThat(result.getScope()).isEqualTo(output.getScope()); |
||||
|
assertThat(result.getResult()).isEqualTo( |
||||
|
JacksonUtil.newObjectNode() |
||||
|
.put("allowedZonesStatus", "INSIDE") |
||||
|
.put("restrictedZonesStatus", "OUTSIDE") |
||||
|
); |
||||
|
|
||||
|
SingleValueArgumentEntry newLatitude = new SingleValueArgumentEntry(System.currentTimeMillis(), new DoubleDataEntry("latitude", 50.4760), 146L); |
||||
|
SingleValueArgumentEntry newLongitude = new SingleValueArgumentEntry(System.currentTimeMillis(), new DoubleDataEntry("longitude", 30.5110), 166L); |
||||
|
|
||||
|
// move the device to new coordinates → leaves allowed, enters restricted
|
||||
|
state.updateState(ctx, Map.of(ENTITY_ID_LATITUDE_ARGUMENT_KEY, newLatitude, ENTITY_ID_LONGITUDE_ARGUMENT_KEY, newLongitude)); |
||||
|
|
||||
|
CalculatedFieldResult result2 = state.performCalculation(ctx.getEntityId(), ctx).get(); |
||||
|
|
||||
|
assertThat(result2).isNotNull(); |
||||
|
assertThat(result2.getType()).isEqualTo(output.getType()); |
||||
|
assertThat(result2.getScope()).isEqualTo(output.getScope()); |
||||
|
assertThat(result2.getResult().get("values")).isEqualTo( |
||||
|
JacksonUtil.newObjectNode() |
||||
|
.put("allowedZonesStatus", "OUTSIDE") |
||||
|
.put("restrictedZonesStatus", "INSIDE") |
||||
|
); |
||||
|
|
||||
|
// Check relations are created and deleted correctly for both iterations.
|
||||
|
ArgumentCaptor<EntityRelation> saveCaptor = ArgumentCaptor.forClass(EntityRelation.class); |
||||
|
verify(relationService, times(2)).saveRelationAsync(eq(ctx.getTenantId()), saveCaptor.capture()); |
||||
|
List<EntityRelation> saveValues = saveCaptor.getAllValues(); |
||||
|
assertThat(saveValues).hasSize(2); |
||||
|
|
||||
|
EntityRelation relationFromFirstIteration = saveValues.get(0); |
||||
|
assertThat(relationFromFirstIteration.getTo()).isEqualTo(ctx.getEntityId()); |
||||
|
assertThat(relationFromFirstIteration.getFrom()).isEqualTo(ZONE_1_ID); |
||||
|
assertThat(relationFromFirstIteration.getType()).isEqualTo("CurrentZone"); |
||||
|
|
||||
|
EntityRelation relationFromSecondIteration = saveValues.get(1); |
||||
|
assertThat(relationFromSecondIteration.getTo()).isEqualTo(ctx.getEntityId()); |
||||
|
assertThat(relationFromSecondIteration.getFrom()).isEqualTo(ZONE_2_ID); |
||||
|
assertThat(relationFromSecondIteration.getType()).isEqualTo("CurrentZone"); |
||||
|
|
||||
|
ArgumentCaptor<EntityRelation> deleteCaptor = ArgumentCaptor.forClass(EntityRelation.class); |
||||
|
verify(relationService).deleteRelationAsync(eq(ctx.getTenantId()), deleteCaptor.capture()); |
||||
|
EntityRelation leftRelation = deleteCaptor.getValue(); |
||||
|
assertThat(leftRelation.getFrom()).isEqualTo(ZONE_1_ID); |
||||
|
assertThat(leftRelation.getTo()).isEqualTo(ctx.getEntityId()); |
||||
|
} |
||||
|
|
||||
|
private CalculatedField getCalculatedField() { |
||||
|
return getCalculatedField(getCalculatedFieldConfig(REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS)); |
||||
|
} |
||||
|
|
||||
|
private CalculatedField getCalculatedField(CalculatedFieldConfiguration configuration) { |
||||
|
CalculatedField calculatedField = new CalculatedField(); |
||||
|
calculatedField.setTenantId(TENANT_ID); |
||||
|
calculatedField.setEntityId(DEVICE_ID); |
||||
|
calculatedField.setType(CalculatedFieldType.GEOFENCING); |
||||
|
calculatedField.setName("Test Geofencing Calculated Field"); |
||||
|
calculatedField.setConfigurationVersion(1); |
||||
|
calculatedField.setConfiguration(configuration); |
||||
|
calculatedField.setVersion(1L); |
||||
|
return calculatedField; |
||||
|
} |
||||
|
|
||||
|
private CalculatedFieldConfiguration getCalculatedFieldConfig(GeofencingReportStrategy reportStrategy) { |
||||
|
var config = new GeofencingCalculatedFieldConfiguration(); |
||||
|
|
||||
|
EntityCoordinates entityCoordinates = new EntityCoordinates("latitude", "longitude"); |
||||
|
config.setEntityCoordinates(entityCoordinates); |
||||
|
|
||||
|
ZoneGroupConfiguration allowedZonesGroup = new ZoneGroupConfiguration("zone", reportStrategy, true); |
||||
|
var allowedZoneDynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); |
||||
|
allowedZoneDynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.TO, "AllowedZone"))); |
||||
|
allowedZonesGroup.setRefDynamicSourceConfiguration(allowedZoneDynamicSourceConfiguration); |
||||
|
allowedZonesGroup.setRelationType("CurrentZone"); |
||||
|
allowedZonesGroup.setDirection(EntitySearchDirection.TO); |
||||
|
|
||||
|
ZoneGroupConfiguration restrictedZonesGroup = new ZoneGroupConfiguration("zone", reportStrategy, true); |
||||
|
var restrictedZoneDynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); |
||||
|
restrictedZoneDynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.TO, "RestrictedZone"))); |
||||
|
restrictedZonesGroup.setRefDynamicSourceConfiguration(restrictedZoneDynamicSourceConfiguration); |
||||
|
restrictedZonesGroup.setRelationType("CurrentZone"); |
||||
|
restrictedZonesGroup.setDirection(EntitySearchDirection.TO); |
||||
|
|
||||
|
config.setZoneGroups(Map.of("allowedZones", allowedZonesGroup, "restrictedZones", restrictedZonesGroup)); |
||||
|
|
||||
|
Output output = new Output(); |
||||
|
output.setType(OutputType.TIME_SERIES); |
||||
|
config.setOutput(output); |
||||
|
return config; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,184 @@ |
|||||
|
/** |
||||
|
* 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 io.hypersistence.utils.hibernate.type.json.internal.JacksonUtil; |
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.thingsboard.common.util.geo.PerimeterDefinition; |
||||
|
import org.thingsboard.server.common.data.id.AssetId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.JsonDataEntry; |
||||
|
import org.thingsboard.server.common.data.kv.StringDataEntry; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingZoneState; |
||||
|
|
||||
|
import java.util.Map; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.assertj.core.api.Assertions.assertThatThrownBy; |
||||
|
|
||||
|
public class GeofencingValueArgumentEntryTest { |
||||
|
|
||||
|
private final AssetId ZONE_1_ID = new AssetId(UUID.fromString("c0e3031c-7df1-45e4-9590-cfd621a4d714")); |
||||
|
private final AssetId ZONE_2_ID = new AssetId(UUID.fromString("e7da6200-2096-4038-a343-ade9ea4fa3e4")); |
||||
|
|
||||
|
private final JsonDataEntry allowedZoneDataEntry = new JsonDataEntry("zone", "[[50.472000, 30.504000], [50.472000, 30.506000], [50.474000, 30.506000], [50.474000, 30.504000]]"); |
||||
|
private final BaseAttributeKvEntry allowedZoneAttributeKvEntry = new BaseAttributeKvEntry(allowedZoneDataEntry, 363L, 155L); |
||||
|
|
||||
|
private final JsonDataEntry restrictedZoneDataEntry = new JsonDataEntry("zone", "[[50.475000, 30.510000], [50.475000, 30.512000], [50.477000, 30.512000], [50.477000, 30.510000]]"); |
||||
|
private final BaseAttributeKvEntry restrictedZoneAttributeKvEntry = new BaseAttributeKvEntry(restrictedZoneDataEntry, 363L, 155L); |
||||
|
|
||||
|
private GeofencingArgumentEntry entry; |
||||
|
|
||||
|
@BeforeEach |
||||
|
void setUp() { |
||||
|
entry = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, allowedZoneAttributeKvEntry, ZONE_2_ID, restrictedZoneAttributeKvEntry)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testArgumentEntryType() { |
||||
|
assertThat(entry.getType()).isEqualTo(ArgumentEntryType.GEOFENCING); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateEntryWhenSingleEntryPassed() { |
||||
|
assertThatThrownBy(() -> entry.updateEntry(new SingleValueArgumentEntry())) |
||||
|
.isInstanceOf(IllegalArgumentException.class) |
||||
|
.hasMessage("Unsupported argument entry type for geofencing argument entry: SINGLE_VALUE"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateEntryWhenRollingEntryPassed() { |
||||
|
assertThatThrownBy(() -> entry.updateEntry(new TsRollingArgumentEntry(5, 30000L))) |
||||
|
.isInstanceOf(IllegalArgumentException.class) |
||||
|
.hasMessage("Unsupported argument entry type for geofencing argument entry: TS_ROLLING"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateEntryWithTheSameTs() { |
||||
|
BaseAttributeKvEntry differentValueSameTs = new BaseAttributeKvEntry(new JsonDataEntry("zone", "[[50.472001, 30.504001], [50.472001, 30.506001], [50.474001, 30.506001], [50.474001, 30.504001]]"), 363L, 156L); |
||||
|
var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, differentValueSameTs, ZONE_2_ID, restrictedZoneAttributeKvEntry)); |
||||
|
assertThat(entry.updateEntry(updated)).isFalse(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
@SuppressWarnings("unchecked") |
||||
|
void testUpdateEntryWhenNewVersionIsNull() { |
||||
|
BaseAttributeKvEntry differentValueNewVersionIsNull = new BaseAttributeKvEntry(new JsonDataEntry("zone", "[[50.472001, 30.504001], [50.472001, 30.506001], [50.474001, 30.506001], [50.474001, 30.504001]]"), 364L, null); |
||||
|
var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, differentValueNewVersionIsNull, ZONE_2_ID, restrictedZoneAttributeKvEntry)); |
||||
|
|
||||
|
assertThat(entry.updateEntry(updated)).isTrue(); |
||||
|
assertThat(entry.getValue()).isInstanceOf(Map.class); |
||||
|
|
||||
|
Map<EntityId, GeofencingZoneState> value = (Map<EntityId, GeofencingZoneState>) entry.getValue(); |
||||
|
assertThat(value).hasSize(2); |
||||
|
assertThat(value.get(ZONE_1_ID).getVersion()).isNull(); |
||||
|
assertThat(value.get(ZONE_1_ID).getTs()).isEqualTo(364L); |
||||
|
assertThat(value.get(ZONE_1_ID).getPerimeterDefinition()) |
||||
|
.isEqualTo(JacksonUtil.fromString(differentValueNewVersionIsNull.getJsonValue().get(), PerimeterDefinition.class)); |
||||
|
|
||||
|
assertThat(value.get(ZONE_2_ID).getVersion()).isEqualTo(155L); |
||||
|
assertThat(value.get(ZONE_2_ID).getTs()).isEqualTo(363L); |
||||
|
assertThat(value.get(ZONE_2_ID).getPerimeterDefinition()) |
||||
|
.isEqualTo(JacksonUtil.fromString(restrictedZoneAttributeKvEntry.getJsonValue().get(), PerimeterDefinition.class)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
@SuppressWarnings("unchecked") |
||||
|
void testUpdateEntryWhenNewVersionIsGreaterThanCurrent() { |
||||
|
BaseAttributeKvEntry differentValueNewVersionIsSet = new BaseAttributeKvEntry(new JsonDataEntry("zone", "[[50.472001, 30.504001], [50.472001, 30.506001], [50.474001, 30.506001], [50.474001, 30.504001]]"), 364L, 156L); |
||||
|
var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, differentValueNewVersionIsSet, ZONE_2_ID, restrictedZoneAttributeKvEntry)); |
||||
|
|
||||
|
assertThat(entry.updateEntry(updated)).isTrue(); |
||||
|
assertThat(entry.getValue()).isInstanceOf(Map.class); |
||||
|
|
||||
|
Map<EntityId, GeofencingZoneState> value = (Map<EntityId, GeofencingZoneState>) entry.getValue(); |
||||
|
assertThat(value).hasSize(2); |
||||
|
assertThat(value.get(ZONE_1_ID).getVersion()).isEqualTo(156L); |
||||
|
assertThat(value.get(ZONE_1_ID).getTs()).isEqualTo(364L); |
||||
|
assertThat(value.get(ZONE_1_ID).getPerimeterDefinition()) |
||||
|
.isEqualTo(JacksonUtil.fromString(differentValueNewVersionIsSet.getJsonValue().get(), PerimeterDefinition.class)); |
||||
|
|
||||
|
assertThat(value.get(ZONE_2_ID).getVersion()).isEqualTo(155L); |
||||
|
assertThat(value.get(ZONE_2_ID).getTs()).isEqualTo(363L); |
||||
|
assertThat(value.get(ZONE_2_ID).getPerimeterDefinition()) |
||||
|
.isEqualTo(JacksonUtil.fromString(restrictedZoneAttributeKvEntry.getJsonValue().get(), PerimeterDefinition.class)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateEntryWhenNewVersionIsLessThanCurrent() { |
||||
|
BaseAttributeKvEntry differentValueNewVersionIsSet = new BaseAttributeKvEntry(new JsonDataEntry("zone", "[[50.472001, 30.504001], [50.472001, 30.506001], [50.474001, 30.506001], [50.474001, 30.504001]]"), 364L, 154L); |
||||
|
var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, differentValueNewVersionIsSet, ZONE_2_ID, restrictedZoneAttributeKvEntry)); |
||||
|
|
||||
|
assertThat(entry.updateEntry(updated)).isFalse(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateEntryWhenNewTsAndVersionIsGreaterThenCurrentAndValueWasNotChanged() { |
||||
|
BaseAttributeKvEntry newTsAndTheSameValue = new BaseAttributeKvEntry(allowedZoneDataEntry, 364L, 156L); |
||||
|
var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, newTsAndTheSameValue, ZONE_2_ID, restrictedZoneAttributeKvEntry)); |
||||
|
|
||||
|
assertThat(entry.updateEntry(updated)).isTrue(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateEntryWithOldTs() { |
||||
|
BaseAttributeKvEntry oldTsAndTheSameValue = new BaseAttributeKvEntry(allowedZoneDataEntry, 362L, 156L); |
||||
|
var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, oldTsAndTheSameValue, ZONE_2_ID, restrictedZoneAttributeKvEntry)); |
||||
|
|
||||
|
assertThat(entry.updateEntry(updated)).isFalse(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testUpdateEntryWithNewZone() { |
||||
|
final AssetId NEW_ZONE_ID = new AssetId(UUID.fromString("a3eacf1a-6af3-4e9f-87c4-502bb25c7dc3")); |
||||
|
BaseAttributeKvEntry newZone = new BaseAttributeKvEntry(new JsonDataEntry("zone", "[[50.472001, 30.504001], [50.472001, 30.506001], [50.474001, 30.506001], [50.474001, 30.504001]]"), 364L, 156L); |
||||
|
var updated = new GeofencingArgumentEntry(Map.of(ZONE_1_ID, allowedZoneAttributeKvEntry, ZONE_2_ID, restrictedZoneAttributeKvEntry, NEW_ZONE_ID, newZone)); |
||||
|
assertThat(entry.updateEntry(updated)).isTrue(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testIsEmpty() { |
||||
|
GeofencingArgumentEntry geofencingArgumentEntry = new GeofencingArgumentEntry(); |
||||
|
assertThat(geofencingArgumentEntry.isEmpty()).isTrue(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testIsEmptyWithEmptyMap() { |
||||
|
GeofencingArgumentEntry geofencingArgumentEntry = new GeofencingArgumentEntry(Map.of()); |
||||
|
assertThat(geofencingArgumentEntry.isEmpty()).isTrue(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testInvalidKvEntryDataTypeForZoneResultInEmptyArgument() { |
||||
|
BaseAttributeKvEntry invalidZoneEntry = new BaseAttributeKvEntry(new StringDataEntry("zone", "someString"), 363L, 155L); |
||||
|
assertThatThrownBy(() -> new GeofencingArgumentEntry(Map.of(ZONE_1_ID, invalidZoneEntry))) |
||||
|
.isExactlyInstanceOf(IllegalArgumentException.class) |
||||
|
.hasMessage("The given string value cannot be transformed to Json object: someString"); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void testNotParsableToPerimeterJsonKvEntryResultInExceptionTrowed() { |
||||
|
BaseAttributeKvEntry invalidZoneEntry = new BaseAttributeKvEntry(new JsonDataEntry("zone", "\"{}\""), 363L, 155L); |
||||
|
assertThatThrownBy(() -> new GeofencingArgumentEntry(Map.of(ZONE_1_ID, invalidZoneEntry))) |
||||
|
.isExactlyInstanceOf(IllegalArgumentException.class) |
||||
|
.hasMessage("The given string value cannot be transformed to Json object: \"{}\""); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,169 @@ |
|||||
|
/** |
||||
|
* 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 org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.thingsboard.common.util.geo.Coordinates; |
||||
|
import org.thingsboard.server.common.data.id.AssetId; |
||||
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.JsonDataEntry; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingEvalResult; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingZoneState; |
||||
|
|
||||
|
import java.util.UUID; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingPresenceStatus.INSIDE; |
||||
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingPresenceStatus.OUTSIDE; |
||||
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingTransitionEvent.ENTERED; |
||||
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingTransitionEvent.LEFT; |
||||
|
|
||||
|
public class GeofencingZoneStateTest { |
||||
|
|
||||
|
private final AssetId ZONE_ID = new AssetId(UUID.fromString("628730fd-d625-417f-9c6d-ae9fe4addbdb")); |
||||
|
|
||||
|
private GeofencingZoneState state; |
||||
|
|
||||
|
@BeforeEach |
||||
|
void setUp() { |
||||
|
String POLYGON = "[[50.472000, 30.504000], [50.472000, 30.506000], [50.474000, 30.506000], [50.474000, 30.504000]]"; |
||||
|
state = new GeofencingZoneState(ZONE_ID, new BaseAttributeKvEntry(new JsonDataEntry("zone", POLYGON), 100L, 1L)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void evaluate_initialInside_thenInsideAgain() { |
||||
|
var inside = new Coordinates(50.4730, 30.5050); |
||||
|
// first evaluation: no prior state -> ENTERED
|
||||
|
assertThat(state.evaluate(inside)).isEqualTo(new GeofencingEvalResult(ENTERED, INSIDE)); |
||||
|
// same position again -> INSIDE (steady state)
|
||||
|
assertThat(state.evaluate(inside)).isEqualTo(new GeofencingEvalResult(null, INSIDE)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void evaluate_initialOutside_thenOutsideAgain() { |
||||
|
var outside = new Coordinates(50.4760, 30.5110); |
||||
|
// first evaluation: no prior state -> OUTSIDE
|
||||
|
assertThat(state.evaluate(outside)).isEqualTo(new GeofencingEvalResult(null, OUTSIDE)); |
||||
|
// same position again -> OUTSIDE (steady state)
|
||||
|
assertThat(state.evaluate(outside)).isEqualTo(new GeofencingEvalResult(null, OUTSIDE)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void evaluate_inside_thenLeave() { |
||||
|
var inside = new Coordinates(50.4730, 30.5050); |
||||
|
var outside = new Coordinates(50.4760, 30.5110); |
||||
|
// enter
|
||||
|
assertThat(state.evaluate(inside)).isEqualTo(new GeofencingEvalResult(ENTERED, INSIDE)); |
||||
|
// leave -> LEFT
|
||||
|
assertThat(state.evaluate(outside)).isEqualTo(new GeofencingEvalResult(LEFT, OUTSIDE)); |
||||
|
// still outside -> OUTSIDE
|
||||
|
assertThat(state.evaluate(outside)).isEqualTo(new GeofencingEvalResult(null, OUTSIDE)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void evaluate_outside_thenEnter() { |
||||
|
var outside = new Coordinates(50.4760, 30.5110); |
||||
|
var inside = new Coordinates(50.4730, 30.5050); |
||||
|
// start outside
|
||||
|
assertThat(state.evaluate(outside)).isEqualTo(new GeofencingEvalResult(null, OUTSIDE)); |
||||
|
// cross boundary -> ENTERED
|
||||
|
assertThat(state.evaluate(inside)).isEqualTo(new GeofencingEvalResult(ENTERED, INSIDE)); |
||||
|
// remain inside -> INSIDE
|
||||
|
assertThat(state.evaluate(inside)).isEqualTo(new GeofencingEvalResult(null, INSIDE)); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void update_withNewerVersion_updatesState_andResetsPresence() { |
||||
|
// arrange: establish a prior presence to ensure it’s reset on update
|
||||
|
var inside = new Coordinates(50.4730, 30.5050); |
||||
|
assertThat(state.evaluate(inside)).isNotNull(); // sets lastPresence internally
|
||||
|
|
||||
|
String NEW_POLYGON = "[[50.470000, 30.502000], [50.470000, 30.503000], [50.471000, 30.503000], [50.471000, 30.502000]]"; |
||||
|
GeofencingZoneState newer = new GeofencingZoneState( |
||||
|
ZONE_ID, |
||||
|
new BaseAttributeKvEntry(new JsonDataEntry("zone", NEW_POLYGON), 200L, 2L) |
||||
|
); |
||||
|
|
||||
|
// act
|
||||
|
boolean changed = state.update(newer); |
||||
|
|
||||
|
// assert
|
||||
|
assertThat(changed).isTrue(); |
||||
|
assertThat(state.getTs()).isEqualTo(200L); |
||||
|
assertThat(state.getVersion()).isEqualTo(2L); |
||||
|
assertThat(state.getPerimeterDefinition()).isNotNull(); |
||||
|
assertThat(state.getLastPresence()).isNull(); // must be reset on successful update
|
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void update_withEqualVersion_doesNothing() { |
||||
|
// arrange: same version (1L) but different ts/polygon should still be ignored
|
||||
|
String SOME_POLYGON = "[[50.472500, 30.504500], [50.472500, 30.505500], [50.473500, 30.505500], [50.473500, 30.504500]]"; |
||||
|
GeofencingZoneState sameVersion = new GeofencingZoneState( |
||||
|
ZONE_ID, |
||||
|
new BaseAttributeKvEntry(new JsonDataEntry("zone", SOME_POLYGON), 300L, 1L) |
||||
|
); |
||||
|
|
||||
|
// act
|
||||
|
boolean changed = state.update(sameVersion); |
||||
|
|
||||
|
// assert: nothing changes
|
||||
|
assertThat(changed).isFalse(); |
||||
|
assertThat(state.getTs()).isEqualTo(100L); |
||||
|
assertThat(state.getVersion()).isEqualTo(1L); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void update_withNullNewVersion_alwaysApplies_andCopiesNull() { |
||||
|
// arrange: the implementation updates if newVersion == null
|
||||
|
String OTHER_POLYGON = "[[50.471000, 30.506000], [50.471000, 30.507000], [50.472000, 30.507000], [50.472000, 30.506000]]"; |
||||
|
GeofencingZoneState nullVersion = new GeofencingZoneState( |
||||
|
ZONE_ID, |
||||
|
new BaseAttributeKvEntry(new JsonDataEntry("zone", OTHER_POLYGON), 400L, null) |
||||
|
); |
||||
|
|
||||
|
// act
|
||||
|
boolean changed = state.update(nullVersion); |
||||
|
|
||||
|
// assert: applied and version copied as null
|
||||
|
assertThat(changed).isTrue(); |
||||
|
assertThat(state.getTs()).isEqualTo(400L); |
||||
|
assertThat(state.getVersion()).isNull(); |
||||
|
assertThat(state.getLastPresence()).isNull(); |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
void update_withNewVersionWhenExistingIsNull_alwaysApplies_andCopiesNew() { |
||||
|
// arrange: the implementation updates if newVersion == null
|
||||
|
String OTHER_POLYGON = "[[50.471000, 30.506000], [50.471000, 30.507000], [50.472000, 30.507000], [50.472000, 30.506000]]"; |
||||
|
GeofencingZoneState newVersion = new GeofencingZoneState( |
||||
|
ZONE_ID, |
||||
|
new BaseAttributeKvEntry(new JsonDataEntry("zone", OTHER_POLYGON), 400L, 2L) |
||||
|
); |
||||
|
state.setVersion(null); |
||||
|
|
||||
|
// act
|
||||
|
boolean changed = state.update(newVersion); |
||||
|
|
||||
|
// assert: applied and version copied as null
|
||||
|
assertThat(changed).isTrue(); |
||||
|
assertThat(state.getTs()).isEqualTo(400L); |
||||
|
assertThat(state.getVersion()).isEqualTo(2); |
||||
|
assertThat(state.getLastPresence()).isNull(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,83 @@ |
|||||
|
/** |
||||
|
* 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.resource; |
||||
|
|
||||
|
import org.junit.Test; |
||||
|
import org.springframework.beans.factory.annotation.Autowired; |
||||
|
import org.springframework.test.context.bean.override.mockito.MockitoSpyBean; |
||||
|
import org.thingsboard.common.util.JacksonUtil; |
||||
|
import org.thingsboard.server.common.data.GeneralFileDescriptor; |
||||
|
import org.thingsboard.server.common.data.ResourceType; |
||||
|
import org.thingsboard.server.common.data.TbResource; |
||||
|
import org.thingsboard.server.common.data.TbResourceDataInfo; |
||||
|
import org.thingsboard.server.common.data.TbResourceInfo; |
||||
|
import org.thingsboard.server.controller.AbstractControllerTest; |
||||
|
import org.thingsboard.server.dao.resource.ResourceService; |
||||
|
import org.thingsboard.server.dao.resource.TbResourceDataCache; |
||||
|
import org.thingsboard.server.dao.service.DaoSqlTest; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.Mockito.clearInvocations; |
||||
|
import static org.mockito.Mockito.timeout; |
||||
|
import static org.mockito.Mockito.verify; |
||||
|
import static org.mockito.Mockito.verifyNoMoreInteractions; |
||||
|
|
||||
|
@DaoSqlTest |
||||
|
public class DefaultResourceDataCacheTest extends AbstractControllerTest { |
||||
|
|
||||
|
@MockitoSpyBean |
||||
|
private ResourceService resourceService; |
||||
|
@Autowired |
||||
|
private TbResourceService tbResourceService; |
||||
|
@MockitoSpyBean |
||||
|
private TbResourceDataCache resourceDataCache; |
||||
|
|
||||
|
@Test |
||||
|
public void testGetCachedResourceData() throws Exception { |
||||
|
loginTenantAdmin(); |
||||
|
|
||||
|
TbResource resource = new TbResource(); |
||||
|
resource.setTenantId(tenantId); |
||||
|
resource.setTitle("File for AI request"); |
||||
|
resource.setResourceType(ResourceType.GENERAL); |
||||
|
resource.setFileName("myTestJson.json"); |
||||
|
GeneralFileDescriptor descriptor = new GeneralFileDescriptor("application/json"); |
||||
|
resource.setDescriptorValue(descriptor); |
||||
|
byte[] data = "This is a test prompt for AI request.".getBytes(); |
||||
|
resource.setData(data); |
||||
|
TbResourceInfo savedResource = tbResourceService.save(resource); |
||||
|
verify(resourceDataCache, timeout(2000).times(1)).evictResourceData(tenantId, savedResource.getId()); |
||||
|
|
||||
|
TbResourceDataInfo cachedData = resourceDataCache.getResourceDataInfoAsync(tenantId, savedResource.getId()).get(); |
||||
|
assertThat(cachedData.getData()).isEqualTo(data); |
||||
|
assertThat(JacksonUtil.treeToValue(cachedData.getDescriptor(), GeneralFileDescriptor.class)).isEqualTo(descriptor); |
||||
|
verify(resourceService).getResourceDataInfo(tenantId, savedResource.getId()); |
||||
|
|
||||
|
// retrieve resource data second time
|
||||
|
clearInvocations(resourceService); |
||||
|
TbResourceDataInfo cachedData2 = resourceDataCache.getResourceDataInfoAsync(tenantId, savedResource.getId()).get(); |
||||
|
assertThat(cachedData2.getData()).isEqualTo(data); |
||||
|
verifyNoMoreInteractions(resourceService); |
||||
|
|
||||
|
// delete resource, check cache
|
||||
|
TbResource resourceById = resourceService.findResourceById(tenantId, savedResource.getId()); |
||||
|
tbResourceService.delete(resourceById, true, null); |
||||
|
verify(resourceDataCache, timeout(2000).times(2)).evictResourceData(tenantId, savedResource.getId()); |
||||
|
TbResourceDataInfo cachedDataAfterDeletion = resourceDataCache.getResourceDataInfoAsync(tenantId, savedResource.getId()).get(); |
||||
|
assertThat(cachedDataAfterDeletion).isEqualTo(null); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,109 @@ |
|||||
|
/** |
||||
|
* 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.utils; |
||||
|
|
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.junit.jupiter.api.extension.ExtendWith; |
||||
|
import org.mockito.junit.jupiter.MockitoExtension; |
||||
|
import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingPresenceStatus; |
||||
|
import org.thingsboard.server.common.data.id.AssetId; |
||||
|
import org.thingsboard.server.common.data.id.CalculatedFieldId; |
||||
|
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.kv.BaseAttributeKvEntry; |
||||
|
import org.thingsboard.server.common.data.kv.JsonDataEntry; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldStateProto; |
||||
|
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.geofencing.GeofencingArgumentEntry; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; |
||||
|
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingZoneState; |
||||
|
|
||||
|
import java.util.LinkedHashMap; |
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
import java.util.UUID; |
||||
|
|
||||
|
import static org.assertj.core.api.Assertions.assertThat; |
||||
|
import static org.mockito.BDDMockito.given; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto; |
||||
|
|
||||
|
@ExtendWith(MockitoExtension.class) |
||||
|
class CalculatedFieldUtilsTest { |
||||
|
|
||||
|
private static final TenantId TENANT_ID = TenantId.fromUUID(UUID.fromString("0a69e1e2-fcbc-4234-a4cd-3844bf54035c")); |
||||
|
private static final CalculatedFieldId CF_ID = CalculatedFieldId.fromString("ec0e91b9-6f27-4e93-946a-5fbc2707d8bc"); |
||||
|
private static final DeviceId DEVICE_ID = DeviceId.fromString("1e03bd38-2010-4739-9362-160c288e36c4"); |
||||
|
|
||||
|
@Test |
||||
|
void toProtoAndFromProto_shouldMapGeofencingArgumentsAndZones() { |
||||
|
// given
|
||||
|
CalculatedFieldEntityCtxId stateId = mock(CalculatedFieldEntityCtxId.class); |
||||
|
given(stateId.tenantId()).willReturn(TENANT_ID); |
||||
|
given(stateId.cfId()).willReturn(CF_ID); |
||||
|
given(stateId.entityId()).willReturn(DEVICE_ID); |
||||
|
|
||||
|
// Build a geofencing argument with two zones (one with inside=true, one with inside=null)
|
||||
|
GeofencingArgumentEntry geofencingArgumentEntry = new GeofencingArgumentEntry(); |
||||
|
Map<EntityId, GeofencingZoneState> zoneStates = new LinkedHashMap<>(); |
||||
|
|
||||
|
UUID zoneId1 = UUID.fromString("624a8fff-71a2-4847-a100-ff1cf52dbe71"); |
||||
|
UUID zoneId2 = UUID.fromString("e2adf6ce-9478-40b1-b0e9-4a6860cc46bb"); |
||||
|
|
||||
|
AssetId z1 = new AssetId(zoneId1); |
||||
|
AssetId z2 = new AssetId(zoneId2); |
||||
|
|
||||
|
JsonDataEntry zone1 = new JsonDataEntry("zone", "[[50.472000, 30.504000], [50.472000, 30.506000], [50.474000, 30.506000], [50.474000, 30.504000]]"); |
||||
|
JsonDataEntry zone2 = new JsonDataEntry("zone", "[[50.475000, 30.510000], [50.475000, 30.512000], [50.477000, 30.512000], [50.477000, 30.510000]]"); |
||||
|
|
||||
|
BaseAttributeKvEntry zone1PerimeterAttribute = new BaseAttributeKvEntry(zone1, System.currentTimeMillis(), 0L); |
||||
|
BaseAttributeKvEntry zone2PerimeterAttribute = new BaseAttributeKvEntry(zone2, System.currentTimeMillis(), 0L); |
||||
|
|
||||
|
GeofencingZoneState s1 = new GeofencingZoneState(z1, zone1PerimeterAttribute); |
||||
|
s1.setLastPresence(GeofencingPresenceStatus.INSIDE); |
||||
|
GeofencingZoneState s2 = new GeofencingZoneState(z2, zone2PerimeterAttribute); |
||||
|
|
||||
|
zoneStates.put(z1, s1); |
||||
|
zoneStates.put(z2, s2); |
||||
|
geofencingArgumentEntry.setZoneStates(zoneStates); |
||||
|
|
||||
|
// Create cf state with the geofencing argument and add it to the state map
|
||||
|
CalculatedFieldState state = new GeofencingCalculatedFieldState(List.of("geofencingArgumentTest")); |
||||
|
state.updateState(mock(CalculatedFieldCtx.class), Map.of("geofencingArgumentTest", geofencingArgumentEntry)); |
||||
|
|
||||
|
// when
|
||||
|
CalculatedFieldStateProto proto = toProto(stateId, state); |
||||
|
|
||||
|
// then
|
||||
|
CalculatedFieldState fromProto = CalculatedFieldUtils.fromProto(proto); |
||||
|
assertThat(fromProto) |
||||
|
.usingRecursiveComparison() |
||||
|
.ignoringFields("requiredArguments") |
||||
|
.isEqualTo(state); |
||||
|
|
||||
|
ArgumentEntry fromProtoArgument = fromProto.getArguments().get("geofencingArgumentTest"); |
||||
|
assertThat(fromProtoArgument).isInstanceOf(GeofencingArgumentEntry.class); |
||||
|
GeofencingArgumentEntry fromProtoGeoArgument = (GeofencingArgumentEntry) fromProtoArgument; |
||||
|
assertThat(fromProtoGeoArgument.getZoneStates()).hasSize(2); |
||||
|
assertThat(fromProtoGeoArgument.getZoneStates().get(z1).getLastPresence()).isEqualTo(GeofencingPresenceStatus.INSIDE); |
||||
|
assertThat(fromProtoGeoArgument.getZoneStates().get(z2).getLastPresence()).isNull(); |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,28 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2025 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.dao.resource; |
||||
|
|
||||
|
import com.google.common.util.concurrent.FluentFuture; |
||||
|
import org.thingsboard.server.common.data.TbResourceDataInfo; |
||||
|
import org.thingsboard.server.common.data.id.TbResourceId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
|
||||
|
public interface TbResourceDataCache { |
||||
|
|
||||
|
FluentFuture<TbResourceDataInfo> getResourceDataInfoAsync(TenantId tenantId, TbResourceId resourceId); |
||||
|
|
||||
|
void evictResourceData(TenantId tenantId, TbResourceId resourceId); |
||||
|
} |
||||
@ -0,0 +1,29 @@ |
|||||
|
/** |
||||
|
* 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; |
||||
|
|
||||
|
import lombok.AllArgsConstructor; |
||||
|
import lombok.Data; |
||||
|
import lombok.EqualsAndHashCode; |
||||
|
import lombok.NoArgsConstructor; |
||||
|
|
||||
|
@Data |
||||
|
@EqualsAndHashCode |
||||
|
@AllArgsConstructor |
||||
|
@NoArgsConstructor |
||||
|
public class GeneralFileDescriptor { |
||||
|
private String mediaType; |
||||
|
} |
||||
@ -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; |
||||
|
|
||||
|
import com.fasterxml.jackson.databind.JsonNode; |
||||
|
import lombok.AllArgsConstructor; |
||||
|
import lombok.Data; |
||||
|
import lombok.NoArgsConstructor; |
||||
|
|
||||
|
@Data |
||||
|
@AllArgsConstructor |
||||
|
@NoArgsConstructor |
||||
|
public class TbResourceDataInfo { |
||||
|
|
||||
|
private byte[] data; |
||||
|
private JsonNode descriptor; |
||||
|
|
||||
|
} |
||||
@ -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.common.data.ai.model.chat; |
||||
|
|
||||
|
import dev.langchain4j.model.chat.ChatModel; |
||||
|
import jakarta.validation.Valid; |
||||
|
import jakarta.validation.constraints.Max; |
||||
|
import jakarta.validation.constraints.NotBlank; |
||||
|
import jakarta.validation.constraints.NotNull; |
||||
|
import jakarta.validation.constraints.Positive; |
||||
|
import jakarta.validation.constraints.PositiveOrZero; |
||||
|
import lombok.Builder; |
||||
|
import lombok.With; |
||||
|
import org.thingsboard.server.common.data.ai.provider.AiProvider; |
||||
|
import org.thingsboard.server.common.data.ai.provider.OllamaProviderConfig; |
||||
|
|
||||
|
@Builder |
||||
|
public record OllamaChatModelConfig( |
||||
|
@NotNull @Valid OllamaProviderConfig providerConfig, |
||||
|
@NotBlank String modelId, |
||||
|
@PositiveOrZero Double temperature, |
||||
|
@Positive @Max(1) Double topP, |
||||
|
@PositiveOrZero Integer topK, |
||||
|
Integer contextLength, |
||||
|
Integer maxOutputTokens, |
||||
|
@With @Positive Integer timeoutSeconds, |
||||
|
@With @PositiveOrZero Integer maxRetries |
||||
|
) implements AiChatModelConfig<OllamaChatModelConfig> { |
||||
|
|
||||
|
@Override |
||||
|
public AiProvider provider() { |
||||
|
return AiProvider.OLLAMA; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public ChatModel configure(Langchain4jChatModelConfigurer configurer) { |
||||
|
return configurer.configureChatModel(this); |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public boolean supportsJsonMode() { |
||||
|
return true; |
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,48 @@ |
|||||
|
/** |
||||
|
* 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.ai.provider; |
||||
|
|
||||
|
import com.fasterxml.jackson.annotation.JsonSubTypes; |
||||
|
import com.fasterxml.jackson.annotation.JsonTypeInfo; |
||||
|
import jakarta.validation.Valid; |
||||
|
import jakarta.validation.constraints.NotNull; |
||||
|
|
||||
|
public record OllamaProviderConfig( |
||||
|
@NotNull String baseUrl, |
||||
|
@NotNull @Valid OllamaAuth auth |
||||
|
) implements AiProviderConfig { |
||||
|
|
||||
|
@JsonTypeInfo( |
||||
|
use = JsonTypeInfo.Id.NAME, |
||||
|
include = JsonTypeInfo.As.PROPERTY, |
||||
|
property = "type" |
||||
|
) |
||||
|
@JsonSubTypes({ |
||||
|
@JsonSubTypes.Type(value = OllamaAuth.None.class, name = "NONE"), |
||||
|
@JsonSubTypes.Type(value = OllamaAuth.Basic.class, name = "BASIC"), |
||||
|
@JsonSubTypes.Type(value = OllamaAuth.Token.class, name = "TOKEN") |
||||
|
}) |
||||
|
public sealed interface OllamaAuth { |
||||
|
|
||||
|
record None() implements OllamaAuth {} |
||||
|
|
||||
|
record Basic(@NotNull String username, @NotNull String password) implements OllamaAuth {} |
||||
|
|
||||
|
record Token(@NotNull String token) implements OllamaAuth {} |
||||
|
|
||||
|
} |
||||
|
|
||||
|
} |
||||
@ -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.common.data.cf.configuration; |
||||
|
|
||||
|
import java.util.Map; |
||||
|
|
||||
|
public interface ArgumentsBasedCalculatedFieldConfiguration extends CalculatedFieldConfiguration { |
||||
|
|
||||
|
Map<String, Argument> getArguments(); |
||||
|
|
||||
|
} |
||||
@ -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 = RelationPathQueryDynamicSourceConfiguration.class, name = "RELATION_PATH_QUERY") |
||||
|
}) |
||||
|
@JsonIgnoreProperties(ignoreUnknown = true) |
||||
|
public interface CfArgumentDynamicSourceConfiguration { |
||||
|
|
||||
|
@JsonIgnore |
||||
|
CFArgumentDynamicSourceType getType(); |
||||
|
|
||||
|
default void validate() {} |
||||
|
|
||||
|
} |
||||
@ -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.common.data.cf.configuration; |
||||
|
|
||||
|
public interface ExpressionBasedCalculatedFieldConfiguration extends ArgumentsBasedCalculatedFieldConfiguration { |
||||
|
|
||||
|
String getExpression(); |
||||
|
|
||||
|
void setExpression(String expression); |
||||
|
|
||||
|
} |
||||
@ -0,0 +1,71 @@ |
|||||
|
/** |
||||
|
* 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 lombok.Data; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.relation.EntityRelation; |
||||
|
import org.thingsboard.server.common.data.relation.EntityRelationPathQuery; |
||||
|
import org.thingsboard.server.common.data.relation.EntitySearchDirection; |
||||
|
import org.thingsboard.server.common.data.relation.RelationPathLevel; |
||||
|
import org.thingsboard.server.common.data.util.CollectionsUtil; |
||||
|
|
||||
|
import java.util.List; |
||||
|
import java.util.NoSuchElementException; |
||||
|
|
||||
|
@Data |
||||
|
public class RelationPathQueryDynamicSourceConfiguration implements CfArgumentDynamicSourceConfiguration { |
||||
|
|
||||
|
private List<RelationPathLevel> levels; |
||||
|
|
||||
|
@Override |
||||
|
public CFArgumentDynamicSourceType getType() { |
||||
|
return CFArgumentDynamicSourceType.RELATION_PATH_QUERY; |
||||
|
} |
||||
|
|
||||
|
@Override |
||||
|
public void validate() { |
||||
|
if (CollectionsUtil.isEmpty(levels)) { |
||||
|
throw new IllegalArgumentException("At least one relation level must be specified!"); |
||||
|
} |
||||
|
levels.forEach(RelationPathLevel::validate); |
||||
|
} |
||||
|
|
||||
|
public List<EntityId> resolveEntityIds(List<EntityRelation> relations) { |
||||
|
EntitySearchDirection lastLevelDirection = getLastLevel().direction(); |
||||
|
return switch (lastLevelDirection) { |
||||
|
case FROM -> relations.stream().map(EntityRelation::getTo).toList(); |
||||
|
case TO -> relations.stream().map(EntityRelation::getFrom).toList(); |
||||
|
}; |
||||
|
} |
||||
|
|
||||
|
public void validateMaxRelationLevel(String argumentName, int maxAllowedRelationLevel) { |
||||
|
if (levels.size() > maxAllowedRelationLevel) { |
||||
|
throw new IllegalArgumentException("Max relation level is greater than configured " + |
||||
|
"maximum allowed relation level in tenant profile: " + maxAllowedRelationLevel + " for argument: " + argumentName); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
public EntityRelationPathQuery toRelationPathQuery(EntityId entityId) { |
||||
|
return new EntityRelationPathQuery(entityId, levels); |
||||
|
} |
||||
|
|
||||
|
private RelationPathLevel getLastLevel() { |
||||
|
return levels.get(levels.size() - 1); |
||||
|
} |
||||
|
|
||||
|
} |
||||
Some files were not shown because too many files changed in this diff
Loading…
Reference in new issue