Browse Source

Simplified impl of geofencing zone state & resolved TODOs

pull/13920/head
dshvaika 1 year ago
parent
commit
5e9921905f
  1. 19
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  2. 2
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  3. 1
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
  4. 13
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingCalculatedFieldState.java
  5. 21
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/GeofencingZoneState.java
  6. 9
      application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java
  7. 4
      common/proto/src/main/proto/queue.proto

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

@ -232,10 +232,16 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
CalculatedFieldCtx cfCtx = msg.getCfCtx(); CalculatedFieldCtx cfCtx = msg.getCfCtx();
CalculatedFieldId cfId = cfCtx.getCfId(); CalculatedFieldId cfId = cfCtx.getCfId();
log.debug("[{}][{}] Processing CF check for updates msg.", entityId, cfId); log.debug("[{}][{}] Processing CF check for updates msg.", entityId, cfId);
CalculatedFieldState currentState = states.get(cfId);
try { try {
var state = updateStateFromDb(cfCtx); var stateFromDb = getStateFromDb(cfCtx);
if (state.isSizeOk()) { if (currentState.equals(stateFromDb)) {
processStateIfReady(cfCtx, Collections.singletonList(cfId), state, null, null, msg.getCallback()); 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 { } else {
throw new RuntimeException(cfCtx.getSizeExceedsLimitMessage()); throw new RuntimeException(cfCtx.getSizeExceedsLimitMessage());
} }
@ -298,6 +304,12 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
} }
private CalculatedFieldState updateStateFromDb(CalculatedFieldCtx ctx) throws InterruptedException, ExecutionException, TimeoutException { 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<CalculatedFieldState> stateFuture = cfService.fetchStateFromDb(ctx, entityId); ListenableFuture<CalculatedFieldState> stateFuture = cfService.fetchStateFromDb(ctx, entityId);
// Ugly but necessary. We do not expect to often fetch data from DB. Only once per <Entity, CalculatedField> pair lifetime. // Ugly but necessary. We do not expect to often fetch data from DB. Only once per <Entity, CalculatedField> pair lifetime.
// This call happens while processing the CF pack from the queue consumer. So the timeout should be relatively low. // 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. // but this will significantly complicate the code.
CalculatedFieldState state = stateFuture.get(1, TimeUnit.MINUTES); CalculatedFieldState state = stateFuture.get(1, TimeUnit.MINUTES);
state.checkStateSize(new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId), ctx.getMaxStateSize()); state.checkStateSize(new CalculatedFieldEntityCtxId(tenantId, ctx.getCfId(), entityId), ctx.getMaxStateSize());
states.put(ctx.getCfId(), state);
return state; return state;
} }

2
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java

@ -409,7 +409,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
if (existingTask != null) { if (existingTask != null) {
existingTask.cancel(false); existingTask.cancel(false);
String reason = cfDeleted ? "removal" : "update"; 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);
} }
} }

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

