Browse Source

Properly close CF state on removal

pull/14193/head
VIacheslavKlimov 12 months ago
parent
commit
8a3b8410ce
  1. 33
      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. 2
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java
  4. 2
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java
  5. 6
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java
  6. 6
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java
  7. 23
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java

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

@ -104,6 +104,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
"[{}][{}] Stopping entity actor due to change partition event." :
"[{}][{}] Stopping entity actor.",
tenantId, entityId);
states.values().forEach(this::closeState);
states.clear();
actorCtx.stop(actorCtx.getSelf());
}
@ -123,7 +124,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
state.setPartition(msg.getPartition());
states.put(cfId, state);
} else {
states.remove(cfId);
removeState(cfId);
}
}
@ -141,7 +142,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
var ctx = msg.getCtx();
CalculatedFieldState state;
if (msg.getStateAction() == StateAction.RECREATE) {
states.remove(ctx.getCfId());
removeState(ctx.getCfId());
state = null;
} else {
state = states.get(ctx.getCfId());
@ -195,14 +196,14 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
msg.getCallback().onSuccess();
} else {
MultipleTbCallback multipleTbCallback = new MultipleTbCallback(states.size(), msg.getCallback());
states.forEach((cfId, state) -> cfStateService.removeState(new CalculatedFieldEntityCtxId(tenantId, cfId, entityId), multipleTbCallback));
states.forEach((cfId, state) -> cfStateService.deleteState(new CalculatedFieldEntityCtxId(tenantId, cfId, entityId), multipleTbCallback));
actorCtx.stop(actorCtx.getSelf());
}
} else {
var cfId = new CalculatedFieldId(msg.getEntityId().getId());
var state = states.remove(cfId);
var state = removeState(cfId);
if (state != null) {
cfStateService.removeState(new CalculatedFieldEntityCtxId(tenantId, cfId, entityId), msg.getCallback());
cfStateService.deleteState(new CalculatedFieldEntityCtxId(tenantId, cfId, entityId), msg.getCallback());
} else {
msg.getCallback().onSuccess();
}
@ -423,14 +424,30 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
if (state.isSizeOk()) {
cfStateService.persistState(ctxId, state, callback);
} else {
removeStateAndRaiseSizeException(ctxId, CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).errorMessage(ctx.getSizeExceedsLimitMessage()).build(), callback);
deleteStateAndRaiseSizeException(ctxId, CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).errorMessage(ctx.getSizeExceedsLimitMessage()).build(), callback);
}
}
}
private void removeStateAndRaiseSizeException(CalculatedFieldEntityCtxId ctxId, CalculatedFieldException ex, TbCallback callback) throws CalculatedFieldException {
private CalculatedFieldState removeState(CalculatedFieldId cfId) {
CalculatedFieldState state = states.remove(cfId);
closeState(state);
return state;
}
private void closeState(CalculatedFieldState state) {
if (state != null) {
try {
state.close();
} catch (Exception e) {
log.warn("[{}][{}] Failed to close CF state", tenantId, state.getEntityId(), e);
}
}
}
private void deleteStateAndRaiseSizeException(CalculatedFieldEntityCtxId ctxId, CalculatedFieldException ex, TbCallback callback) throws CalculatedFieldException {
// We remove the state, but remember that it is over-sized in a local map.
cfStateService.removeState(ctxId, new TbCallback() {
cfStateService.deleteState(ctxId, new TbCallback() {
@Override
public void onSuccess() {
callback.onFailure(ex);

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

@ -141,7 +141,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
log.debug("Pushing CF state restore msg to specific actor [{}]", msg.getId().entityId());
getOrCreateActor(msg.getId().entityId()).tell(msg);
} else {
cfStateService.removeState(msg.getId(), msg.getCallback());
cfStateService.deleteState(msg.getId(), msg.getCallback());
}
}

2
application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java

@ -62,7 +62,7 @@ public abstract class AbstractCalculatedFieldStateService implements CalculatedF
protected abstract void doPersist(CalculatedFieldEntityCtxId stateId, CalculatedFieldStateProto stateMsgProto, TbCallback callback);
@Override
public final void removeState(CalculatedFieldEntityCtxId stateId, TbCallback callback) {
public final void deleteState(CalculatedFieldEntityCtxId stateId, TbCallback callback) {
doRemove(stateId, callback);
}

2
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java

@ -33,7 +33,7 @@ public interface CalculatedFieldStateService {
void persistState(CalculatedFieldEntityCtxId stateId, CalculatedFieldState state, TbCallback callback) throws CalculatedFieldStateException;
void removeState(CalculatedFieldEntityCtxId stateId, TbCallback callback);
void deleteState(CalculatedFieldEntityCtxId stateId, TbCallback callback);
void restore(QueueKey queueKey, Set<TopicPartitionInfo> partitions);

6
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java

@ -23,13 +23,14 @@ import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import org.thingsboard.server.utils.CalculatedFieldUtils;
import java.io.Closeable;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@Getter
public abstract class BaseCalculatedFieldState implements CalculatedFieldState {
public abstract class BaseCalculatedFieldState implements CalculatedFieldState, Closeable {
protected final EntityId entityId;
protected CalculatedFieldCtx ctx;
@ -117,6 +118,9 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState {
}
}
@Override
public void close() {}
protected void validateNewEntry(String key, ArgumentEntry newEntry) {}
private void updateLastUpdateTimestamp(ArgumentEntry entry) {

6
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java

@ -22,6 +22,7 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.actors.TbActorRef;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.service.cf.CalculatedFieldResult;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
@ -29,6 +30,7 @@ import org.thingsboard.server.service.cf.ctx.state.alarm.AlarmCalculatedFieldSta
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState;
import java.io.Closeable;
import java.util.Map;
import static org.thingsboard.server.utils.CalculatedFieldUtils.toSingleValueArgumentProto;
@ -40,11 +42,13 @@ import static org.thingsboard.server.utils.CalculatedFieldUtils.toSingleValueArg
@Type(value = GeofencingCalculatedFieldState.class, name = "GEOFENCING"),
@Type(value = AlarmCalculatedFieldState.class, name = "ALARM")
})
public interface CalculatedFieldState {
public interface CalculatedFieldState extends Closeable {
@JsonIgnore
CalculatedFieldType getType();
EntityId getEntityId();
Map<String, ArgumentEntry> getArguments();
long getLatestTimestamp();

23
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java

@ -90,6 +90,8 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState {
private Alarm currentAlarm;
private boolean initialFetchDone;
// TODO: deprecate device profile node, describe the differences and improvements
public AlarmCalculatedFieldState(EntityId entityId) {
super(entityId);
}
@ -107,7 +109,7 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState {
}
@Override
public void init() { // todo: properly close state!
public void init() {
super.init();
AtomicBoolean reevalNeeded = new AtomicBoolean(false);
Map<AlarmSeverity, AlarmRule> createRules = configuration.getCreateRules();
@ -143,7 +145,6 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState {
AlarmEvalResult evalResult = state.reeval(System.currentTimeMillis());
if (evalResult.getStatus() == TRUE || evalResult.getStatus() == NOT_YET_TRUE) {
ScheduledFuture<?> future = ctx.scheduleReevaluation(evalResult.getLeftDuration(), actorCtx);
// TODO: use single task for multiple durations if durations are close enough. but be careful when cancelling the task in one of the states
if (future != null) {
state.setDurationCheckFuture(future);
}
@ -167,17 +168,21 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState {
return ruleState;
}
@Override
public Map<String, ArgumentEntry> update(Map<String, ArgumentEntry> argumentValues, CalculatedFieldCtx ctx) {
return super.update(argumentValues, ctx);
}
@Override
public void reset() {
super.reset();
configuration = null;
}
@Override
public void close() {
super.close();
for (AlarmRuleState state : createRuleStates.values()) {
clearState(state);
}
clearState(clearRuleState);
}
@Override
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) {
initCurrentAlarm(ctx);
@ -186,10 +191,8 @@ public class AlarmCalculatedFieldState extends BaseCalculatedFieldState {
boolean newEvent = !updatedArgs.isEmpty();
AlarmEvalResult evalResult = state.eval(newEvent, ctx);
if (evalResult.getStatus() == NOT_YET_TRUE && evalResult.getLeftDuration() > 0) {
// rounding up to the closest second
// long leftDuration = (long) Math.ceil(evalResult.getLeftDuration() / 1000.0) * 1000;
long leftDuration = evalResult.getLeftDuration();
ScheduledFuture<?> future = ctx.scheduleReevaluation(leftDuration, actorCtx); // TODO: use single task for multiple durations if durations are close enough. but be careful when cancelling the task in one of the states
ScheduledFuture<?> future = ctx.scheduleReevaluation(leftDuration, actorCtx);
if (future != null) {
state.setDurationCheckFuture(future);
}

Loading…
Cancel
Save