From 5e9921905f1928d416874bd47071b9490163bfd2 Mon Sep 17 00:00:00 2001 From: dshvaika Date: Mon, 4 Aug 2025 14:41:46 +0300 Subject: [PATCH] Simplified impl of geofencing zone state & resolved TODOs --- ...CalculatedFieldEntityMessageProcessor.java | 19 +++++++++++++---- ...alculatedFieldManagerMessageProcessor.java | 2 +- ...faultCalculatedFieldProcessingService.java | 1 - .../state/GeofencingCalculatedFieldState.java | 13 +----------- .../cf/ctx/state/GeofencingZoneState.java | 21 +++++++++---------- .../server/utils/CalculatedFieldUtils.java | 9 ++------ common/proto/src/main/proto/queue.proto | 4 +--- 7 files changed, 30 insertions(+), 39 deletions(-) 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 d41feac542..95b705adfc 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 @@ -232,10 +232,16 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM CalculatedFieldCtx cfCtx = msg.getCfCtx(); CalculatedFieldId cfId = cfCtx.getCfId(); log.debug("[{}][{}] Processing CF check for updates msg.", entityId, cfId); + CalculatedFieldState currentState = states.get(cfId); try { - var state = updateStateFromDb(cfCtx); - if (state.isSizeOk()) { - processStateIfReady(cfCtx, Collections.singletonList(cfId), state, null, null, msg.getCallback()); + var stateFromDb = getStateFromDb(cfCtx); + if (currentState.equals(stateFromDb)) { + log.debug("[{}][{}] CF state is up-to-date.", entityId, cfId); + return; + } + states.put(cfId, stateFromDb); + if (stateFromDb.isSizeOk()) { + processStateIfReady(cfCtx, Collections.singletonList(cfId), stateFromDb, null, null, msg.getCallback()); } else { throw new RuntimeException(cfCtx.getSizeExceedsLimitMessage()); } @@ -298,6 +304,12 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } private CalculatedFieldState updateStateFromDb(CalculatedFieldCtx ctx) throws InterruptedException, ExecutionException, TimeoutException { + CalculatedFieldState stateFromDb = getStateFromDb(ctx); + states.put(ctx.getCfId(), stateFromDb); + return stateFromDb; + } + + private CalculatedFieldState getStateFromDb(CalculatedFieldCtx ctx) throws InterruptedException, ExecutionException, TimeoutException { ListenableFuture stateFuture = cfService.fetchStateFromDb(ctx, entityId); // Ugly but necessary. We do not expect to often fetch data from DB. Only once per pair lifetime. // This call happens while processing the CF pack from the queue consumer. So the timeout should be relatively low. @@ -305,7 +317,6 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM // but this will significantly complicate the code. CalculatedFieldState state = stateFuture.get(1, TimeUnit.MINUTES); state.checkStateSize(new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId), ctx.getMaxStateSize()); - states.put(ctx.getCfId(), state); return state; } 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 f557a632b4..de6c38b9b9 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 @@ -409,7 +409,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware if (existingTask != null) { existingTask.cancel(false); String reason = cfDeleted ? "removal" : "update"; - log.debug("[{}][{}] Cancelled check for update task for CF due to: " + reason + "!", tenantId, cfId); + log.debug("[{}][{}] Cancelled check for update task due to CF " + reason + "!", tenantId, cfId); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java index 0284428c46..75ee581274 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java @@ -305,7 +305,6 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP } private ListenableFuture fetchGeofencingKvEntry(TenantId tenantId, List geofencingEntities, Argument argument) { - // TODO: Should we handle any other case? if (argument.getRefEntityKey().getType() != ArgumentType.ATTRIBUTE) { throw new IllegalStateException("Unsupported argument key type: " + argument.getRefEntityKey().getType()); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldState.java index bc7460784b..0d6772b676 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldState.java @@ -104,7 +104,6 @@ public class GeofencingCalculatedFieldState implements CalculatedFieldState { } if (entryUpdated) { stateUpdated = true; - updateLastUpdateTimestamp(newEntry); } } return stateUpdated; @@ -138,18 +137,8 @@ public class GeofencingCalculatedFieldState implements CalculatedFieldState { } } - private void updateLastUpdateTimestamp(ArgumentEntry entry) { - long newTs = this.latestTimestamp; - if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { - newTs = singleValueArgumentEntry.getTs(); - } - this.latestTimestamp = Math.max(this.latestTimestamp, newTs); - } - - // TODO: Ensure all cases are covered based on rule node logic. private List updateGeofencingZonesState(CalculatedFieldCtx ctx, boolean restricted) { var results = new ArrayList(); - long stateSwitchTime = System.currentTimeMillis(); double latitude = (double) arguments.get(ENTITY_ID_LATITUDE_ARGUMENT_KEY).getValue(); double longitude = (double) arguments.get(ENTITY_ID_LONGITUDE_ARGUMENT_KEY).getValue(); @@ -159,7 +148,7 @@ public class GeofencingCalculatedFieldState implements CalculatedFieldState { for (var zoneEntry : zonesEntry.getZoneStates().entrySet()) { GeofencingZoneState state = zoneEntry.getValue(); - String event = state.evaluate(entityCoordinates, stateSwitchTime); + String event = state.evaluate(entityCoordinates); ObjectNode stateNode = JacksonUtil.newObjectNode(); stateNode.put("entityId", ctx.getEntityId().toString()); stateNode.put("zoneId", state.getZoneId().toString()); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingZoneState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingZoneState.java index abee7cabb6..9e27907b73 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingZoneState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingZoneState.java @@ -16,10 +16,10 @@ package org.thingsboard.server.service.cf.ctx.state; import lombok.Data; +import lombok.EqualsAndHashCode; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.geo.Coordinates; import org.thingsboard.common.util.geo.PerimeterDefinition; -import org.thingsboard.rule.engine.geo.EntityGeofencingState; import org.thingsboard.rule.engine.util.GpsGeofencingEvents; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; @@ -39,7 +39,8 @@ public class GeofencingZoneState { private Long version; private PerimeterDefinition perimeterDefinition; - private EntityGeofencingState state; + @EqualsAndHashCode.Exclude + private Boolean inside; public GeofencingZoneState(EntityId zoneId, KvEntry entry) { this.zoneId = zoneId; @@ -59,7 +60,9 @@ public class GeofencingZoneState { this.ts = proto.getTs(); this.version = proto.getVersion(); this.perimeterDefinition = JacksonUtil.fromString(proto.getPerimeterDefinition(), PerimeterDefinition.class); - this.state = new EntityGeofencingState(proto.getInside(), proto.getStateSwitchTime(), proto.getStayed()); + if (proto.hasInside()) { + this.inside = proto.getInside(); + } } public boolean update(GeofencingZoneState newZoneState) { @@ -72,20 +75,16 @@ public class GeofencingZoneState { this.version = newVersion; this.perimeterDefinition = newZoneState.getPerimeterDefinition(); // TODO: should we reinitialize state if zone changed? + // this.inside = null; return true; } return false; } - public String evaluate(Coordinates entityCoordinates, long currentTs) { + public String evaluate(Coordinates entityCoordinates) { boolean inside = perimeterDefinition.checkMatches(entityCoordinates); - if (state == null) { - state = new EntityGeofencingState(inside, ts, false); - } - if (state.getStateSwitchTime() == 0L || state.isInside() != inside) { - state.setInside(inside); - state.setStateSwitchTime(currentTs); - state.setStayed(false); + if (this.inside == null || this.inside != inside) { + this.inside = inside; return inside ? GpsGeofencingEvents.ENTERED : GpsGeofencingEvents.LEFT; } return inside ? GpsGeofencingEvents.INSIDE : GpsGeofencingEvents.OUTSIDE; diff --git a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java index 302ced0df1..370f7883f0 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java @@ -16,7 +16,6 @@ package org.thingsboard.server.utils; import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.rule.engine.geo.EntityGeofencingState; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.id.CalculatedFieldId; @@ -139,12 +138,8 @@ public class CalculatedFieldUtils { .setTs(zoneState.getTs()) .setVersion(zoneState.getVersion()) .setPerimeterDefinition(JacksonUtil.toString(zoneState.getPerimeterDefinition())); - if (zoneState.getState() != null) { - EntityGeofencingState state = zoneState.getState(); - builder.setInside(state.isInside()) - .setStayed(state.isStayed()) - .setStateSwitchTime(state.getStateSwitchTime()); - + if (zoneState.getInside() != null) { + builder.setInside(zoneState.getInside()); } return builder.build(); } diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 885bad2cca..de55cd549d 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -905,9 +905,7 @@ message GeofencingZoneProto { int64 ts = 2; string perimeterDefinition = 3; int64 version = 4; - bool inside = 5; - int64 stateSwitchTime = 6; - bool stayed = 7; + optional bool inside = 5; } message GeofencingArgumentProto {