|
|
|
@ -36,8 +36,6 @@ import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.Argument; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.CFArgumentDynamicSourceType; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.GeofencingCalculatedFieldConfiguration; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.GeofencingZoneGroupConfiguration; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.OutputType; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.RelationQueryDynamicSourceConfiguration; |
|
|
|
import org.thingsboard.server.common.data.id.CalculatedFieldId; |
|
|
|
@ -131,18 +129,18 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP |
|
|
|
|
|
|
|
@Override |
|
|
|
public ListenableFuture<CalculatedFieldState> fetchStateFromDb(CalculatedFieldCtx ctx, EntityId entityId) { |
|
|
|
Map<String, ListenableFuture<ArgumentEntry>> argFutures = new HashMap<>(); |
|
|
|
|
|
|
|
if (ctx.getCalculatedField().getType().equals(CalculatedFieldType.GEOFENCING)) { |
|
|
|
fetchGeofencingCalculatedFieldArguments(ctx, entityId, argFutures, false); |
|
|
|
} else { |
|
|
|
for (var entry : ctx.getArguments().entrySet()) { |
|
|
|
var argEntityId = resolveEntityId(entityId, entry); |
|
|
|
var argValueFuture = fetchKvEntry(ctx.getTenantId(), argEntityId, entry.getValue()); |
|
|
|
argFutures.put(entry.getKey(), argValueFuture); |
|
|
|
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); |
|
|
|
var argValueFuture = fetchKvEntry(ctx.getTenantId(), argEntityId, entry.getValue()); |
|
|
|
futures.put(entry.getKey(), argValueFuture); |
|
|
|
} |
|
|
|
yield futures; |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
}; |
|
|
|
return Futures.whenAllComplete(argFutures.values()).call(() -> { |
|
|
|
var result = createStateByType(ctx); |
|
|
|
result.updateState(ctx, resolveArgumentFutures(argFutures)); |
|
|
|
@ -156,14 +154,11 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP |
|
|
|
if (!ctx.getCalculatedField().getType().equals(CalculatedFieldType.GEOFENCING)) { |
|
|
|
return Map.of(); |
|
|
|
} |
|
|
|
Map<String, ListenableFuture<ArgumentEntry>> argFutures = new HashMap<>(); |
|
|
|
fetchGeofencingCalculatedFieldArguments(ctx, entityId, argFutures, true); |
|
|
|
return resolveArgumentFutures(argFutures); |
|
|
|
return resolveArgumentFutures(fetchGeofencingCalculatedFieldArguments(ctx, entityId, true)); |
|
|
|
} |
|
|
|
|
|
|
|
private void fetchGeofencingCalculatedFieldArguments(CalculatedFieldCtx ctx, EntityId entityId, Map<String, ListenableFuture<ArgumentEntry>> argFutures, boolean dynamicArgumentsOnly) { |
|
|
|
var configuration = (GeofencingCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); |
|
|
|
var zoneGroupConfigs = configuration.getGeofencingZoneGroupConfigurations(); |
|
|
|
private Map<String, ListenableFuture<ArgumentEntry>> fetchGeofencingCalculatedFieldArguments(CalculatedFieldCtx ctx, EntityId entityId, boolean dynamicArgumentsOnly) { |
|
|
|
Map<String, ListenableFuture<ArgumentEntry>> argFutures = new HashMap<>(); |
|
|
|
Set<Entry<String, Argument>> entries = ctx.getArguments().entrySet(); |
|
|
|
if (dynamicArgumentsOnly) { |
|
|
|
entries = entries.stream() |
|
|
|
@ -175,13 +170,13 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP |
|
|
|
case ENTITY_ID_LATITUDE_ARGUMENT_KEY, ENTITY_ID_LONGITUDE_ARGUMENT_KEY -> |
|
|
|
argFutures.put(entry.getKey(), fetchKvEntry(ctx.getTenantId(), resolveEntityId(entityId, entry), entry.getValue())); |
|
|
|
default -> { |
|
|
|
var zoneGroupConfiguration = zoneGroupConfigs.get(entry.getKey()); |
|
|
|
var resolvedEntityIdsFuture = resolveGeofencingEntityIds(ctx.getTenantId(), entityId, entry); |
|
|
|
argFutures.put(entry.getKey(), Futures.transformAsync(resolvedEntityIdsFuture, resolvedEntityIds -> |
|
|
|
fetchGeofencingKvEntry(ctx.getTenantId(), resolvedEntityIds, entry.getValue(), zoneGroupConfiguration), calculatedFieldCallbackExecutor)); |
|
|
|
fetchGeofencingKvEntry(ctx.getTenantId(), resolvedEntityIds, entry.getValue()), calculatedFieldCallbackExecutor)); |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
return argFutures; |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -321,12 +316,11 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP |
|
|
|
}; |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchGeofencingKvEntry(TenantId tenantId, List<EntityId> geofencingEntities, |
|
|
|
Argument argument, GeofencingZoneGroupConfiguration zoneGroupConfiguration) { |
|
|
|
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() |
|
|
|
List<ListenableFuture<Entry<EntityId, AttributeKvEntry>>> kvFutures = geofencingEntities.stream() |
|
|
|
.map(entityId -> { |
|
|
|
var attributesFuture = attributesService.find( |
|
|
|
tenantId, |
|
|
|
@ -341,10 +335,10 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP |
|
|
|
); |
|
|
|
}).collect(Collectors.toList()); |
|
|
|
|
|
|
|
ListenableFuture<List<Map.Entry<EntityId, AttributeKvEntry>>> allFutures = Futures.allAsList(kvFutures); |
|
|
|
ListenableFuture<List<Entry<EntityId, AttributeKvEntry>>> allFutures = Futures.allAsList(kvFutures); |
|
|
|
|
|
|
|
return Futures.transform(allFutures, entries -> ArgumentEntry.createGeofencingValueArgument(entries.stream() |
|
|
|
.collect(Collectors.toMap(Entry::getKey, Entry::getValue)), zoneGroupConfiguration), |
|
|
|
.collect(Collectors.toMap(Entry::getKey, Entry::getValue))), |
|
|
|
calculatedFieldCallbackExecutor |
|
|
|
); |
|
|
|
} |
|
|
|
|