@ -305,7 +305,6 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP
} }
private ListenableFuture<ArgumentEntry> fetchGeofencingKvEntry(TenantId tenantId, List<EntityId> geofencingEntities, Argument argument) { private ListenableFuture<ArgumentEntry> fetchGeofencingKvEntry(TenantId tenantId, List<EntityId> geofencingEntities, Argument argument) {
// TODO: Should we handle any other case?
if (argument.getRefEntityKey().getType() != ArgumentType.ATTRIBUTE) { if (argument.getRefEntityKey().getType() != ArgumentType.ATTRIBUTE) {
throw new IllegalStateException("Unsupported argument key type: " + argument.getRefEntityKey().getType()); throw new IllegalStateException("Unsupported argument key type: " + argument.getRefEntityKey().getType());
} }

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

@ -104,7 +104,6 @@ public class GeofencingCalculatedFieldState implements CalculatedFieldState {
} }
if (entryUpdated) { if (entryUpdated) {
stateUpdated = true; stateUpdated = true;
updateLastUpdateTimestamp(newEntry);
} }
} }
return stateUpdated; 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<CalculatedFieldResult> updateGeofencingZonesState(CalculatedFieldCtx ctx, boolean restricted) { private List<CalculatedFieldResult> updateGeofencingZonesState(CalculatedFieldCtx ctx, boolean restricted) {
var results = new ArrayList<CalculatedFieldResult>(); var results = new ArrayList<CalculatedFieldResult>();
long stateSwitchTime = System.currentTimeMillis();
double latitude = (double) arguments.get(ENTITY_ID_LATITUDE_ARGUMENT_KEY).getValue(); double latitude = (double) arguments.get(ENTITY_ID_LATITUDE_ARGUMENT_KEY).getValue();
double longitude = (double) arguments.get(ENTITY_ID_LONGITUDE_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()) { for (var zoneEntry : zonesEntry.getZoneStates().entrySet()) {
GeofencingZoneState state = zoneEntry.getValue(); GeofencingZoneState state = zoneEntry.getValue();
String event = state.evaluate(entityCoordinates, stateSwitchTime); String event = state.evaluate(entityCoordinates);
ObjectNode stateNode = JacksonUtil.newObjectNode(); ObjectNode stateNode = JacksonUtil.newObjectNode();
stateNode.put("entityId", ctx.getEntityId().toString()); stateNode.put("entityId", ctx.getEntityId().toString());
stateNode.put("zoneId", state.getZoneId().toString()); stateNode.put("zoneId", state.getZoneId().toString());

21
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; package org.thingsboard.server.service.cf.ctx.state;
import lombok.Data; import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.geo.Coordinates; import org.thingsboard.common.util.geo.Coordinates;
import org.thingsboard.common.util.geo.PerimeterDefinition; import org.thingsboard.common.util.geo.PerimeterDefinition;
import org.thingsboard.rule.engine.geo.EntityGeofencingState;
import org.thingsboard.rule.engine.util.GpsGeofencingEvents; import org.thingsboard.rule.engine.util.GpsGeofencingEvents;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.EntityIdFactory;
@ -39,7 +39,8 @@ public class GeofencingZoneState {
private Long version; private Long version;
private PerimeterDefinition perimeterDefinition; private PerimeterDefinition perimeterDefinition;
private EntityGeofencingState state; @EqualsAndHashCode.Exclude
private Boolean inside;
public GeofencingZoneState(EntityId zoneId, KvEntry entry) { public GeofencingZoneState(EntityId zoneId, KvEntry entry) {
this.zoneId = zoneId; this.zoneId = zoneId;
@ -59,7 +60,9 @@ public class GeofencingZoneState {
this.ts = proto.getTs(); this.ts = proto.getTs();
this.version = proto.getVersion(); this.version = proto.getVersion();
this.perimeterDefinition = JacksonUtil.fromString(proto.getPerimeterDefinition(), PerimeterDefinition.class); 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) { public boolean update(GeofencingZoneState newZoneState) {
@ -72,20 +75,16 @@ public class GeofencingZoneState {
this.version = newVersion; this.version = newVersion;
this.perimeterDefinition = newZoneState.getPerimeterDefinition(); this.perimeterDefinition = newZoneState.getPerimeterDefinition();
// TODO: should we reinitialize state if zone changed? // TODO: should we reinitialize state if zone changed?
// this.inside = null;
return true; return true;
} }
return false; return false;
} }
public String evaluate(Coordinates entityCoordinates, long currentTs) { public String evaluate(Coordinates entityCoordinates) {
boolean inside = perimeterDefinition.checkMatches(entityCoordinates); boolean inside = perimeterDefinition.checkMatches(entityCoordinates);
if (state == null) { if (this.inside == null || this.inside != inside) {
state = new EntityGeofencingState(inside, ts, false); this.inside = inside;
}
if (state.getStateSwitchTime() == 0L || state.isInside() != inside) {
state.setInside(inside);
state.setStateSwitchTime(currentTs);
state.setStayed(false);
return inside ? GpsGeofencingEvents.ENTERED : GpsGeofencingEvents.LEFT; return inside ? GpsGeofencingEvents.ENTERED : GpsGeofencingEvents.LEFT;
} }
return inside ? GpsGeofencingEvents.INSIDE : GpsGeofencingEvents.OUTSIDE; return inside ? GpsGeofencingEvents.INSIDE : GpsGeofencingEvents.OUTSIDE;

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

@ -16,7 +16,6 @@
package org.thingsboard.server.utils; package org.thingsboard.server.utils;
import org.thingsboard.common.util.JacksonUtil; 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.StringUtils;
import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CalculatedFieldId;
@ -139,12 +138,8 @@ public class CalculatedFieldUtils {
.setTs(zoneState.getTs()) .setTs(zoneState.getTs())
.setVersion(zoneState.getVersion()) .setVersion(zoneState.getVersion())
.setPerimeterDefinition(JacksonUtil.toString(zoneState.getPerimeterDefinition())); .setPerimeterDefinition(JacksonUtil.toString(zoneState.getPerimeterDefinition()));
if (zoneState.getState() != null) { if (zoneState.getInside() != null) {
EntityGeofencingState state = zoneState.getState(); builder.setInside(zoneState.getInside());
builder.setInside(state.isInside())
.setStayed(state.isStayed())
.setStateSwitchTime(state.getStateSwitchTime());
} }
return builder.build(); return builder.build();
} }

4
common/proto/src/main/proto/queue.proto

@ -905,9 +905,7 @@ message GeofencingZoneProto {
int64 ts = 2; int64 ts = 2;
string perimeterDefinition = 3; string perimeterDefinition = 3;
int64 version = 4; int64 version = 4;
bool inside = 5; optional bool inside = 5;
int64 stateSwitchTime = 6;
bool stayed = 7;
} }
message GeofencingArgumentProto { message GeofencingArgumentProto {

Loading…
Cancel
Save