From 2c06aa475f3241bc474a9ce8aee46490eec02bc3 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Tue, 30 Sep 2025 16:26:46 +0300 Subject: [PATCH 1/6] geofencing cf bugfixes --- .../main/data/upgrade/basic/schema_update.sql | 2 +- ...ulatedFieldDynamicArgumentsRefreshMsg.java | 35 -------- .../CalculatedFieldEntityActor.java | 3 - ...CalculatedFieldEntityMessageProcessor.java | 50 +++++------ .../CalculatedFieldManagerActor.java | 3 - ...alculatedFieldManagerMessageProcessor.java | 62 ------------- ...ulatedFieldDynamicArgumentsRefreshMsg.java | 37 -------- ...tractCalculatedFieldProcessingService.java | 4 + .../ctx/state/BaseCalculatedFieldState.java | 4 +- .../cf/ctx/state/CalculatedFieldCtx.java | 42 +++++++-- .../cf/ctx/state/CalculatedFieldState.java | 7 -- .../ctx/state/SimpleCalculatedFieldState.java | 5 +- .../geofencing/GeofencingArgumentEntry.java | 10 ++- .../GeofencingCalculatedFieldState.java | 87 ++++++++----------- .../cf/CalculatedFieldIntegrationTest.java | 10 ++- .../server/controller/AbstractWebTest.java | 12 +-- .../GeofencingCalculatedFieldStateTest.java | 6 +- .../state/SimpleCalculatedFieldStateTest.java | 3 +- ...SupportedCalculatedFieldConfiguration.java | 3 - ...eofencingCalculatedFieldConfiguration.java | 7 +- .../DefaultTenantProfileConfiguration.java | 2 +- ...ncingCalculatedFieldConfigurationTest.java | 27 ------ .../server/common/msg/MsgType.java | 5 +- .../service/CalculatedFieldServiceTest.java | 72 ++++----------- 24 files changed, 148 insertions(+), 350 deletions(-) delete mode 100644 application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldDynamicArgumentsRefreshMsg.java delete mode 100644 application/src/main/java/org/thingsboard/server/actors/calculatedField/EntityCalculatedFieldDynamicArgumentsRefreshMsg.java diff --git a/application/src/main/data/upgrade/basic/schema_update.sql b/application/src/main/data/upgrade/basic/schema_update.sql index 320d3e5bdd..0add4c0545 100644 --- a/application/src/main/data/upgrade/basic/schema_update.sql +++ b/application/src/main/data/upgrade/basic/schema_update.sql @@ -27,7 +27,7 @@ SET profile_data = jsonb_set( CASE WHEN (profile_data -> 'configuration') ? 'minAllowedScheduledUpdateIntervalInSecForCF' THEN NULL - ELSE to_jsonb(3600) + ELSE to_jsonb(60) END, 'maxRelationLevelPerCfArgument', CASE diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldDynamicArgumentsRefreshMsg.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldDynamicArgumentsRefreshMsg.java deleted file mode 100644 index 301fe22dfb..0000000000 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldDynamicArgumentsRefreshMsg.java +++ /dev/null @@ -1,35 +0,0 @@ -/** - * Copyright © 2016-2025 The Thingsboard Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.thingsboard.server.actors.calculatedField; - -import lombok.Data; -import org.thingsboard.server.common.data.id.CalculatedFieldId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.msg.MsgType; -import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; - -@Data -public class CalculatedFieldDynamicArgumentsRefreshMsg implements ToCalculatedFieldSystemMsg { - - private final TenantId tenantId; - private final CalculatedFieldId cfId; - - @Override - public MsgType getMsgType() { - return MsgType.CF_DYNAMIC_ARGUMENTS_REFRESH_MSG; - } - -} diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java index 2a5f3c3cfd..c57984ef3d 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java @@ -75,9 +75,6 @@ public class CalculatedFieldEntityActor extends AbstractCalculatedFieldActor { case CF_LINKED_TELEMETRY_MSG: processor.process((EntityCalculatedFieldLinkedTelemetryMsg) msg); break; - case CF_ENTITY_DYNAMIC_ARGUMENTS_REFRESH_MSG: - processor.process((EntityCalculatedFieldDynamicArgumentsRefreshMsg) msg); - break; default: return false; } diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java index 7513ca41e2..f8b61a082f 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java @@ -49,6 +49,7 @@ 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.SingleValueArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry; +import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; import java.util.ArrayList; import java.util.Collection; @@ -227,18 +228,6 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } } - public void process(EntityCalculatedFieldDynamicArgumentsRefreshMsg msg) throws CalculatedFieldException { - log.debug("[{}][{}] Processing CF dynamic arguments refresh msg.", entityId, msg.getCfId()); - CalculatedFieldState currentState = states.get(msg.getCfId()); - if (currentState == null) { - log.debug("[{}][{}] Failed to find CF state for entity.", entityId, msg.getCfId()); - } else { - currentState.setDirty(true); - log.debug("[{}][{}] CF state marked as dirty.", entityId, msg.getCfId()); - } - msg.getCallback().onSuccess(); - } - private void processTelemetry(CalculatedFieldCtx ctx, CalculatedFieldTelemetryMsgProto proto, List cfIdList, MultipleTbCallback callback) throws CalculatedFieldException { processArgumentValuesUpdate(ctx, cfIdList, callback, mapToArguments(ctx, proto.getTsDataList()), toTbMsgId(proto), toTbMsgType(proto)); } @@ -266,12 +255,13 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM if (state == null) { state = getOrInitState(ctx); justRestored = true; - } else if (state.isDirty()) { - log.debug("[{}][{}] Going to update dirty CF state.", entityId, ctx.getCfId()); + } else if (ctx.shouldFetchDynamicArgumentsFromDb(state)) { + log.debug("[{}][{}] Going to update dynamic arguments for CF.", entityId, ctx.getCfId()); try { Map dynamicArgsFromDb = cfService.fetchDynamicArgsFromDb(ctx, entityId); dynamicArgsFromDb.forEach(newArgValues::putIfAbsent); - state.setDirty(false); + var geofencingState = (GeofencingCalculatedFieldState) state; + geofencingState.setLastDynamicArgumentsRefreshTs(System.currentTimeMillis()); } catch (Exception e) { throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).cause(e).build(); } @@ -403,7 +393,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM return mapToArguments(entityId, argNames, geofencingArgumentNames, scope, attrDataList); } - private Map mapToArguments(EntityId entityId, Map argNames, List geoArgNames, AttributeScopeProto scope, List attrDataList) { + private Map mapToArguments(EntityId entityId, Map argNames, List geofencingArgNames, AttributeScopeProto scope, List attrDataList) { Map arguments = new HashMap<>(); for (AttributeValueProto item : attrDataList) { ReferencedEntityKey key = new ReferencedEntityKey(item.getKey(), ArgumentType.ATTRIBUTE, AttributeScope.valueOf(scope.name())); @@ -411,7 +401,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM if (argName == null) { continue; } - if (geoArgNames.contains(argName)) { + if (geofencingArgNames.contains(argName)) { arguments.put(argName, new GeofencingArgumentEntry(entityId, item)); continue; } @@ -425,26 +415,32 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM if (argNames.isEmpty()) { return Collections.emptyMap(); } - return mapToArgumentsWithDefaultValue(argNames, ctx.getArguments(), scope, removedAttrKeys); + List geofencingArgumentNames = ctx.getLinkedEntityGeofencingArgumentNames(); + return mapToArgumentsWithDefaultValue(argNames, ctx.getArguments(), geofencingArgumentNames, scope, removedAttrKeys); } private Map mapToArgumentsWithDefaultValue(CalculatedFieldCtx ctx, AttributeScopeProto scope, List removedAttrKeys) { - return mapToArgumentsWithDefaultValue(ctx.getMainEntityArguments(), ctx.getArguments(), scope, removedAttrKeys); + return mapToArgumentsWithDefaultValue(ctx.getMainEntityArguments(), ctx.getArguments(), ctx.getMainEntityGeofencingArgumentNames(), scope, removedAttrKeys); } - private Map mapToArgumentsWithDefaultValue(Map argNames, Map configArguments, AttributeScopeProto scope, List removedAttrKeys) { + private Map mapToArgumentsWithDefaultValue(Map argNames, Map configArguments, List geofencingArgNames, AttributeScopeProto scope, List removedAttrKeys) { Map arguments = new HashMap<>(); for (String removedKey : removedAttrKeys) { ReferencedEntityKey key = new ReferencedEntityKey(removedKey, ArgumentType.ATTRIBUTE, AttributeScope.valueOf(scope.name())); String argName = argNames.get(key); - if (argName != null) { - Argument argument = configArguments.get(argName); - String defaultValue = (argument != null) ? argument.getDefaultValue() : null; - arguments.put(argName, StringUtils.isNotEmpty(defaultValue) - ? new SingleValueArgumentEntry(System.currentTimeMillis(), new StringDataEntry(removedKey, defaultValue), null) - : new SingleValueArgumentEntry()); - + if (argName == null) { + continue; } + if (geofencingArgNames.contains(argName)) { + arguments.put(argName, new GeofencingArgumentEntry()); + continue; + } + Argument argument = configArguments.get(argName); + String defaultValue = (argument != null) ? argument.getDefaultValue() : null; + arguments.put(argName, StringUtils.isNotEmpty(defaultValue) + ? new SingleValueArgumentEntry(System.currentTimeMillis(), new StringDataEntry(removedKey, defaultValue), null) + : new SingleValueArgumentEntry()); + } return arguments; } diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerActor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerActor.java index ab6cb34176..7d2ae0ff44 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerActor.java @@ -79,9 +79,6 @@ public class CalculatedFieldManagerActor extends AbstractCalculatedFieldActor { case CF_LINKED_TELEMETRY_MSG: processor.onLinkedTelemetryMsg((CalculatedFieldLinkedTelemetryMsg) msg); break; - case CF_DYNAMIC_ARGUMENTS_REFRESH_MSG: - processor.onDynamicArgumentsRefreshMsg((CalculatedFieldDynamicArgumentsRefreshMsg) msg); - break; default: return false; } diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java index 75ca0b4c9b..7a76cb9821 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java @@ -27,7 +27,6 @@ import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.ProfileEntityIdInfo; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.CalculatedFieldLink; -import org.thingsboard.server.common.data.cf.configuration.ScheduledUpdateSupportedCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.DeviceId; @@ -57,10 +56,7 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; -import java.util.concurrent.ScheduledFuture; -import java.util.concurrent.TimeUnit; import java.util.function.BiConsumer; import static org.thingsboard.server.utils.CalculatedFieldUtils.fromProto; @@ -74,7 +70,6 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware private final Map calculatedFields = new HashMap<>(); private final Map> entityIdCalculatedFields = new HashMap<>(); private final Map> entityIdCalculatedFieldLinks = new HashMap<>(); - private final Map> cfDynamicArgumentsRefreshTasks = new ConcurrentHashMap<>(); private final CalculatedFieldProcessingService cfExecService; private final CalculatedFieldStateService cfStateService; @@ -113,8 +108,6 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware calculatedFields.clear(); entityIdCalculatedFields.clear(); entityIdCalculatedFieldLinks.clear(); - cfDynamicArgumentsRefreshTasks.values().forEach(future -> future.cancel(true)); - cfDynamicArgumentsRefreshTasks.clear(); ctx.stop(ctx.getSelf()); } @@ -274,7 +267,6 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware // Alternative approach would be to use any list but avoid modifications to the list (change the complete map value instead) entityIdCalculatedFields.computeIfAbsent(cf.getEntityId(), id -> new CopyOnWriteArrayList<>()).add(cfCtx); addLinks(cf); - scheduleDynamicArgumentsRefreshTaskForCfIfNeeded(cfCtx); applyToTargetCfEntityActors(cfCtx, callback, (id, cb) -> initCfForEntity(id, cfCtx, false, cb)); } } @@ -304,12 +296,6 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware calculatedFields.put(newCf.getId(), newCfCtx); List oldCfList = entityIdCalculatedFields.get(newCf.getEntityId()); - boolean hasSchedulingConfigChanges = newCfCtx.hasSchedulingConfigChanges(oldCfCtx); - if (hasSchedulingConfigChanges) { - cancelCfDynamicArgumentsRefreshTaskIfExists(cfId, false); - scheduleDynamicArgumentsRefreshTaskForCfIfNeeded(newCfCtx); - } - List newCfList = new CopyOnWriteArrayList<>(); boolean found = false; for (CalculatedFieldCtx oldCtx : oldCfList) { @@ -350,19 +336,9 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware } entityIdCalculatedFields.get(cfCtx.getEntityId()).remove(cfCtx); deleteLinks(cfCtx); - cancelCfDynamicArgumentsRefreshTaskIfExists(cfId, true); applyToTargetCfEntityActors(cfCtx, callback, (id, cb) -> deleteCfForEntity(id, cfId, cb)); } - private void cancelCfDynamicArgumentsRefreshTaskIfExists(CalculatedFieldId cfId, boolean cfDeleted) { - var existingTask = cfDynamicArgumentsRefreshTasks.remove(cfId); - if (existingTask != null) { - existingTask.cancel(false); - String reason = cfDeleted ? "deletion" : "update"; - log.debug("[{}][{}] Cancelled dynamic arguments refresh task due to CF {}!", tenantId, cfId, reason); - } - } - public void onTelemetryMsg(CalculatedFieldTelemetryMsg msg) { EntityId entityId = msg.getEntityId(); log.debug("Received telemetry msg from entity [{}]", entityId); @@ -442,43 +418,6 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware return result; } - private void scheduleDynamicArgumentsRefreshTaskForCfIfNeeded(CalculatedFieldCtx cfCtx) { - CalculatedField cf = cfCtx.getCalculatedField(); - if (!(cf.getConfiguration() instanceof ScheduledUpdateSupportedCalculatedFieldConfiguration scheduledCfConfig)) { - return; - } - if (!scheduledCfConfig.isScheduledUpdateEnabled()) { - return; - } - if (cfDynamicArgumentsRefreshTasks.containsKey(cf.getId())) { - log.debug("[{}][{}] Dynamic arguments refresh task for CF already exists!", tenantId, cf.getId()); - return; - } - long refreshDynamicSourceInterval = TimeUnit.SECONDS.toMillis(scheduledCfConfig.getScheduledUpdateInterval()); - var scheduledMsg = new CalculatedFieldDynamicArgumentsRefreshMsg(tenantId, cfCtx.getCfId()); - - ScheduledFuture scheduledFuture = systemContext - .schedulePeriodicMsgWithDelay(ctx, scheduledMsg, refreshDynamicSourceInterval, refreshDynamicSourceInterval); - cfDynamicArgumentsRefreshTasks.put(cf.getId(), scheduledFuture); - log.debug("[{}][{}] Scheduled dynamic arguments refresh task for CF!", tenantId, cf.getId()); - } - - public void onDynamicArgumentsRefreshMsg(CalculatedFieldDynamicArgumentsRefreshMsg msg) { - log.debug("[{}] [{}] Processing CF dynamic arguments refresh task.", tenantId, msg.getCfId()); - CalculatedFieldCtx cfCtx = calculatedFields.get(msg.getCfId()); - if (cfCtx == null) { - log.debug("[{}][{}] Failed to find CF context, going to stop dynamic arguments refresh task for CF.", tenantId, msg.getCfId()); - cancelCfDynamicArgumentsRefreshTaskIfExists(msg.getCfId(), true); - return; - } - applyToTargetCfEntityActors(cfCtx, msg.getCallback(), (id, cb) -> refreshDynamicArgumentsForEntity(id, msg.getCfId(), cb)); - } - - private void refreshDynamicArgumentsForEntity(EntityId entityId, CalculatedFieldId cfId, TbCallback callback) { - log.debug("Pushing CF dynamic arguments refresh msg to specific actor [{}]", entityId); - getOrCreateActor(entityId).tell(new EntityCalculatedFieldDynamicArgumentsRefreshMsg(tenantId, cfId, callback)); - } - private void linkedTelemetryMsgForEntity(EntityId entityId, EntityCalculatedFieldLinkedTelemetryMsg msg) { log.debug("Pushing linked telemetry msg to specific actor [{}]", entityId); getOrCreateActor(entityId).tell(msg); @@ -565,7 +504,6 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware // We use copy on write lists to safely pass the reference to another actor for the iteration. // Alternative approach would be to use any list but avoid modifications to the list (change the complete map value instead) entityIdCalculatedFields.computeIfAbsent(cf.getEntityId(), id -> new CopyOnWriteArrayList<>()).add(cfCtx); - scheduleDynamicArgumentsRefreshTaskForCfIfNeeded(cfCtx); } private void initCalculatedFieldLink(CalculatedFieldLink link) { diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/EntityCalculatedFieldDynamicArgumentsRefreshMsg.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/EntityCalculatedFieldDynamicArgumentsRefreshMsg.java deleted file mode 100644 index fdf864611f..0000000000 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/EntityCalculatedFieldDynamicArgumentsRefreshMsg.java +++ /dev/null @@ -1,37 +0,0 @@ -/** - * Copyright © 2016-2025 The Thingsboard Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.thingsboard.server.actors.calculatedField; - -import lombok.Data; -import org.thingsboard.server.common.data.id.CalculatedFieldId; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.msg.MsgType; -import org.thingsboard.server.common.msg.ToCalculatedFieldSystemMsg; -import org.thingsboard.server.common.msg.queue.TbCallback; - -@Data -public class EntityCalculatedFieldDynamicArgumentsRefreshMsg implements ToCalculatedFieldSystemMsg { - - private final TenantId tenantId; - private final CalculatedFieldId cfId; - private final TbCallback callback; - - @Override - public MsgType getMsgType() { - return MsgType.CF_ENTITY_DYNAMIC_ARGUMENTS_REFRESH_MSG; - } - -} diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java index 45305ca9e3..8709e0cb68 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java @@ -45,6 +45,7 @@ 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; @@ -102,6 +103,9 @@ public abstract class AbstractCalculatedFieldProcessingService { return Futures.whenAllComplete(argFutures.values()).call(() -> { var result = createStateByType(ctx); result.updateState(ctx, resolveArgumentFutures(argFutures)); + if (ctx.hasRelationQueryDynamicArguments() && result instanceof GeofencingCalculatedFieldState geofencingCalculatedFieldState) { + geofencingCalculatedFieldState.setLastDynamicArgumentsRefreshTs(System.currentTimeMillis()); + } return result; }, MoreExecutors.directExecutor()); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java index 9a1d06cf24..6d877331bd 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java @@ -62,7 +62,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { boolean entryUpdated; if (existingEntry == null || newEntry.isForceResetPrevious()) { - validateNewEntry(newEntry); + validateNewEntry(key, newEntry); arguments.put(key, newEntry); entryUpdated = true; } else { @@ -93,7 +93,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { } } - protected void validateNewEntry(ArgumentEntry newEntry) {} + protected void validateNewEntry(String key, ArgumentEntry newEntry) {} private void updateLastUpdateTimestamp(ArgumentEntry entry) { long newTs = this.latestTimestamp; diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java index c2cc083853..13012a028a 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java @@ -43,11 +43,13 @@ import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.usagerecord.ApiLimitService; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; +import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; import java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import static org.thingsboard.common.util.ExpressionFunctionsUtil.userDefinedFunctions; @@ -78,9 +80,12 @@ public class CalculatedFieldCtx { private long maxStateSize; private long maxSingleValueArgumentSize; + private boolean relationQueryDynamicArguments; private List mainEntityGeofencingArgumentNames; private List linkedEntityGeofencingArgumentNames; + private long scheduledUpdateIntervalMillis; + public CalculatedFieldCtx(CalculatedField calculatedField, TbelInvokeService tbelInvokeService, ApiLimitService apiLimitService, RelationService relationService) { this.calculatedField = calculatedField; @@ -101,6 +106,7 @@ public class CalculatedFieldCtx { var refId = entry.getValue().getRefEntityId(); var refKey = entry.getValue().getRefEntityKey(); if (refId == null && entry.getValue().hasDynamicSource()) { + relationQueryDynamicArguments = true; continue; } if (refId == null || refId.equals(calculatedField.getEntityId())) { @@ -126,6 +132,9 @@ public class CalculatedFieldCtx { }); } } + if (calculatedField.getConfiguration() instanceof ScheduledUpdateSupportedCalculatedFieldConfiguration scheduledConfig) { + this.scheduledUpdateIntervalMillis = scheduledConfig.isScheduledUpdateEnabled() ? TimeUnit.SECONDS.toMillis(scheduledConfig.getScheduledUpdateInterval()) : -1L; + } this.tbelInvokeService = tbelInvokeService; this.relationService = relationService; @@ -329,25 +338,42 @@ public class CalculatedFieldCtx { public boolean hasOtherSignificantChanges(CalculatedFieldCtx other) { boolean expressionChanged = calculatedField.getConfiguration() instanceof ExpressionBasedCalculatedFieldConfiguration && !expression.equals(other.expression); boolean outputChanged = !output.equals(other.output); - return expressionChanged || outputChanged; + boolean scheduledUpdatesConfigChanged = scheduledUpdateIntervalMillis != other.scheduledUpdateIntervalMillis; + return expressionChanged || outputChanged || scheduledUpdatesConfigChanged; } public boolean hasStateChanges(CalculatedFieldCtx other) { boolean typeChanged = !cfType.equals(other.cfType); boolean argumentsChanged = !arguments.equals(other.arguments); - return typeChanged || argumentsChanged; + boolean geoZoneGroupsConfigChanged = hasGeofencingZoneGroupConfigurationChanges(other); + return typeChanged || argumentsChanged || geoZoneGroupsConfigChanged; } - public boolean hasSchedulingConfigChanges(CalculatedFieldCtx other) { - if (calculatedField.getConfiguration() instanceof ScheduledUpdateSupportedCalculatedFieldConfiguration thisConfig - && other.calculatedField.getConfiguration() instanceof ScheduledUpdateSupportedCalculatedFieldConfiguration otherConfig) { - boolean refreshTriggerChanged = thisConfig.isScheduledUpdateEnabled() != otherConfig.isScheduledUpdateEnabled(); - boolean refreshIntervalChanged = thisConfig.getScheduledUpdateInterval() != otherConfig.getScheduledUpdateInterval(); - return refreshTriggerChanged || refreshIntervalChanged; + private boolean hasGeofencingZoneGroupConfigurationChanges(CalculatedFieldCtx other) { + if (calculatedField.getConfiguration() instanceof GeofencingCalculatedFieldConfiguration thisConfig + && other.calculatedField.getConfiguration() instanceof GeofencingCalculatedFieldConfiguration otherConfig) { + return !thisConfig.getZoneGroups().equals(otherConfig.getZoneGroups()); } return false; } + public boolean hasRelationQueryDynamicArguments() { + return relationQueryDynamicArguments && scheduledUpdateIntervalMillis != -1; + } + + public boolean shouldFetchDynamicArgumentsFromDb(CalculatedFieldState state) { + if (!hasRelationQueryDynamicArguments()) { + return false; + } + if (!(state instanceof GeofencingCalculatedFieldState geofencingState)) { + return false; + } + if (geofencingState.getLastDynamicArgumentsRefreshTs() == -1L) { + return true; + } + return geofencingState.getLastDynamicArgumentsRefreshTs() < System.currentTimeMillis() - scheduledUpdateIntervalMillis; + } + public String getSizeExceedsLimitMessage() { return "Failed to init CF state. State size exceeds limit of " + (maxStateSize / 1024) + "Kb!"; } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java index 5f8e7538c4..3e0964bfd2 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java @@ -50,13 +50,6 @@ public interface CalculatedFieldState { long getLatestTimestamp(); - default void setDirty(boolean dirty) { - } - - default boolean isDirty() { - return false; - } - void setRequiredArguments(List requiredArguments); boolean updateState(CalculatedFieldCtx ctx, Map argumentValues); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java index 80b650fc7c..5a437b52f3 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java @@ -48,9 +48,10 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { } @Override - protected void validateNewEntry(ArgumentEntry newEntry) { + protected void validateNewEntry(String key, ArgumentEntry newEntry) { if (newEntry instanceof TsRollingArgumentEntry) { - throw new IllegalArgumentException("Rolling argument entry is not supported for simple calculated fields."); + throw new IllegalArgumentException("Unsupported argument type detected for argument: " + key + ". " + + "Rolling argument entry is not supported for simple calculated fields."); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingArgumentEntry.java index f526cc00ab..53e5c19e72 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingArgumentEntry.java @@ -41,7 +41,7 @@ public class GeofencingArgumentEntry implements ArgumentEntry { } public GeofencingArgumentEntry(EntityId entityId, TransportProtos.AttributeValueProto entry) { - this.zoneStates = toZones(Map.of(entityId, ProtoUtils.fromProto(entry))); + this(Map.of(entityId, ProtoUtils.fromProto(entry))); } public GeofencingArgumentEntry(Map entityIdkvEntryMap) { @@ -63,6 +63,10 @@ public class GeofencingArgumentEntry implements ArgumentEntry { 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)) { @@ -97,6 +101,10 @@ public class GeofencingArgumentEntry implements ArgumentEntry { zoneStates.put(zoneId, newZoneState); return true; } + if (newZoneState.getPerimeterDefinition() == null) { + zoneStates.remove(zoneId); + return true; + } return existingZoneState.update(newZoneState); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java index 506ddcff78..398de0c20b 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java @@ -15,16 +15,19 @@ */ 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; @@ -39,7 +42,6 @@ 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.HashMap; import java.util.List; import java.util.Map; import java.util.stream.Collectors; @@ -51,18 +53,14 @@ import static org.thingsboard.server.common.data.cf.configuration.geofencing.Geo @Data @Slf4j +@NoArgsConstructor @EqualsAndHashCode(callSuper = true) public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { - private boolean dirty; + private long lastDynamicArgumentsRefreshTs = -1; - public GeofencingCalculatedFieldState() { - super(new ArrayList<>(), new HashMap<>(), false, -1); - this.dirty = false; - } - - public GeofencingCalculatedFieldState(List argNames) { - super(argNames); + public GeofencingCalculatedFieldState(List requiredArguments) { + super(requiredArguments); } @Override @@ -71,49 +69,21 @@ public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { } @Override - public boolean updateState(CalculatedFieldCtx ctx, Map argumentValues) { - if (arguments == null) { - arguments = new HashMap<>(); - } - - boolean stateUpdated = false; - - for (var entry : argumentValues.entrySet()) { - String key = entry.getKey(); - ArgumentEntry newEntry = entry.getValue(); - - checkArgumentSize(key, newEntry, ctx); - - ArgumentEntry existingEntry = arguments.get(key); - boolean entryUpdated; - - if (existingEntry == null || newEntry.isForceResetPrevious()) { - entryUpdated = switch (key) { - case ENTITY_ID_LATITUDE_ARGUMENT_KEY, ENTITY_ID_LONGITUDE_ARGUMENT_KEY -> { - if (!(newEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry)) { - throw new IllegalArgumentException("Unsupported argument entry type for " + key + " argument: " + newEntry.getType() + ". " + - "Only SINGLE_VALUE type is allowed."); - } - arguments.put(key, singleValueArgumentEntry); - yield true; - } - default -> { - if (!(newEntry instanceof GeofencingArgumentEntry geofencingArgumentEntry)) { - throw new IllegalArgumentException("Unsupported argument entry type for " + key + " argument: " + newEntry.getType() + ". " + - "Only GEOFENCING type is allowed."); - } - arguments.put(key, geofencingArgumentEntry); - yield true; - } - }; - } else { - entryUpdated = existingEntry.updateEntry(newEntry); + 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."); + } } - if (entryUpdated) { - stateUpdated = true; + default -> { + if (!(newEntry instanceof GeofencingArgumentEntry)) { + throw new IllegalArgumentException("Unsupported argument entry type for " + key + " argument: " + newEntry.getType() + ". " + + "Only GEOFENCING type is allowed."); + } } } - return stateUpdated; } @Override @@ -125,7 +95,7 @@ public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { var geofencingCfg = (GeofencingCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); Map zoneGroups = geofencingCfg.getZoneGroups(); - ObjectNode resultNode = JacksonUtil.newObjectNode(); + ObjectNode valuesNode = JacksonUtil.newObjectNode(); List> relationFutures = new ArrayList<>(); getGeofencingArguments().forEach((argumentKey, argumentEntry) -> { @@ -154,10 +124,11 @@ public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { relationFutures.add(f); } }); - updateResultNode(argumentKey, zoneResults, zoneGroupCfg.getReportStrategy(), resultNode); + updateValuesNode(argumentKey, zoneResults, zoneGroupCfg.getReportStrategy(), valuesNode); }); - var result = new CalculatedFieldResult(ctx.getOutput().getType(), ctx.getOutput().getScope(), resultNode); + 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); } @@ -171,7 +142,7 @@ public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { .collect(Collectors.toMap(Map.Entry::getKey, entry -> (GeofencingArgumentEntry) entry.getValue())); } - private void updateResultNode(String argumentKey, List zoneResults, GeofencingReportStrategy geofencingReportStrategy, ObjectNode resultNode) { + private void updateValuesNode(String argumentKey, List zoneResults, GeofencingReportStrategy geofencingReportStrategy, ObjectNode resultNode) { GeofencingEvalResult aggregationResult = aggregateZoneGroup(zoneResults); final String eventKey = argumentKey + "Event"; final String statusKey = argumentKey + "Status"; @@ -185,6 +156,16 @@ public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState { } } + 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 zoneResults) { boolean nowInside = zoneResults.stream().anyMatch(r -> INSIDE.equals(r.status())); boolean prevInside = zoneResults.stream() diff --git a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java index b6a1faa1ed..ffd876575c 100644 --- a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java @@ -831,10 +831,11 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes TenantProfile foundTenantProfile = doGet("/api/tenantProfile/" + tenantProfileEntityInfo.getId().getId().toString(), TenantProfile.class); assertThat(foundTenantProfile).isNotNull(); assertThat(foundTenantProfile.getDefaultProfileConfiguration()).isNotNull(); - foundTenantProfile.getDefaultProfileConfiguration().setMinAllowedScheduledUpdateIntervalInSecForCF(TIMEOUT / 10); + int minAllowedScheduledUpdateIntervalInSecForCF = TIMEOUT / 10; + foundTenantProfile.getDefaultProfileConfiguration().setMinAllowedScheduledUpdateIntervalInSecForCF(minAllowedScheduledUpdateIntervalInSecForCF); TenantProfile savedTenantProfile = doPost("/api/tenantProfile", foundTenantProfile, TenantProfile.class); assertThat(savedTenantProfile).isNotNull(); - assertThat(savedTenantProfile.getDefaultProfileConfiguration().getMinAllowedScheduledUpdateIntervalInSecForCF()).isEqualTo(TIMEOUT / 10); + assertThat(savedTenantProfile.getDefaultProfileConfiguration().getMinAllowedScheduledUpdateIntervalInSecForCF()).isEqualTo(minAllowedScheduledUpdateIntervalInSecForCF); loginTenantAdmin(); // --- Arrange entities --- @@ -884,7 +885,8 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes cfg.setOutput(out); // Enable scheduled refresh with a 6-second interval - cfg.setScheduledUpdateInterval(6); + cfg.setScheduledUpdateInterval(minAllowedScheduledUpdateIntervalInSecForCF); + cfg.setScheduledUpdateEnabled(true); cf.setConfiguration(cfg); CalculatedField savedCalculatedField = doPost("/api/calculatedField", cf, CalculatedField.class); @@ -935,7 +937,7 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes relAllowedB.setType("AllowedZone"); doPost("/api/relation", relAllowedB).andExpect(status().isOk()); - awaitForCalculatedFieldEntityMessageProcessorToRegisterCfStateAsDirty(device.getId(), savedCalculatedField.getId()); + awaitForCalculatedFieldEntityMessageProcessorToRegisterCfStateAsReadyToRefreshDynamicArguments(device.getId(), savedCalculatedField.getId(), minAllowedScheduledUpdateIntervalInSecForCF); // --- Same coordinates as before, but now we expect ENTERED since a new zone is registered --- doPost("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/timeseries/unusedScope", diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java index fd01581e36..d7c9ad4590 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractWebTest.java @@ -155,6 +155,7 @@ import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.queue.memory.InMemoryStorage; import org.thingsboard.server.service.cf.CfRocksDb; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; +import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; import org.thingsboard.server.service.entitiy.tenant.profile.TbTenantProfileService; import org.thingsboard.server.service.security.auth.jwt.RefreshTokenRequest; import org.thingsboard.server.service.security.auth.rest.LoginRequest; @@ -1104,14 +1105,15 @@ public abstract class AbstractWebTest extends AbstractInMemoryStorageTest { }); } - protected void awaitForCalculatedFieldEntityMessageProcessorToRegisterCfStateAsDirty(EntityId entityId, CalculatedFieldId cfId) { + protected void awaitForCalculatedFieldEntityMessageProcessorToRegisterCfStateAsReadyToRefreshDynamicArguments(EntityId entityId, CalculatedFieldId cfId, int scheduledUpdateInterval) { CalculatedFieldEntityMessageProcessor processor = getCalculatedFieldEntityMessageProcessor(entityId); Map statesMap = (Map) ReflectionTestUtils.getField(processor, "states"); - Awaitility.await("CF state for entity actor marked as dirty").atMost(5, TimeUnit.SECONDS).until(() -> { + Awaitility.await("CF state for entity actor ready to refresh dynamic arguments").atMost(TIMEOUT, TimeUnit.SECONDS).until(() -> { CalculatedFieldState calculatedFieldState = statesMap.get(cfId); - boolean stateDirty = calculatedFieldState != null && calculatedFieldState.isDirty(); - log.warn("entityId {}, cfId {}, state dirty == {}", entityId, cfId, stateDirty); - return stateDirty; + boolean isReady = calculatedFieldState != null && ((GeofencingCalculatedFieldState) calculatedFieldState).getLastDynamicArgumentsRefreshTs() + < System.currentTimeMillis() - TimeUnit.SECONDS.toMillis(scheduledUpdateInterval); + log.warn("entityId {}, cfId {}, state ready to refresh == {}", entityId, cfId, isReady); + return isReady; }); } diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldStateTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldStateTest.java index 691a1f7ec4..a86af77555 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldStateTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldStateTest.java @@ -257,7 +257,7 @@ public class GeofencingCalculatedFieldStateTest { assertThat(result2).isNotNull(); assertThat(result2.getType()).isEqualTo(output.getType()); assertThat(result2.getScope()).isEqualTo(output.getScope()); - assertThat(result2.getResult()).isEqualTo( + assertThat(result2.getResult().get("values")).isEqualTo( JacksonUtil.newObjectNode() .put("allowedZonesEvent", "LEFT") .put("allowedZonesStatus", "OUTSIDE") @@ -329,7 +329,7 @@ public class GeofencingCalculatedFieldStateTest { assertThat(result2).isNotNull(); assertThat(result2.getType()).isEqualTo(output.getType()); assertThat(result2.getScope()).isEqualTo(output.getScope()); - assertThat(result2.getResult()).isEqualTo( + assertThat(result2.getResult().get("values")).isEqualTo( JacksonUtil.newObjectNode() .put("allowedZonesEvent", "LEFT") .put("restrictedZonesEvent", "ENTERED") @@ -401,7 +401,7 @@ public class GeofencingCalculatedFieldStateTest { assertThat(result2).isNotNull(); assertThat(result2.getType()).isEqualTo(output.getType()); assertThat(result2.getScope()).isEqualTo(output.getScope()); - assertThat(result2.getResult()).isEqualTo( + assertThat(result2.getResult().get("values")).isEqualTo( JacksonUtil.newObjectNode() .put("allowedZonesStatus", "OUTSIDE") .put("restrictedZonesStatus", "INSIDE") diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldStateTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldStateTest.java index 8c631ecf6f..3aef8896de 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldStateTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldStateTest.java @@ -123,7 +123,8 @@ public class SimpleCalculatedFieldStateTest { Map newArgs = Map.of("key3", new TsRollingArgumentEntry(10, 30000L)); assertThatThrownBy(() -> state.updateState(ctx, newArgs)) .isInstanceOf(IllegalArgumentException.class) - .hasMessage("Rolling argument entry is not supported for simple calculated fields."); + .hasMessage("Unsupported argument type detected for argument: key3. " + + "Rolling argument entry is not supported for simple calculated fields."); } @Test diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ScheduledUpdateSupportedCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ScheduledUpdateSupportedCalculatedFieldConfiguration.java index 7902a9cf5b..d0c5786f62 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ScheduledUpdateSupportedCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/ScheduledUpdateSupportedCalculatedFieldConfiguration.java @@ -15,11 +15,8 @@ */ package org.thingsboard.server.common.data.cf.configuration; -import com.fasterxml.jackson.annotation.JsonIgnore; - public interface ScheduledUpdateSupportedCalculatedFieldConfiguration extends CalculatedFieldConfiguration { - @JsonIgnore boolean isScheduledUpdateEnabled(); int getScheduledUpdateInterval(); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfiguration.java index dc331f5876..b331abc50b 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfiguration.java @@ -34,6 +34,8 @@ public class GeofencingCalculatedFieldConfiguration implements ArgumentsBasedCal private EntityCoordinates entityCoordinates; private Map zoneGroups; + + private boolean scheduledUpdateEnabled; private int scheduledUpdateInterval; private Output output; @@ -61,11 +63,6 @@ public class GeofencingCalculatedFieldConfiguration implements ArgumentsBasedCal return output; } - @Override - public boolean isScheduledUpdateEnabled() { - return scheduledUpdateInterval > 0 && zoneGroups.values().stream().anyMatch(ZoneGroupConfiguration::hasDynamicSource); - } - @Override public void validate() { if (entityCoordinates == null) { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java index 4c8b9e06bd..c6bd9a7f38 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java @@ -173,7 +173,7 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura @Schema(example = "10") private long maxArgumentsPerCF = 10; @Schema(example = "3600") - private int minAllowedScheduledUpdateIntervalInSecForCF = 3600; + private int minAllowedScheduledUpdateIntervalInSecForCF = 60; @Schema(example = "10") private int maxRelationLevelPerCfArgument = 10; @Builder.Default diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfigurationTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfigurationTest.java index 91a47aac57..c5fb1f8953 100644 --- a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfigurationTest.java +++ b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/geofencing/GeofencingCalculatedFieldConfigurationTest.java @@ -101,33 +101,6 @@ public class GeofencingCalculatedFieldConfigurationTest { verify(zoneGroupConfigurationB).validate(zoneGroupBName); } - @Test - void scheduledUpdateDisabledWhenIntervalIsZero() { - var cfg = new GeofencingCalculatedFieldConfiguration(); - cfg.setScheduledUpdateInterval(0); - assertThat(cfg.isScheduledUpdateEnabled()).isFalse(); - } - - @Test - void scheduledUpdateDisabledWhenIntervalIsGreaterThanZeroButNoZonesWithDynamicArguments() { - var cfg = new GeofencingCalculatedFieldConfiguration(); - var zoneGroupConfigurationMock = mock(ZoneGroupConfiguration.class); - when(zoneGroupConfigurationMock.hasDynamicSource()).thenReturn(false); - cfg.setZoneGroups(Map.of("someGroupName", zoneGroupConfigurationMock)); - cfg.setScheduledUpdateInterval(60); - assertThat(cfg.isScheduledUpdateEnabled()).isFalse(); - } - - @Test - void scheduledUpdateEnabledWhenIntervalIsGreaterThanZeroAndDynamicArgumentsPresent() { - var cfg = new GeofencingCalculatedFieldConfiguration(); - var zoneGroupConfigurationMock = mock(ZoneGroupConfiguration.class); - when(zoneGroupConfigurationMock.hasDynamicSource()).thenReturn(true); - cfg.setZoneGroups(Map.of("someGroupName", zoneGroupConfigurationMock)); - cfg.setScheduledUpdateInterval(60); - assertThat(cfg.isScheduledUpdateEnabled()).isTrue(); - } - @Test void testGetArgumentsOverride() { var cfg = new GeofencingCalculatedFieldConfiguration(); diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java b/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java index 48b07af29b..20043582d7 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/MsgType.java @@ -147,10 +147,7 @@ public enum MsgType { /* CF Manager Actor -> CF Entity actor */ CF_ENTITY_TELEMETRY_MSG, CF_ENTITY_INIT_CF_MSG, - CF_ENTITY_DELETE_MSG, - - CF_DYNAMIC_ARGUMENTS_REFRESH_MSG, - CF_ENTITY_DYNAMIC_ARGUMENTS_REFRESH_MSG; + CF_ENTITY_DELETE_MSG; @Getter private final boolean ignoreOnStart; diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java index 49bdcca0fb..7f563ea436 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java @@ -99,53 +99,7 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { } @Test - public void testSaveGeofencingCalculatedField_shouldNotChangeScheduledInterval() { - // Arrange a device - Device device = createTestDevice(); - - // Build a valid Geofencing configuration - GeofencingCalculatedFieldConfiguration cfg = new GeofencingCalculatedFieldConfiguration(); - - // Coordinates: TS_LATEST, no dynamic source - EntityCoordinates entityCoordinates = new EntityCoordinates("latitude", "longitude"); - cfg.setEntityCoordinates(entityCoordinates); - - // Zone-group argument (ATTRIBUTE) — no DYNAMIC configuration, so no scheduling even if the scheduled interval is set - ZoneGroupConfiguration zoneGroupConfiguration = new ZoneGroupConfiguration("allowed", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false); - zoneGroupConfiguration.setRefEntityId(device.getId()); - cfg.setZoneGroups(Map.of("allowed", zoneGroupConfiguration)); - - // Set a scheduled interval to some value - cfg.setScheduledUpdateInterval(600); - - // Create & save Calculated Field - CalculatedField cf = new CalculatedField(); - cf.setTenantId(tenantId); - cf.setEntityId(device.getId()); - cf.setType(CalculatedFieldType.GEOFENCING); - cf.setName("GF clamp test"); - cf.setConfigurationVersion(0); - cf.setConfiguration(cfg); - - CalculatedField saved = calculatedFieldService.save(cf); - - assertThat(saved).isNotNull(); - assertThat(saved.getConfiguration()).isInstanceOf(GeofencingCalculatedFieldConfiguration.class); - - var geofencingCalculatedFieldConfiguration = (GeofencingCalculatedFieldConfiguration) saved.getConfiguration(); - - // Assert: the interval is saved, but scheduling is not enabled - int savedInterval = geofencingCalculatedFieldConfiguration.getScheduledUpdateInterval(); - boolean scheduledUpdateEnabled = geofencingCalculatedFieldConfiguration.isScheduledUpdateEnabled(); - - assertThat(savedInterval).isEqualTo(600); - assertThat(scheduledUpdateEnabled).isFalse(); - - calculatedFieldService.deleteCalculatedField(tenantId, saved.getId()); - } - - @Test - public void testSaveGeofencingCalculatedField_shouldThrowWhenScheduledIntervalIsLessThanMinAllowedIntervalInTenantProfile() { + public void testSaveGeofencingCalculatedField_shouldThrowWhenScheduledIntervalLessThanMinAllowedIntervalInTenantProfile() { // Arrange a device Device device = createTestDevice(); @@ -165,15 +119,22 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { zoneGroupConfiguration.setRefDynamicSourceConfiguration(dynamicSourceConfiguration); cfg.setZoneGroups(Map.of("allowed", zoneGroupConfiguration)); + // Get tenant profile min. + int min = tbTenantProfileCache.get(tenantId) + .getDefaultProfileConfiguration() + .getMinAllowedScheduledUpdateIntervalInSecForCF(); + int valueFromConfig = min - 10; + // Enable scheduling with an interval below tenant min - cfg.setScheduledUpdateInterval(600); + cfg.setScheduledUpdateEnabled(true); + cfg.setScheduledUpdateInterval(valueFromConfig); // Create & save Calculated Field CalculatedField cf = new CalculatedField(); cf.setTenantId(tenantId); cf.setEntityId(device.getId()); cf.setType(CalculatedFieldType.GEOFENCING); - cf.setName("GF clamp test"); + cf.setName("GF min allowed scheduled update interval test"); cf.setConfigurationVersion(0); cf.setConfiguration(cfg); @@ -185,7 +146,7 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { } @Test - public void testSaveGeofencingCalculatedField_shouldThrowWhenRelationLevelIsGreaterThanMaxAllowedRelationLevelInTenantProfile() { + public void testSaveGeofencingCalculatedField_shouldThrowWhenRelationLevelGreaterThanMaxAllowedRelationLevelInTenantProfile() { // Arrange a device Device device = createTestDevice(); @@ -210,7 +171,7 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { cf.setTenantId(tenantId); cf.setEntityId(device.getId()); cf.setType(CalculatedFieldType.GEOFENCING); - cf.setName("GF clamp test"); + cf.setName("GF max relation level test"); cf.setConfigurationVersion(0); cf.setConfiguration(cfg); @@ -221,7 +182,7 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { } @Test - public void testSaveGeofencingCalculatedField_shouldUseScheduledIntervalFromConfig() { + public void testSaveGeofencingCalculatedField_shouldSaveWithoutDataValidationExceptionOnScheduledUpdateInterval() { // Arrange a device Device device = createTestDevice(); @@ -245,10 +206,10 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { int min = tbTenantProfileCache.get(tenantId) .getDefaultProfileConfiguration() .getMinAllowedScheduledUpdateIntervalInSecForCF(); - + int valueFromConfig = min + 100; // Enable scheduling with an interval greater than tenant min - int valueFromConfig = min + 100; + cfg.setScheduledUpdateEnabled(true); cfg.setScheduledUpdateInterval(valueFromConfig); // Create & save Calculated Field @@ -256,7 +217,7 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { cf.setTenantId(tenantId); cf.setEntityId(device.getId()); cf.setType(CalculatedFieldType.GEOFENCING); - cf.setName("GF no clamp test"); + cf.setName("GF no validation error test"); cf.setConfigurationVersion(0); cf.setConfiguration(cfg); @@ -267,7 +228,6 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { var geofencingCalculatedFieldConfiguration = (GeofencingCalculatedFieldConfiguration) saved.getConfiguration(); - // Assert: the interval is clamped up to tenant profile min (or stays >= original if already >= min) int savedInterval = geofencingCalculatedFieldConfiguration.getScheduledUpdateInterval(); assertThat(savedInterval).isEqualTo(valueFromConfig); From 909497703a71449d5d7a388d96da45bc745d13bb Mon Sep 17 00:00:00 2001 From: dshvaika Date: Wed, 1 Oct 2025 15:35:55 +0300 Subject: [PATCH 2/6] new relation path query --- ...tractCalculatedFieldProcessingService.java | 6 ++ .../server/dao/relation/RelationService.java | 3 + .../CFArgumentDynamicSourceType.java | 2 +- .../CfArgumentDynamicSourceConfiguration.java | 3 +- ...onPathQueryDynamicSourceConfiguration.java | 73 ++++++++++++++ .../cf/configuration/RelationQueryBased.java | 29 ++++++ ...lationQueryDynamicSourceConfiguration.java | 9 +- .../relation/EntityRelationPathQuery.java | 24 +++++ .../data/relation/RelationPathLevel.java | 30 ++++++ .../dao/relation/BaseRelationService.java | 21 ++++ .../server/dao/relation/RelationDao.java | 3 + .../CalculatedFieldDataValidator.java | 6 +- .../dao/sql/relation/JpaRelationDao.java | 99 +++++++++++++++++++ 13 files changed, 295 insertions(+), 13 deletions(-) create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryBased.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelationPathQuery.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/relation/RelationPathLevel.java diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java index 8709e0cb68..bd3546082a 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java @@ -26,6 +26,7 @@ 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.cf.configuration.RelationQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -178,6 +179,11 @@ public abstract class AbstractCalculatedFieldProcessingService { yield Futures.transform(relationService.findByQuery(tenantId, configuration.toEntityRelationsQuery(entityId)), configuration::resolveEntityIds, calculatedFieldCallbackExecutor); } + case RELATION_PATH_QUERY -> { + var configuration = (RelationPathQueryDynamicSourceConfiguration) refDynamicSourceConfiguration; + yield Futures.transform(relationService.findByRelationPathQueryAsync(tenantId, configuration.toRelationPathQuery(entityId)), + configuration::resolveEntityIds, calculatedFieldCallbackExecutor); + } }; } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/relation/RelationService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/relation/RelationService.java index 1db5739e94..a0bc9a72e6 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/relation/RelationService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/relation/RelationService.java @@ -20,6 +20,7 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelationInfo; +import org.thingsboard.server.common.data.relation.EntityRelationPathQuery; import org.thingsboard.server.common.data.relation.EntityRelationsQuery; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleChainType; @@ -83,6 +84,8 @@ public interface RelationService { List findRuleNodeToRuleChainRelations(TenantId tenantId, RuleChainType ruleChainType, int limit); + ListenableFuture> findByRelationPathQueryAsync(TenantId tenantId, EntityRelationPathQuery relationPathQuery); + // TODO: This method may be useful for some validations in the future // ListenableFuture checkRecursiveRelation(EntityId from, EntityId to); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CFArgumentDynamicSourceType.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CFArgumentDynamicSourceType.java index bd2e9b0c00..35c6cdf562 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CFArgumentDynamicSourceType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CFArgumentDynamicSourceType.java @@ -17,6 +17,6 @@ package org.thingsboard.server.common.data.cf.configuration; public enum CFArgumentDynamicSourceType { - RELATION_QUERY + RELATION_QUERY, RELATION_PATH_QUERY } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CfArgumentDynamicSourceConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CfArgumentDynamicSourceConfiguration.java index f36071615e..397b1d016e 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CfArgumentDynamicSourceConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CfArgumentDynamicSourceConfiguration.java @@ -26,7 +26,8 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo; property = "type" ) @JsonSubTypes({ - @JsonSubTypes.Type(value = RelationQueryDynamicSourceConfiguration.class, name = "RELATION_QUERY") + @JsonSubTypes.Type(value = RelationQueryDynamicSourceConfiguration.class, name = "RELATION_QUERY"), + @JsonSubTypes.Type(value = RelationPathQueryDynamicSourceConfiguration.class, name = "RELATION_PATH_QUERY") }) @JsonIgnoreProperties(ignoreUnknown = true) public interface CfArgumentDynamicSourceConfiguration { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java new file mode 100644 index 0000000000..c7889be19d --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java @@ -0,0 +1,73 @@ +/** + * Copyright © 2016-2025 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.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, RelationQueryBased { + + private List 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 resolveEntityIds(List relations) { + EntitySearchDirection lastLevelDirection = getLastLevel().direction(); + return switch (lastLevelDirection) { + case FROM -> relations.stream().map(EntityRelation::getTo).toList(); + case TO -> relations.stream().map(EntityRelation::getFrom).toList(); + }; + } + + @Override + @JsonIgnore + public int getMaxLevel() { + return levels != null ? levels.size() : 0; + } + + public EntityRelationPathQuery toRelationPathQuery(EntityId entityId) { + return new EntityRelationPathQuery(entityId, levels); + } + + private RelationPathLevel getLastLevel() { + if (CollectionsUtil.isEmpty(levels)) { + throw new NoSuchElementException(); + } + return levels.get(levels.size() - 1); + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryBased.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryBased.java new file mode 100644 index 0000000000..e3248b264f --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryBased.java @@ -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.cf.configuration; + +public interface RelationQueryBased { + + int getMaxLevel(); + + default void validateMaxRelationLevel(String argumentName, int maxAllowedRelationLevel) { + if (getMaxLevel() > maxAllowedRelationLevel) { + throw new IllegalArgumentException("Max relation level is greater than configured " + + "maximum allowed relation level in tenant profile: " + maxAllowedRelationLevel + " for argument: " + argumentName); + } + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfiguration.java index 4e9b4252c9..120d04d40d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfiguration.java @@ -29,7 +29,7 @@ import java.util.Collections; import java.util.List; @Data -public class RelationQueryDynamicSourceConfiguration implements CfArgumentDynamicSourceConfiguration { +public class RelationQueryDynamicSourceConfiguration implements CfArgumentDynamicSourceConfiguration, RelationQueryBased { private int maxLevel; private boolean fetchLastLevelOnly; @@ -59,13 +59,6 @@ public class RelationQueryDynamicSourceConfiguration implements CfArgumentDynami return maxLevel == 1; } - public void validateMaxRelationLevel(String argumentName, int maxAllowedRelationLevel) { - if (maxLevel > maxAllowedRelationLevel) { - throw new IllegalArgumentException("Max relation level is greater than configured " + - "maximum allowed relation level in tenant profile: " + maxAllowedRelationLevel + " for argument: " + argumentName); - } - } - public EntityRelationsQuery toEntityRelationsQuery(EntityId rootEntityId) { if (isSimpleRelation()) { throw new IllegalArgumentException("Entity relations query can't be created for a simple relation!"); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelationPathQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelationPathQuery.java new file mode 100644 index 0000000000..a81e3e0c86 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/relation/EntityRelationPathQuery.java @@ -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.relation; + +import org.thingsboard.server.common.data.id.EntityId; + +import java.util.List; + +public record EntityRelationPathQuery(EntityId rootEntityId, List levels) { + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/relation/RelationPathLevel.java b/common/data/src/main/java/org/thingsboard/server/common/data/relation/RelationPathLevel.java new file mode 100644 index 0000000000..c28135204f --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/relation/RelationPathLevel.java @@ -0,0 +1,30 @@ +/** + * 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.relation; + +import org.thingsboard.server.common.data.StringUtils; + +public record RelationPathLevel(EntitySearchDirection direction, String relationType) { + + public void validate() { + if (direction == null) { + throw new IllegalArgumentException("Direction must be specified!"); + } + if (StringUtils.isBlank(relationType)) { + throw new IllegalArgumentException("Relation type must be specified!"); + } + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java b/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java index 87f7e44a1f..0ec1257bcc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java @@ -32,17 +32,21 @@ import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.event.TransactionalEventListener; import org.springframework.transaction.support.TransactionSynchronizationManager; +import org.springframework.util.CollectionUtils; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.server.cache.TbTransactionalCache; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.id.UUIDBased; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelationInfo; +import org.thingsboard.server.common.data.relation.EntityRelationPathQuery; import org.thingsboard.server.common.data.relation.EntityRelationsQuery; import org.thingsboard.server.common.data.relation.EntitySearchDirection; import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter; +import org.thingsboard.server.common.data.relation.RelationPathLevel; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.relation.RelationsSearchParameters; import org.thingsboard.server.common.data.rule.RuleChainType; @@ -495,6 +499,23 @@ public class BaseRelationService implements RelationService { return relationDao.findRuleNodeToRuleChainRelations(ruleChainType, limit); } + @Override + public ListenableFuture> findByRelationPathQueryAsync(TenantId tenantId, EntityRelationPathQuery relationPathQuery) { + log.trace("Executing findByRelationPathQuery, tenantId [{}], relationPathQuery {}", tenantId, relationPathQuery); + validateId(tenantId, id -> "Invalid tenant id: " + id); + validate(relationPathQuery); + return executor.submit(() -> relationDao.findByRelationPathQuery(tenantId, relationPathQuery)); + } + + private void validate(EntityRelationPathQuery relationPathQuery) { + validateId((UUIDBased) relationPathQuery.rootEntityId(), id -> "Invalid root entity id: " + id); + List levels = relationPathQuery.levels(); + if (CollectionUtils.isEmpty(levels)) { + throw new DataValidationException("Relation path levels should be specified!"); + } + levels.forEach(RelationPathLevel::validate); + } + protected void validate(EntityRelation relation) { if (relation == null) { throw new DataValidationException("Relation type should be specified!"); diff --git a/dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java b/dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java index c061faed5e..ad53164ad7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/relation/RelationDao.java @@ -19,6 +19,7 @@ import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.common.data.relation.EntityRelationPathQuery; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleChainType; @@ -71,4 +72,6 @@ public interface RelationDao { List findRuleNodeToRuleChainRelations(RuleChainType ruleChainType, int limit); + List findByRelationPathQuery(TenantId tenantId, EntityRelationPathQuery relationPathQuery); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java index 05b782c26c..85c30ca74e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java @@ -19,7 +19,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.configuration.ArgumentsBasedCalculatedFieldConfiguration; -import org.thingsboard.server.common.data.cf.configuration.RelationQueryDynamicSourceConfiguration; +import org.thingsboard.server.common.data.cf.configuration.RelationQueryBased; import org.thingsboard.server.common.data.cf.configuration.ScheduledUpdateSupportedCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; @@ -98,10 +98,10 @@ public class CalculatedFieldDataValidator extends DataValidator if (!(calculatedField.getConfiguration() instanceof ArgumentsBasedCalculatedFieldConfiguration argumentsBasedCfg)) { return; } - Map relationQueryBasedArguments = argumentsBasedCfg.getArguments().entrySet() + Map relationQueryBasedArguments = argumentsBasedCfg.getArguments().entrySet() .stream() .filter(entry -> entry.getValue().hasDynamicSource()) - .collect(Collectors.toMap(Map.Entry::getKey, entry -> (RelationQueryDynamicSourceConfiguration) entry.getValue().getRefDynamicSourceConfiguration())); + .collect(Collectors.toMap(Map.Entry::getKey, entry -> (RelationQueryBased) entry.getValue().getRefDynamicSourceConfiguration())); if (relationQueryBasedArguments.isEmpty()) { return; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java index 7417418f54..782c9b9d53 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java @@ -25,6 +25,9 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; 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.relation.RelationTypeGroup; import org.thingsboard.server.common.data.rule.RuleChainType; import org.thingsboard.server.dao.DaoUtil; @@ -43,6 +46,7 @@ import java.util.stream.Collectors; import static org.thingsboard.server.dao.model.ModelConstants.RELATION_FROM_ID_PROPERTY; import static org.thingsboard.server.dao.model.ModelConstants.RELATION_FROM_TYPE_PROPERTY; +import static org.thingsboard.server.dao.model.ModelConstants.RELATION_TABLE_NAME; import static org.thingsboard.server.dao.model.ModelConstants.RELATION_TO_ID_PROPERTY; import static org.thingsboard.server.dao.model.ModelConstants.RELATION_TO_TYPE_PROPERTY; import static org.thingsboard.server.dao.model.ModelConstants.RELATION_TYPE_GROUP_PROPERTY; @@ -293,4 +297,99 @@ public class JpaRelationDao extends JpaAbstractDaoListeningExecutorService imple public List findRuleNodeToRuleChainRelations(RuleChainType ruleChainType, int limit) { return DaoUtil.convertDataList(relationRepository.findRuleNodeToRuleChainRelations(ruleChainType, PageRequest.of(0, limit))); } + + @Override + public List findByRelationPathQuery(TenantId tenantId, EntityRelationPathQuery query) { + List levels = query.levels(); + if (levels == null || levels.isEmpty()) { + return Collections.emptyList(); + } + String sql = buildRelationPathSql(query); + Object[] params = buildRelationPathParams(query); + + log.info("[{}] relation path query: {}", tenantId, sql); + + return jdbcTemplate.queryForList(sql, params).stream() + .map(row -> { + var entityRelation = new EntityRelation(); + var fromId = (UUID) row.get(RELATION_FROM_ID_PROPERTY); + var fromType = (String) row.get(RELATION_FROM_TYPE_PROPERTY); + var toId = (UUID) row.get(RELATION_TO_ID_PROPERTY); + var toType = (String) row.get(RELATION_TO_TYPE_PROPERTY); + var grp = (String) row.get(RELATION_TYPE_GROUP_PROPERTY); + var type = (String) row.get(RELATION_TYPE_PROPERTY); + var version = (Long) row.get(VERSION_COLUMN); + + entityRelation.setFrom(EntityIdFactory.getByTypeAndUuid(fromType, fromId)); + entityRelation.setTo(EntityIdFactory.getByTypeAndUuid(toType, toId)); + entityRelation.setType(type); + entityRelation.setTypeGroup(RelationTypeGroup.valueOf(grp)); + entityRelation.setVersion(version); + return entityRelation; + }) + .collect(Collectors.toList()); + } + + private Object[] buildRelationPathParams(EntityRelationPathQuery query) { + final List params = new ArrayList<>(); + // seed + params.add(query.rootEntityId().getId()); + params.add(query.rootEntityId().getEntityType().name()); + + // levels + for (var lvl : query.levels()) { + params.add(lvl.relationType()); + } + return params.toArray(); + } + + private static String buildRelationPathSql(EntityRelationPathQuery query) { + List levels = query.levels(); + StringBuilder sb = new StringBuilder(); + + sb.append("WITH seed AS (\n") + .append(" SELECT ?::uuid AS id, ?::varchar AS type\n") + .append(")"); + + String prev = "seed"; + for (int i = 0; i < levels.size() - 1; i++) { + RelationPathLevel lvl = levels.get(i); + boolean down = lvl.direction() == EntitySearchDirection.FROM; + + String cur = "lvl" + (i + 1); + String joinCond = down + ? "r.from_id = p.id AND r.from_type = p.type" + : "r.to_id = p.id AND r.to_type = p.type"; + String selectNext = down + ? "r.to_id AS id, r.to_type AS type" + : "r.from_id AS id, r.from_type AS type"; + + sb.append(",\n").append(cur).append(" AS (\n") + .append(" SELECT ").append(selectNext).append("\n") + .append(" FROM ").append(RELATION_TABLE_NAME).append(" r\n") + .append(" JOIN ").append(prev).append(" p ON ").append(joinCond).append("\n") + .append(" WHERE r.relation_type_group = '").append(RelationTypeGroup.COMMON).append("'\n") + .append(" AND r.relation_type = ?\n") + .append(")"); + prev = cur; + } + + RelationPathLevel last = levels.get(levels.size() - 1); + boolean lastDown = last.direction() == EntitySearchDirection.FROM; + String prevForLast = (levels.size() == 1) ? "seed" : prev; + String lastJoin = lastDown + ? "r.from_id = p.id AND r.from_type = p.type" + : "r.to_id = p.id AND r.to_type = p.type"; + + sb.append("\n") + .append("SELECT r.from_id, r.from_type, r.to_id, r.to_type,\n") + .append(" r.relation_type_group, r.relation_type, r.version\n") + .append("FROM ").append(RELATION_TABLE_NAME).append(" r\n") + .append("JOIN ").append(prevForLast).append(" p ON ").append(lastJoin).append("\n") + .append("WHERE r.relation_type_group = '").append(RelationTypeGroup.COMMON).append("'\n") + .append(" AND r.relation_type = ?"); + + return sb.toString(); + } + } From 2cb05c9d2b975930db26371f81dcefeb36824c40 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Thu, 2 Oct 2025 13:35:51 +0300 Subject: [PATCH 3/6] Added tests & temporary removed RelationQueryDynamicSourceConfiguration from interface --- .../CfArgumentDynamicSourceConfiguration.java | 1 - ...thQueryDynamicSourceConfigurationTest.java | 126 ++++++++++++++++++ ...onQueryDynamicSourceConfigurationTest.java | 14 +- .../dao/service/RelationServiceTest.java | 53 +++++++- 4 files changed, 183 insertions(+), 11 deletions(-) create mode 100644 common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfigurationTest.java diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CfArgumentDynamicSourceConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CfArgumentDynamicSourceConfiguration.java index 397b1d016e..ff746161e9 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CfArgumentDynamicSourceConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CfArgumentDynamicSourceConfiguration.java @@ -26,7 +26,6 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo; property = "type" ) @JsonSubTypes({ - @JsonSubTypes.Type(value = RelationQueryDynamicSourceConfiguration.class, name = "RELATION_QUERY"), @JsonSubTypes.Type(value = RelationPathQueryDynamicSourceConfiguration.class, name = "RELATION_PATH_QUERY") }) @JsonIgnoreProperties(ignoreUnknown = true) diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfigurationTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfigurationTest.java new file mode 100644 index 0000000000..1a6ea5de3b --- /dev/null +++ b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfigurationTest.java @@ -0,0 +1,126 @@ +/** + * 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 org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.NullAndEmptySource; +import org.mockito.junit.jupiter.MockitoExtension; +import org.thingsboard.server.common.data.id.EntityId; +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 java.util.ArrayList; +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@ExtendWith(MockitoExtension.class) +public class RelationPathQueryDynamicSourceConfigurationTest { + + @Test + void typeShouldBeRelationQuery() { + var cfg = new RelationPathQueryDynamicSourceConfiguration(); + assertThat(cfg.getType()).isEqualTo(CFArgumentDynamicSourceType.RELATION_PATH_QUERY); + } + + @ParameterizedTest + @NullAndEmptySource + void validateShouldThrowWhenLevelsIsNull(List levels) { + var cfg = new RelationPathQueryDynamicSourceConfiguration(); + cfg.setLevels(levels); + + assertThatThrownBy(cfg::validate) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("At least one relation level must be specified!"); + } + + @Test + void validateShouldCallValidateForPathLevels() { + List levels = new ArrayList<>(); + + RelationPathLevel lvl1 = mock(RelationPathLevel.class); + RelationPathLevel lvl2 = mock(RelationPathLevel.class); + levels.add(lvl1); + levels.add(lvl2); + + var cfg = new RelationPathQueryDynamicSourceConfiguration(); + cfg.setLevels(levels); + + assertThatCode(cfg::validate).doesNotThrowAnyException(); + + verify(lvl1).validate(); + verify(lvl2).validate(); + } + + @Test + void resolveEntityIds_whenDirectionFROM_thenReturnsToIds() { + List levels = new ArrayList<>(); + + RelationPathLevel lvl1 = mock(RelationPathLevel.class); + RelationPathLevel lvl2 = mock(RelationPathLevel.class); + levels.add(lvl1); + levels.add(lvl2); + + when(lvl2.direction()).thenReturn(EntitySearchDirection.FROM); + + EntityRelation rel1 = mock(EntityRelation.class); + EntityRelation rel2 = mock(EntityRelation.class); + + when(rel1.getTo()).thenReturn(mock(EntityId.class)); + when(rel2.getTo()).thenReturn(mock(EntityId.class)); + + var cfg = new RelationPathQueryDynamicSourceConfiguration(); + cfg.setLevels(levels); + + var out = cfg.resolveEntityIds(List.of(rel1, rel2)); + + assertThat(out).containsExactly(rel1.getTo(), rel2.getTo()); + } + + @Test + void resolveEntityIds_whenDirectionTO_thenReturnsFromIds() { + List levels = new ArrayList<>(); + + RelationPathLevel lvl1 = mock(RelationPathLevel.class); + RelationPathLevel lvl2 = mock(RelationPathLevel.class); + levels.add(lvl1); + levels.add(lvl2); + + when(lvl2.direction()).thenReturn(EntitySearchDirection.TO); + + EntityRelation rel1 = mock(EntityRelation.class); + EntityRelation rel2 = mock(EntityRelation.class); + + when(rel1.getFrom()).thenReturn(mock(EntityId.class)); + when(rel2.getFrom()).thenReturn(mock(EntityId.class)); + + var cfg = new RelationPathQueryDynamicSourceConfiguration(); + cfg.setLevels(levels); + + var out = cfg.resolveEntityIds(List.of(rel1, rel2)); + + assertThat(out).containsExactly(rel1.getFrom(), rel2.getFrom()); + } + +} diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfigurationTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfigurationTest.java index 86fa52ba66..afd78e47f9 100644 --- a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfigurationTest.java +++ b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfigurationTest.java @@ -123,9 +123,8 @@ public class RelationQueryDynamicSourceConfigurationTest { .hasMessage("Relation query dynamic source configuration relation type must be specified!"); } - @ParameterizedTest - @NullAndEmptySource - void isSimpleRelationTrueWhenLevelIsOneAndEntityTypesEmptyOrNull(List entityTypes) { + @Test + void isSimpleRelationTrueWhenLevelIsOneAndEntityTypesEmptyOrNull() { var cfg = new RelationQueryDynamicSourceConfiguration(); cfg.setMaxLevel(1); assertThat(cfg.isSimpleRelation()).isTrue(); @@ -138,9 +137,8 @@ public class RelationQueryDynamicSourceConfigurationTest { assertThat(cfg.isSimpleRelation()).isFalse(); } - @ParameterizedTest - @NullAndEmptySource - void toEntityRelationsQueryShouldThrowForSimpleRelation(List entityTypes) { + @Test + void toEntityRelationsQueryShouldThrowForSimpleRelation() { var cfg = new RelationQueryDynamicSourceConfiguration(); cfg.setMaxLevel(1); cfg.setFetchLastLevelOnly(false); @@ -177,7 +175,7 @@ public class RelationQueryDynamicSourceConfigurationTest { } @Test - void resolveEntityIdsFromDirectionFROMReturnsToIds() { + void resolveEntityIds_whenDirectionFROM_thenReturnsToIds() { when(rel1.getTo()).thenReturn(mock(EntityId.class)); when(rel2.getTo()).thenReturn(mock(EntityId.class)); @@ -190,7 +188,7 @@ public class RelationQueryDynamicSourceConfigurationTest { } @Test - void resolveEntityIdsFromDirectionTOReturnsFromIds() { + void resolveEntityIds_whenDirectionTO_thenReturnsFromIds() { when(rel1.getFrom()).thenReturn(mock(EntityId.class)); when(rel2.getFrom()).thenReturn(mock(EntityId.class)); diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/RelationServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/RelationServiceTest.java index 063e35ae80..bb7bfad677 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/RelationServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/RelationServiceTest.java @@ -28,9 +28,11 @@ import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.relation.EntityRelation; +import org.thingsboard.server.common.data.relation.EntityRelationPathQuery; import org.thingsboard.server.common.data.relation.EntityRelationsQuery; import org.thingsboard.server.common.data.relation.EntitySearchDirection; import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter; +import org.thingsboard.server.common.data.relation.RelationPathLevel; import org.thingsboard.server.common.data.relation.RelationTypeGroup; import org.thingsboard.server.common.data.relation.RelationsSearchParameters; import org.thingsboard.server.dao.exception.DataValidationException; @@ -42,6 +44,8 @@ import java.util.LinkedList; import java.util.List; import java.util.concurrent.ExecutionException; +import static org.assertj.core.api.Assertions.assertThat; + @DaoSqlTest public class RelationServiceTest extends AbstractServiceTest { @@ -348,14 +352,14 @@ public class RelationServiceTest extends AbstractServiceTest { query.setFilters(Collections.singletonList(new RelationEntityTypeFilter(EntityRelation.CONTAINS_TYPE, Collections.singletonList(EntityType.ASSET)))); List relations = relationService.findByQuery(SYSTEM_TENANT_ID, query).get(); Assert.assertEquals(expected.size(), relations.size()); - for(EntityRelation r : expected){ + for (EntityRelation r : expected) { Assert.assertTrue(relations.contains(r)); } //Test from cache relations = relationService.findByQuery(SYSTEM_TENANT_ID, query).get(); Assert.assertEquals(expected.size(), relations.size()); - for(EntityRelation r : expected){ + for (EntityRelation r : expected) { Assert.assertTrue(relations.contains(r)); } } @@ -623,6 +627,51 @@ public class RelationServiceTest extends AbstractServiceTest { Assert.assertTrue(relations.contains(relationF)); } + @Test + public void testFindByPathQuery() throws Exception { + /* + A + └──[firstLevel, TO]→ B + └──[secondLevel, TO]→ C + ├──[thirdLevel, FROM]→ D + ├──[thirdLevel, FROM]→ E + └──[thirdLevel, FROM]→ F + */ + // rootEntity + AssetId assetA = new AssetId(Uuids.timeBased()); + // firstLevelEntity + AssetId assetB = new AssetId(Uuids.timeBased()); + // secondLevelEntity + AssetId assetC = new AssetId(Uuids.timeBased()); + // thirdLevelEntities + AssetId assetD = new AssetId(Uuids.timeBased()); + AssetId assetE = new AssetId(Uuids.timeBased()); + AssetId assetF = new AssetId(Uuids.timeBased()); + + EntityRelation firstLevelRelation = new EntityRelation(assetB, assetA, "firstLevel"); + EntityRelation secondLevelRelation = new EntityRelation(assetC, assetB, "secondLevel"); + EntityRelation thirdLevelRelation1 = new EntityRelation(assetC, assetD, "thirdLevel"); + EntityRelation thirdLevelRelation2 = new EntityRelation(assetC, assetE, "thirdLevel"); + EntityRelation thirdLevelRelation3 = new EntityRelation(assetC, assetF, "thirdLevel"); + + firstLevelRelation = saveRelation(firstLevelRelation); + secondLevelRelation = saveRelation(secondLevelRelation); + thirdLevelRelation1 = saveRelation(thirdLevelRelation1); + thirdLevelRelation2 = saveRelation(thirdLevelRelation2); + thirdLevelRelation3 = saveRelation(thirdLevelRelation3); + + List expectedRelations = List.of(thirdLevelRelation1, thirdLevelRelation2, thirdLevelRelation3); + + EntityRelationPathQuery relationPathQuery = new EntityRelationPathQuery(assetA, List.of( + new RelationPathLevel(EntitySearchDirection.TO, "firstLevel"), + new RelationPathLevel(EntitySearchDirection.TO, "secondLevel"), + new RelationPathLevel(EntitySearchDirection.FROM, "thirdLevel") + )); + List entityRelations = relationService.findByRelationPathQueryAsync(tenantId, relationPathQuery).get(); + + assertThat(expectedRelations).containsExactlyInAnyOrderElementsOf(entityRelations); + } + @Test public void testFindByQueryLargeHierarchyFetchAllWithUnlimLvl() throws Exception { AssetId rootAsset = new AssetId(Uuids.timeBased()); From fcba7004f94ab92521786bdcf3eeb61bd932ce79 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Thu, 2 Oct 2025 16:10:09 +0300 Subject: [PATCH 4/6] fixed tests to use new dynamic source configuration --- .../cf/CalculatedFieldIntegrationTest.java | 25 +++++-------- .../GeofencingCalculatedFieldStateTest.java | 17 +++------ .../data/cf/configuration/ArgumentTest.java | 2 +- .../ZoneGroupConfigurationTest.java | 4 +- .../dao/sql/relation/JpaRelationDao.java | 2 +- .../service/CalculatedFieldServiceTest.java | 37 +++++++++++-------- .../server/msa/cf/CalculatedFieldTest.java | 18 ++++----- 7 files changed, 48 insertions(+), 57 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java index ffd876575c..efcc38304a 100644 --- a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java @@ -36,7 +36,7 @@ import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfig 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.ReferencedEntityKey; -import org.thingsboard.server.common.data.cf.configuration.RelationQueryDynamicSourceConfiguration; +import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.cf.configuration.ScriptCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates; @@ -47,10 +47,12 @@ import org.thingsboard.server.common.data.id.AssetProfileId; import org.thingsboard.server.common.data.id.EntityId; 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.controller.CalculatedFieldControllerTest; import org.thingsboard.server.dao.service.DaoSqlTest; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.concurrent.TimeUnit; @@ -760,19 +762,13 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes // Zone groups: ATTRIBUTE on specific assets (one zone per group) ZoneGroupConfiguration allowedZonesGroup = new ZoneGroupConfiguration("zone", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false); - var allowedZoneDynamicSourceConfiguration = new RelationQueryDynamicSourceConfiguration(); - allowedZoneDynamicSourceConfiguration.setDirection(EntitySearchDirection.FROM); - allowedZoneDynamicSourceConfiguration.setRelationType("AllowedZone"); - allowedZoneDynamicSourceConfiguration.setMaxLevel(1); - allowedZoneDynamicSourceConfiguration.setFetchLastLevelOnly(true); + var allowedZoneDynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); + allowedZoneDynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.FROM, "AllowedZone"))); allowedZonesGroup.setRefDynamicSourceConfiguration(allowedZoneDynamicSourceConfiguration); ZoneGroupConfiguration restrictedZonesGroup = new ZoneGroupConfiguration("zone", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false); - var restrictedZoneDynamicSourceConfiguration = new RelationQueryDynamicSourceConfiguration(); - restrictedZoneDynamicSourceConfiguration.setDirection(EntitySearchDirection.FROM); - restrictedZoneDynamicSourceConfiguration.setRelationType("RestrictedZone"); - restrictedZoneDynamicSourceConfiguration.setMaxLevel(1); - restrictedZoneDynamicSourceConfiguration.setFetchLastLevelOnly(true); + var restrictedZoneDynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); + restrictedZoneDynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.FROM, "RestrictedZone"))); restrictedZonesGroup.setRefDynamicSourceConfiguration(restrictedZoneDynamicSourceConfiguration); cfg.setZoneGroups(Map.of("allowedZones", allowedZonesGroup, "restrictedZones", restrictedZonesGroup)); @@ -870,11 +866,8 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes cfg.setEntityCoordinates(new EntityCoordinates(ENTITY_ID_LATITUDE_ARGUMENT_KEY, ENTITY_ID_LONGITUDE_ARGUMENT_KEY)); var allowedZonesGroup = new ZoneGroupConfiguration("zone", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false); - var allowedZoneDynamicSourceConfiguration = new RelationQueryDynamicSourceConfiguration(); - allowedZoneDynamicSourceConfiguration.setDirection(EntitySearchDirection.FROM); - allowedZoneDynamicSourceConfiguration.setRelationType("AllowedZone"); - allowedZoneDynamicSourceConfiguration.setMaxLevel(1); - allowedZoneDynamicSourceConfiguration.setFetchLastLevelOnly(true); + var allowedZoneDynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); + allowedZoneDynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.FROM, "AllowedZone"))); allowedZonesGroup.setRefDynamicSourceConfiguration(allowedZoneDynamicSourceConfiguration); cfg.setZoneGroups(Map.of("allowedZones", allowedZonesGroup)); diff --git a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldStateTest.java b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldStateTest.java index a86af77555..1f1ee32df2 100644 --- a/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldStateTest.java +++ b/application/src/test/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldStateTest.java @@ -28,7 +28,7 @@ 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.RelationQueryDynamicSourceConfiguration; +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; @@ -41,6 +41,7 @@ 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; @@ -453,21 +454,15 @@ public class GeofencingCalculatedFieldStateTest { config.setEntityCoordinates(entityCoordinates); ZoneGroupConfiguration allowedZonesGroup = new ZoneGroupConfiguration("zone", reportStrategy, true); - var allowedZoneDynamicSourceConfiguration = new RelationQueryDynamicSourceConfiguration(); - allowedZoneDynamicSourceConfiguration.setDirection(EntitySearchDirection.TO); - allowedZoneDynamicSourceConfiguration.setRelationType("AllowedZone"); - allowedZoneDynamicSourceConfiguration.setMaxLevel(1); - allowedZoneDynamicSourceConfiguration.setFetchLastLevelOnly(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 RelationQueryDynamicSourceConfiguration(); - restrictedZoneDynamicSourceConfiguration.setDirection(EntitySearchDirection.TO); - restrictedZoneDynamicSourceConfiguration.setRelationType("RestrictedZone"); - restrictedZoneDynamicSourceConfiguration.setMaxLevel(1); - restrictedZoneDynamicSourceConfiguration.setFetchLastLevelOnly(true); + var restrictedZoneDynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); + restrictedZoneDynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.TO, "RestrictedZone"))); restrictedZonesGroup.setRefDynamicSourceConfiguration(restrictedZoneDynamicSourceConfiguration); restrictedZonesGroup.setRelationType("CurrentZone"); restrictedZonesGroup.setDirection(EntitySearchDirection.TO); diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/ArgumentTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/ArgumentTest.java index fd59317649..d1108e9eed 100644 --- a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/ArgumentTest.java +++ b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/ArgumentTest.java @@ -31,7 +31,7 @@ public class ArgumentTest { @Test void validateShouldReturnTrueIfDynamicSourceConfigurationIsNotNull() { var argument = new Argument(); - argument.setRefDynamicSourceConfiguration(new RelationQueryDynamicSourceConfiguration()); + argument.setRefDynamicSourceConfiguration(new RelationPathQueryDynamicSourceConfiguration()); assertThat(argument.hasDynamicSource()).isTrue(); } diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/geofencing/ZoneGroupConfigurationTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/geofencing/ZoneGroupConfigurationTest.java index 4eb822d93c..beb4639a31 100644 --- a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/geofencing/ZoneGroupConfigurationTest.java +++ b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/geofencing/ZoneGroupConfigurationTest.java @@ -23,7 +23,7 @@ import org.thingsboard.server.common.data.AttributeScope; 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.ReferencedEntityKey; -import org.thingsboard.server.common.data.cf.configuration.RelationQueryDynamicSourceConfiguration; +import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntitySearchDirection; @@ -100,7 +100,7 @@ public class ZoneGroupConfigurationTest { @Test void whenHasDynamicSourceCalled_shouldReturnTrueIfDynamicSourceConfigurationIsNotNull() { var zoneGroupConfiguration = new ZoneGroupConfiguration("perimeter", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false); - zoneGroupConfiguration.setRefDynamicSourceConfiguration(new RelationQueryDynamicSourceConfiguration()); + zoneGroupConfiguration.setRefDynamicSourceConfiguration(new RelationPathQueryDynamicSourceConfiguration()); assertThat(zoneGroupConfiguration.hasDynamicSource()).isTrue(); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java index 782c9b9d53..b2871313ed 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/relation/JpaRelationDao.java @@ -307,7 +307,7 @@ public class JpaRelationDao extends JpaAbstractDaoListeningExecutorService imple String sql = buildRelationPathSql(query); Object[] params = buildRelationPathParams(query); - log.info("[{}] relation path query: {}", tenantId, sql); + log.trace("[{}] relation path query: {}", tenantId, sql); return jdbcTemplate.queryForList(sql, params).stream() .map(row -> { diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java index 7f563ea436..4b830835fa 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/CalculatedFieldServiceTest.java @@ -31,7 +31,7 @@ import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfig 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.ReferencedEntityKey; -import org.thingsboard.server.common.data.cf.configuration.RelationQueryDynamicSourceConfiguration; +import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates; import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingCalculatedFieldConfiguration; @@ -39,15 +39,19 @@ import org.thingsboard.server.common.data.cf.configuration.geofencing.ZoneGroupC import org.thingsboard.server.common.data.id.EntityId; 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.cf.CalculatedFieldService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; +import java.util.ArrayList; +import java.util.List; import java.util.Map; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.mock; import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingReportStrategy.REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS; @DaoSqlTest @@ -112,10 +116,8 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { // Zone-group argument (ATTRIBUTE) — make it DYNAMIC so scheduling is enabled ZoneGroupConfiguration zoneGroupConfiguration = new ZoneGroupConfiguration("allowed", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false); - var dynamicSourceConfiguration = new RelationQueryDynamicSourceConfiguration(); - dynamicSourceConfiguration.setDirection(EntitySearchDirection.FROM); - dynamicSourceConfiguration.setMaxLevel(1); - dynamicSourceConfiguration.setRelationType(EntityRelation.CONTAINS_TYPE); + var dynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); + dynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.FROM, EntityRelation.CONTAINS_TYPE))); zoneGroupConfiguration.setRefDynamicSourceConfiguration(dynamicSourceConfiguration); cfg.setZoneGroups(Map.of("allowed", zoneGroupConfiguration)); @@ -150,19 +152,26 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { // Arrange a device Device device = createTestDevice(); - // Build a valid Geofencing configuration GeofencingCalculatedFieldConfiguration cfg = new GeofencingCalculatedFieldConfiguration(); // Coordinates: TS_LATEST, no dynamic source EntityCoordinates entityCoordinates = new EntityCoordinates("latitude", "longitude"); cfg.setEntityCoordinates(entityCoordinates); - // Zone-group argument (ATTRIBUTE) — make it DYNAMIC so scheduling is enabled + int maxRelationLevel = tbTenantProfileCache.get(tenantId) + .getDefaultProfileConfiguration() + .getMaxRelationLevelPerCfArgument(); + + // Zone-group argument (ATTRIBUTE) ZoneGroupConfiguration zoneGroupConfiguration = new ZoneGroupConfiguration( "allowed", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false); - var dynamicSourceConfiguration = new RelationQueryDynamicSourceConfiguration(); - dynamicSourceConfiguration.setDirection(EntitySearchDirection.FROM); - dynamicSourceConfiguration.setMaxLevel(Integer.MAX_VALUE); - dynamicSourceConfiguration.setRelationType(EntityRelation.CONTAINS_TYPE); + var dynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); + + List levels = new ArrayList<>(); + for (int i = 0; i < maxRelationLevel + 1; i++) { + levels.add(mock(RelationPathLevel.class)); + } + + dynamicSourceConfiguration.setLevels(levels); zoneGroupConfiguration.setRefDynamicSourceConfiguration(dynamicSourceConfiguration); cfg.setZoneGroups(Map.of("allowed", zoneGroupConfiguration)); @@ -195,10 +204,8 @@ public class CalculatedFieldServiceTest extends AbstractServiceTest { // Zone-group argument (ATTRIBUTE) — make it DYNAMIC so scheduling is enabled ZoneGroupConfiguration zoneGroupConfiguration = new ZoneGroupConfiguration( "allowed", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false); - var dynamicSourceConfiguration = new RelationQueryDynamicSourceConfiguration(); - dynamicSourceConfiguration.setDirection(EntitySearchDirection.FROM); - dynamicSourceConfiguration.setMaxLevel(1); - dynamicSourceConfiguration.setRelationType(EntityRelation.CONTAINS_TYPE); + var dynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); + dynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.FROM, EntityRelation.CONTAINS_TYPE))); zoneGroupConfiguration.setRefDynamicSourceConfiguration(dynamicSourceConfiguration); cfg.setZoneGroups(Map.of("allowed", zoneGroupConfiguration)); diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java index 7f2dfba937..5e8d367538 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/cf/CalculatedFieldTest.java @@ -33,7 +33,7 @@ import org.thingsboard.server.common.data.cf.configuration.ArgumentType; 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.ReferencedEntityKey; -import org.thingsboard.server.common.data.cf.configuration.RelationQueryDynamicSourceConfiguration; +import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.cf.configuration.ScriptCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates; @@ -50,10 +50,12 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.UserId; 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.msa.AbstractContainerTest; import org.thingsboard.server.msa.ui.utils.EntityPrototypes; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.concurrent.TimeUnit; @@ -366,19 +368,13 @@ public class CalculatedFieldTest extends AbstractContainerTest { // Dynamic groups via relations ZoneGroupConfiguration allowedZoneGroupConfiguration = new ZoneGroupConfiguration("zone", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false); - var allowedDynamicSourceConfiguration = new RelationQueryDynamicSourceConfiguration(); - allowedDynamicSourceConfiguration.setDirection(EntitySearchDirection.FROM); - allowedDynamicSourceConfiguration.setMaxLevel(1); - allowedDynamicSourceConfiguration.setFetchLastLevelOnly(true); - allowedDynamicSourceConfiguration.setRelationType("AllowedZone"); + var allowedDynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); + allowedDynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.FROM, "AllowedZone"))); allowedZoneGroupConfiguration.setRefDynamicSourceConfiguration(allowedDynamicSourceConfiguration); ZoneGroupConfiguration restrictedZoneGroupConfiguration = new ZoneGroupConfiguration("zone", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false); - var restrictedDynamicSourceConfiguration = new RelationQueryDynamicSourceConfiguration(); - restrictedDynamicSourceConfiguration.setDirection(EntitySearchDirection.FROM); - restrictedDynamicSourceConfiguration.setMaxLevel(1); - restrictedDynamicSourceConfiguration.setFetchLastLevelOnly(true); - restrictedDynamicSourceConfiguration.setRelationType("RestrictedZone"); + var restrictedDynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); + restrictedDynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.FROM, "RestrictedZone"))); restrictedZoneGroupConfiguration.setRefDynamicSourceConfiguration(restrictedDynamicSourceConfiguration); cfg.setZoneGroups(Map.of("allowedZones", allowedZoneGroupConfiguration, "restrictedZones", restrictedZoneGroupConfiguration)); From eea9e6bf6e43e38c5cd65ef5ca83db4ffb456ab9 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Fri, 3 Oct 2025 11:43:24 +0300 Subject: [PATCH 5/6] Removed RELATION_QUERY type --- ...tractCalculatedFieldProcessingService.java | 16 -- .../CFArgumentDynamicSourceType.java | 2 +- ...onPathQueryDynamicSourceConfiguration.java | 14 +- .../cf/configuration/RelationQueryBased.java | 29 --- ...lationQueryDynamicSourceConfiguration.java | 79 ------- ...onQueryDynamicSourceConfigurationTest.java | 214 ------------------ .../dao/relation/BaseRelationService.java | 7 + .../CalculatedFieldDataValidator.java | 6 +- 8 files changed, 17 insertions(+), 350 deletions(-) delete mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryBased.java delete mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfiguration.java delete mode 100644 common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfigurationTest.java diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java index bd3546082a..f1b39d87d8 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java @@ -27,7 +27,6 @@ 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.cf.configuration.RelationQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.Aggregation; @@ -164,21 +163,6 @@ public abstract class AbstractCalculatedFieldProcessingService { } var refDynamicSourceConfiguration = value.getRefDynamicSourceConfiguration(); return switch (refDynamicSourceConfiguration.getType()) { - case RELATION_QUERY -> { - var configuration = (RelationQueryDynamicSourceConfiguration) refDynamicSourceConfiguration; - if (configuration.isSimpleRelation()) { - yield switch (configuration.getDirection()) { - case FROM -> - Futures.transform(relationService.findByFromAndTypeAsync(tenantId, entityId, configuration.getRelationType(), RelationTypeGroup.COMMON), - configuration::resolveEntityIds, calculatedFieldCallbackExecutor); - case TO -> - Futures.transform(relationService.findByToAndTypeAsync(tenantId, entityId, configuration.getRelationType(), RelationTypeGroup.COMMON), - configuration::resolveEntityIds, calculatedFieldCallbackExecutor); - }; - } - yield Futures.transform(relationService.findByQuery(tenantId, configuration.toEntityRelationsQuery(entityId)), - configuration::resolveEntityIds, calculatedFieldCallbackExecutor); - } case RELATION_PATH_QUERY -> { var configuration = (RelationPathQueryDynamicSourceConfiguration) refDynamicSourceConfiguration; yield Futures.transform(relationService.findByRelationPathQueryAsync(tenantId, configuration.toRelationPathQuery(entityId)), diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CFArgumentDynamicSourceType.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CFArgumentDynamicSourceType.java index 35c6cdf562..dc52287f3e 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CFArgumentDynamicSourceType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/CFArgumentDynamicSourceType.java @@ -17,6 +17,6 @@ package org.thingsboard.server.common.data.cf.configuration; public enum CFArgumentDynamicSourceType { - RELATION_QUERY, RELATION_PATH_QUERY + RELATION_PATH_QUERY } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java index c7889be19d..dc92ff3685 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationPathQueryDynamicSourceConfiguration.java @@ -28,7 +28,7 @@ import java.util.List; import java.util.NoSuchElementException; @Data -public class RelationPathQueryDynamicSourceConfiguration implements CfArgumentDynamicSourceConfiguration, RelationQueryBased { +public class RelationPathQueryDynamicSourceConfiguration implements CfArgumentDynamicSourceConfiguration { private List levels; @@ -53,10 +53,11 @@ public class RelationPathQueryDynamicSourceConfiguration implements CfArgumentDy }; } - @Override - @JsonIgnore - public int getMaxLevel() { - return levels != null ? levels.size() : 0; + 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) { @@ -64,9 +65,6 @@ public class RelationPathQueryDynamicSourceConfiguration implements CfArgumentDy } private RelationPathLevel getLastLevel() { - if (CollectionsUtil.isEmpty(levels)) { - throw new NoSuchElementException(); - } return levels.get(levels.size() - 1); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryBased.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryBased.java deleted file mode 100644 index e3248b264f..0000000000 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryBased.java +++ /dev/null @@ -1,29 +0,0 @@ -/** - * 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 RelationQueryBased { - - int getMaxLevel(); - - default void validateMaxRelationLevel(String argumentName, int maxAllowedRelationLevel) { - if (getMaxLevel() > maxAllowedRelationLevel) { - throw new IllegalArgumentException("Max relation level is greater than configured " + - "maximum allowed relation level in tenant profile: " + maxAllowedRelationLevel + " for argument: " + argumentName); - } - } - -} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfiguration.java deleted file mode 100644 index 120d04d40d..0000000000 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfiguration.java +++ /dev/null @@ -1,79 +0,0 @@ -/** - * 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.StringUtils; -import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.relation.EntityRelation; -import org.thingsboard.server.common.data.relation.EntityRelationsQuery; -import org.thingsboard.server.common.data.relation.EntitySearchDirection; -import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter; -import org.thingsboard.server.common.data.relation.RelationsSearchParameters; - -import java.util.Collections; -import java.util.List; - -@Data -public class RelationQueryDynamicSourceConfiguration implements CfArgumentDynamicSourceConfiguration, RelationQueryBased { - - private int maxLevel; - private boolean fetchLastLevelOnly; - private EntitySearchDirection direction; - private String relationType; - - @Override - public CFArgumentDynamicSourceType getType() { - return CFArgumentDynamicSourceType.RELATION_QUERY; - } - - @Override - public void validate() { - if (maxLevel < 1) { - throw new IllegalArgumentException("Relation query dynamic source configuration max relation level can't be less than 1!"); - } - if (direction == null) { - throw new IllegalArgumentException("Relation query dynamic source configuration direction must be specified!"); - } - if (StringUtils.isBlank(relationType)) { - throw new IllegalArgumentException("Relation query dynamic source configuration relation type must be specified!"); - } - } - - @JsonIgnore - public boolean isSimpleRelation() { - return maxLevel == 1; - } - - public EntityRelationsQuery toEntityRelationsQuery(EntityId rootEntityId) { - if (isSimpleRelation()) { - throw new IllegalArgumentException("Entity relations query can't be created for a simple relation!"); - } - var entityRelationsQuery = new EntityRelationsQuery(); - entityRelationsQuery.setParameters(new RelationsSearchParameters(rootEntityId, direction, maxLevel, fetchLastLevelOnly)); - entityRelationsQuery.setFilters(Collections.singletonList(new RelationEntityTypeFilter(relationType, Collections.emptyList()))); - return entityRelationsQuery; - } - - public List resolveEntityIds(List relations) { - return switch (direction) { - case FROM -> relations.stream().map(EntityRelation::getTo).toList(); - case TO -> relations.stream().map(EntityRelation::getFrom).toList(); - }; - } - -} diff --git a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfigurationTest.java b/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfigurationTest.java deleted file mode 100644 index afd78e47f9..0000000000 --- a/common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/RelationQueryDynamicSourceConfigurationTest.java +++ /dev/null @@ -1,214 +0,0 @@ -/** - * 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 org.junit.jupiter.api.Test; -import org.junit.jupiter.api.extension.ExtendWith; -import org.junit.jupiter.params.ParameterizedTest; -import org.junit.jupiter.params.provider.NullAndEmptySource; -import org.junit.jupiter.params.provider.ValueSource; -import org.mockito.Mock; -import org.mockito.junit.jupiter.MockitoExtension; -import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.relation.EntityRelation; -import org.thingsboard.server.common.data.relation.EntitySearchDirection; -import org.thingsboard.server.common.data.relation.RelationEntityTypeFilter; -import org.thingsboard.server.common.data.relation.RelationsSearchParameters; - -import java.util.List; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.assertj.core.api.Assertions.assertThatCode; -import static org.assertj.core.api.Assertions.assertThatThrownBy; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.when; - -@ExtendWith(MockitoExtension.class) -public class RelationQueryDynamicSourceConfigurationTest { - - @Mock - EntityId rootEntityId; - - @Mock - EntityRelation rel1; - @Mock - EntityRelation rel2; - - @Test - void typeShouldBeRelationQuery() { - var cfg = new RelationQueryDynamicSourceConfiguration(); - assertThat(cfg.getType()).isEqualTo(CFArgumentDynamicSourceType.RELATION_QUERY); - } - - @Test - void validateShouldThrowWhenMaxLevelLessThanOne() { - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setMaxLevel(0); - cfg.setDirection(EntitySearchDirection.FROM); - cfg.setRelationType(EntityRelation.CONTAINS_TYPE); - - assertThatThrownBy(cfg::validate) - .isInstanceOf(IllegalArgumentException.class) - .hasMessage("Relation query dynamic source configuration max relation level can't be less than 1!"); - } - - @Test - void validateShouldThrowWhenMaxLevelGreaterThanMaxAllowedLevelFromTenantProfile() { - int maxAllowedRelationLevel = 2; - int argumentMaxRelationLevel = 3; - - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setMaxLevel(argumentMaxRelationLevel); - cfg.setDirection(EntitySearchDirection.FROM); - cfg.setRelationType(EntityRelation.CONTAINS_TYPE); - - String testRelationArgument = "testRelationArgument"; - assertThatThrownBy(() -> cfg.validateMaxRelationLevel(testRelationArgument, maxAllowedRelationLevel)) - .isInstanceOf(IllegalArgumentException.class) - .hasMessage("Max relation level is greater than configured " + - "maximum allowed relation level in tenant profile: " + maxAllowedRelationLevel + " for argument: " + testRelationArgument); - } - - @Test - void validateShouldPassValidationWhenMaxLevelLessThanMaxAllowedLevelFromTenantProfile() { - int maxAllowedRelationLevel = 5; - int argumentMaxRelationLevel = 2; - - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setMaxLevel(argumentMaxRelationLevel); - cfg.setDirection(EntitySearchDirection.FROM); - cfg.setRelationType(EntityRelation.CONTAINS_TYPE); - - String testRelationArgument = "testRelationArgument"; - assertThatCode(() -> cfg.validateMaxRelationLevel(testRelationArgument, maxAllowedRelationLevel)).doesNotThrowAnyException(); - } - - @Test - void validateShouldThrowWhenDirectionIsNull() { - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setMaxLevel(1); - cfg.setDirection(null); - cfg.setRelationType(EntityRelation.CONTAINS_TYPE); - - assertThatThrownBy(cfg::validate) - .isInstanceOf(IllegalArgumentException.class) - .hasMessage("Relation query dynamic source configuration direction must be specified!"); - } - - @ParameterizedTest - @ValueSource(strings = {" "}) - @NullAndEmptySource - void validateShouldThrowWhenRelationTypeIsNull(String relationType) { - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setMaxLevel(1); - cfg.setDirection(EntitySearchDirection.TO); - cfg.setRelationType(relationType); - - assertThatThrownBy(cfg::validate) - .isInstanceOf(IllegalArgumentException.class) - .hasMessage("Relation query dynamic source configuration relation type must be specified!"); - } - - @Test - void isSimpleRelationTrueWhenLevelIsOneAndEntityTypesEmptyOrNull() { - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setMaxLevel(1); - assertThat(cfg.isSimpleRelation()).isTrue(); - } - - @Test - void isSimpleRelationFalseWhenMaxLevelNotOne() { - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setMaxLevel(2); - assertThat(cfg.isSimpleRelation()).isFalse(); - } - - @Test - void toEntityRelationsQueryShouldThrowForSimpleRelation() { - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setMaxLevel(1); - cfg.setFetchLastLevelOnly(false); - cfg.setDirection(EntitySearchDirection.FROM); - cfg.setRelationType(EntityRelation.CONTAINS_TYPE); - - assertThatThrownBy(() -> cfg.toEntityRelationsQuery(rootEntityId)) - .isInstanceOf(IllegalArgumentException.class) - .hasMessage("Entity relations query can't be created for a simple relation!"); - } - - @Test - void toEntityRelationsQueryShouldBuildQueryForNonSimpleRelation() { - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setMaxLevel(2); - cfg.setFetchLastLevelOnly(true); - cfg.setDirection(EntitySearchDirection.TO); - cfg.setRelationType(EntityRelation.MANAGES_TYPE); - - var query = cfg.toEntityRelationsQuery(rootEntityId); - - assertThat(query).isNotNull(); - RelationsSearchParameters params = query.getParameters(); - assertThat(params).isNotNull(); - assertThat(params.getRootId()).isEqualTo(rootEntityId.getId()); - assertThat(params.getDirection()).isEqualTo(EntitySearchDirection.TO); - assertThat(params.getMaxLevel()).isEqualTo(2); - assertThat(params.isFetchLastLevelOnly()).isTrue(); - - assertThat(query.getFilters()).hasSize(1); - assertThat(query.getFilters().get(0)).isInstanceOf(RelationEntityTypeFilter.class); - RelationEntityTypeFilter filter = query.getFilters().get(0); - assertThat(filter.getRelationType()).isEqualTo(EntityRelation.MANAGES_TYPE); - } - - @Test - void resolveEntityIds_whenDirectionFROM_thenReturnsToIds() { - when(rel1.getTo()).thenReturn(mock(EntityId.class)); - when(rel2.getTo()).thenReturn(mock(EntityId.class)); - - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setDirection(EntitySearchDirection.FROM); - - var out = cfg.resolveEntityIds(List.of(rel1, rel2)); - - assertThat(out).containsExactly(rel1.getTo(), rel2.getTo()); - } - - @Test - void resolveEntityIds_whenDirectionTO_thenReturnsFromIds() { - when(rel1.getFrom()).thenReturn(mock(EntityId.class)); - when(rel2.getFrom()).thenReturn(mock(EntityId.class)); - - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setDirection(EntitySearchDirection.TO); - - var out = cfg.resolveEntityIds(List.of(rel1, rel2)); - - assertThat(out).containsExactly(rel1.getFrom(), rel2.getFrom()); - } - - @Test - void validateShouldPassForValidConfig() { - var cfg = new RelationQueryDynamicSourceConfiguration(); - cfg.setMaxLevel(2); - cfg.setFetchLastLevelOnly(false); - cfg.setDirection(EntitySearchDirection.FROM); - cfg.setRelationType(EntityRelation.CONTAINS_TYPE); - - assertThatCode(cfg::validate).doesNotThrowAnyException(); - } - -} diff --git a/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java b/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java index 0ec1257bcc..18d1806fe4 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/relation/BaseRelationService.java @@ -504,6 +504,13 @@ public class BaseRelationService implements RelationService { log.trace("Executing findByRelationPathQuery, tenantId [{}], relationPathQuery {}", tenantId, relationPathQuery); validateId(tenantId, id -> "Invalid tenant id: " + id); validate(relationPathQuery); + if (relationPathQuery.levels().size() == 1) { + RelationPathLevel relationPathLevel = relationPathQuery.levels().get(0); + return switch (relationPathLevel.direction()) { + case FROM -> findByFromAndTypeAsync(tenantId, relationPathQuery.rootEntityId(), relationPathLevel.relationType(), RelationTypeGroup.COMMON); + case TO -> findByToAndTypeAsync(tenantId, relationPathQuery.rootEntityId(), relationPathLevel.relationType(), RelationTypeGroup.COMMON); + }; + } return executor.submit(() -> relationDao.findByRelationPathQuery(tenantId, relationPathQuery)); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java index 85c30ca74e..67cf32191b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java @@ -19,7 +19,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.configuration.ArgumentsBasedCalculatedFieldConfiguration; -import org.thingsboard.server.common.data.cf.configuration.RelationQueryBased; +import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.cf.configuration.ScheduledUpdateSupportedCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; @@ -98,10 +98,10 @@ public class CalculatedFieldDataValidator extends DataValidator if (!(calculatedField.getConfiguration() instanceof ArgumentsBasedCalculatedFieldConfiguration argumentsBasedCfg)) { return; } - Map relationQueryBasedArguments = argumentsBasedCfg.getArguments().entrySet() + Map relationQueryBasedArguments = argumentsBasedCfg.getArguments().entrySet() .stream() .filter(entry -> entry.getValue().hasDynamicSource()) - .collect(Collectors.toMap(Map.Entry::getKey, entry -> (RelationQueryBased) entry.getValue().getRefDynamicSourceConfiguration())); + .collect(Collectors.toMap(Map.Entry::getKey, entry -> (RelationPathQueryDynamicSourceConfiguration) entry.getValue().getRefDynamicSourceConfiguration())); if (relationQueryBasedArguments.isEmpty()) { return; } From 95583a0703bef23ed36b0da532301931326c2953 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Fri, 3 Oct 2025 12:18:19 +0300 Subject: [PATCH 6/6] Added minor test to CalculatedFieldController and todo for geo CF state --- ...tractCalculatedFieldProcessingService.java | 1 + .../CalculatedFieldControllerTest.java | 66 ++++++++++++++++++- 2 files changed, 64 insertions(+), 3 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java index f1b39d87d8..fa10a49503 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java @@ -103,6 +103,7 @@ public abstract class AbstractCalculatedFieldProcessingService { 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()); } diff --git a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java index af43b34558..27622b347a 100644 --- a/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/CalculatedFieldControllerTest.java @@ -29,15 +29,23 @@ import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfig 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.ReferencedEntityKey; +import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; +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.ZoneGroupConfiguration; import org.thingsboard.server.common.data.id.DeviceId; +import org.thingsboard.server.common.data.relation.EntitySearchDirection; +import org.thingsboard.server.common.data.relation.RelationPathLevel; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.dao.service.DaoSqlTest; +import java.util.List; import java.util.Map; import static org.assertj.core.api.Assertions.assertThat; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; +import static org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingReportStrategy.REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS; @DaoSqlTest public class CalculatedFieldControllerTest extends AbstractControllerTest { @@ -84,7 +92,35 @@ public class CalculatedFieldControllerTest extends AbstractControllerTest { assertThat(savedCalculatedField.getEntityId()).isEqualTo(calculatedField.getEntityId()); assertThat(savedCalculatedField.getType()).isEqualTo(calculatedField.getType()); assertThat(savedCalculatedField.getName()).isEqualTo(calculatedField.getName()); - assertThat(savedCalculatedField.getConfiguration()).isEqualTo(getCalculatedFieldConfig()); + assertThat(savedCalculatedField.getConfiguration()).isEqualTo(getSimpleCalculatedFieldConfig()); + assertThat(savedCalculatedField.getVersion()).isEqualTo(1L); + + savedCalculatedField.setName("Test CF"); + + CalculatedField updatedCalculatedField = doPost("/api/calculatedField", savedCalculatedField, CalculatedField.class); + + assertThat(updatedCalculatedField.getName()).isEqualTo(savedCalculatedField.getName()); + assertThat(updatedCalculatedField.getVersion()).isEqualTo(savedCalculatedField.getVersion() + 1); + + doDelete("/api/calculatedField/" + savedCalculatedField.getId().getId().toString()) + .andExpect(status().isOk()); + } + + @Test + public void testSaveGeofencingCalculatedField() throws Exception { + Device testDevice = createDevice("Test device", "1234567890"); + CalculatedField calculatedField = getCalculatedField(testDevice.getId(), getGeofencingCalculatedFieldConfig()); + + CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); + + assertThat(savedCalculatedField).isNotNull(); + assertThat(savedCalculatedField.getId()).isNotNull(); + assertThat(savedCalculatedField.getCreatedTime()).isGreaterThan(0); + assertThat(savedCalculatedField.getTenantId()).isEqualTo(savedTenant.getId()); + assertThat(savedCalculatedField.getEntityId()).isEqualTo(calculatedField.getEntityId()); + assertThat(savedCalculatedField.getType()).isEqualTo(calculatedField.getType()); + assertThat(savedCalculatedField.getName()).isEqualTo(calculatedField.getName()); + assertThat(savedCalculatedField.getConfiguration()).isEqualTo(getGeofencingCalculatedFieldConfig()); assertThat(savedCalculatedField.getVersion()).isEqualTo(1L); savedCalculatedField.setName("Test CF"); @@ -128,17 +164,41 @@ public class CalculatedFieldControllerTest extends AbstractControllerTest { } private CalculatedField getCalculatedField(DeviceId deviceId) { + return getCalculatedField(deviceId, getSimpleCalculatedFieldConfig()); + } + + private CalculatedField getCalculatedField(DeviceId deviceId, CalculatedFieldConfiguration configuration) { CalculatedField calculatedField = new CalculatedField(); calculatedField.setEntityId(deviceId); calculatedField.setType(CalculatedFieldType.SIMPLE); calculatedField.setName("Test Calculated Field"); calculatedField.setConfigurationVersion(1); - calculatedField.setConfiguration(getCalculatedFieldConfig()); + calculatedField.setConfiguration(configuration); calculatedField.setVersion(1L); return calculatedField; } - private CalculatedFieldConfiguration getCalculatedFieldConfig() { + private CalculatedFieldConfiguration getGeofencingCalculatedFieldConfig() { + var config = new GeofencingCalculatedFieldConfiguration(); + + var refDynamicSourceConfiguration = new RelationPathQueryDynamicSourceConfiguration(); + refDynamicSourceConfiguration.setLevels(List.of(new RelationPathLevel(EntitySearchDirection.TO, "FromSafeArea"))); + + var zoneGroupConfiguration = new ZoneGroupConfiguration("perimeter", REPORT_TRANSITION_EVENTS_AND_PRESENCE_STATUS, false); + zoneGroupConfiguration.setRefDynamicSourceConfiguration(refDynamicSourceConfiguration); + + Output output = new Output(); + output.setType(OutputType.TIME_SERIES); + + config.setEntityCoordinates(new EntityCoordinates("latitide", "longitude")); + config.setZoneGroups(Map.of("safeArea", zoneGroupConfiguration)); + config.setScheduledUpdateEnabled(false); + config.setOutput(output); + + return config; + } + + private CalculatedFieldConfiguration getSimpleCalculatedFieldConfig() { SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); Argument argument = new Argument();