From 8a3b8410ceccb0d77865ebd4996d13fc3d7a6c3b Mon Sep 17 00:00:00 2001 From: VIacheslavKlimov Date: Thu, 9 Oct 2025 16:55:56 +0300 Subject: [PATCH] Properly close CF state on removal --- ...CalculatedFieldEntityMessageProcessor.java | 33 ++++++++++++++----- ...alculatedFieldManagerMessageProcessor.java | 2 +- .../AbstractCalculatedFieldStateService.java | 2 +- .../cf/CalculatedFieldStateService.java | 2 +- .../ctx/state/BaseCalculatedFieldState.java | 6 +++- .../cf/ctx/state/CalculatedFieldState.java | 6 +++- .../alarm/AlarmCalculatedFieldState.java | 23 +++++++------ 7 files changed, 51 insertions(+), 23 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 3e9502bfe9..ccbdcc6b33 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 @@ -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); 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 f0e8aa7906..14713d79b2 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 @@ -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()); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java index c8b99afd9e..e673577742 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldStateService.java +++ b/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); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java index d0b34f18e8..10276ac421 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldStateService.java +++ b/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 partitions); 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 153b83f8d4..e442964280 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 @@ -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) { 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 ff94206220..cc7188e7a5 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 @@ -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 getArguments(); long getLatestTimestamp(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java index bc9ed4e32e..90999e5bbd 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/alarm/AlarmCalculatedFieldState.java +++ b/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 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 update(Map 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 performCalculation(Map 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); }