diff --git a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java index 16562dbbe..281830dc9 100644 --- a/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java +++ b/impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java @@ -53,7 +53,7 @@ public class WorkflowMutableInstance implements WorkflowInstance { protected AtomicReference> futureRef = new AtomicReference<>(); protected Instant completedAt; - private Map additionalObjects = new ConcurrentHashMap<>(); + protected Map additionalObjects = new ConcurrentHashMap<>(); protected final Map iterationsMap = new ConcurrentHashMap<>(); @@ -448,10 +448,6 @@ public void addCancelable(CompletableFuture cancelable) { } } - protected void setMetadata(Map metadata) { - this.additionalObjects = new ConcurrentHashMap<>(metadata); - } - @Override public T addMetadataIfAbsent(String key, Supplier supplier) { return (T) additionalObjects.computeIfAbsent(key, k -> supplier.get()); diff --git a/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/PersistenceTaskInfo.java b/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/PersistenceTaskInfo.java index 93b6526d1..b6fdc2298 100644 --- a/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/PersistenceTaskInfo.java +++ b/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/PersistenceTaskInfo.java @@ -15,4 +15,9 @@ */ package io.serverlessworkflow.impl.persistence; -public interface PersistenceTaskInfo {} +import java.util.Map; + +public interface PersistenceTaskInfo { + + Map additionalObjects(); +} diff --git a/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/RetriedTaskInfo.java b/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/RetriedTaskInfo.java index 2b32c63ec..5894fdd9a 100644 --- a/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/RetriedTaskInfo.java +++ b/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/RetriedTaskInfo.java @@ -17,7 +17,7 @@ import java.util.Map; -public record RetriedTaskInfo(int retryAttempt, Map metadata) +public record RetriedTaskInfo(int retryAttempt, Map additionalObjects) implements PersistenceTaskInfo { public RetriedTaskInfo(int retryAttempt) { this(retryAttempt, Map.of()); diff --git a/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/WorkflowPersistenceInstance.java b/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/WorkflowPersistenceInstance.java index a64fede81..d582bba09 100644 --- a/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/WorkflowPersistenceInstance.java +++ b/impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/WorkflowPersistenceInstance.java @@ -24,12 +24,16 @@ import io.serverlessworkflow.impl.WorkflowMutableInstance; import io.serverlessworkflow.impl.WorkflowStatus; import io.serverlessworkflow.impl.executors.TransitionInfo; +import java.util.HashSet; import java.util.Optional; +import java.util.Set; import java.util.concurrent.CompletableFuture; +import java.util.function.Supplier; public class WorkflowPersistenceInstance extends WorkflowMutableInstance { private final PersistenceWorkflowInfo info; + private Set userMetaKeys = new HashSet<>(); public static WorkflowInstance of(WorkflowDefinition definition, PersistenceWorkflowInfo info) { return definition @@ -48,7 +52,7 @@ private WorkflowPersistenceInstance(WorkflowDefinition definition, PersistenceWo } }); this.startedAt = info.startedAt(); - setMetadata(info.metadata()); + additionalObjects.putAll(info.metadata()); } @Override @@ -84,7 +88,6 @@ public void restoreContext(WorkflowContext workflow, TaskContext context) { : workflow.definition().taskExecutor(completedTaskInfo.nextPosition()), completedTaskInfo.isEndNode())); workflow.context(completedTaskInfo.context()); - setMetadata(completedTaskInfo.additionalObjects()); } else if (taskInfo instanceof RetriedTaskInfo retriedTaskInfo) { if (context.retryAttempt() == 0) { context.retryAttempt(retriedTaskInfo.retryAttempt()); @@ -100,7 +103,44 @@ public void restoreContext(WorkflowContext workflow, TaskContext context) { } searchContext = tryContext.parent(); } - setMetadata(retriedTaskInfo.metadata()); + } + Set retainKeys; + synchronized (userMetaKeys) { + retainKeys = new HashSet<>(userMetaKeys); + taskInfo + .additionalObjects() + .forEach( + (k, v) -> { + if (!retainKeys.contains(k)) { + additionalObjects.put(k, v); + retainKeys.add(k); + } + }); + additionalObjects.keySet().retainAll(retainKeys); + } + } + + @Override + public T addMetadataIfAbsent(String key, Supplier supplier) { + synchronized (userMetaKeys) { + userMetaKeys.add(key); + return super.addMetadataIfAbsent(key, supplier); + } + } + + @Override + public void removeMetadata(String key) { + synchronized (userMetaKeys) { + userMetaKeys.add(key); + super.removeMetadata(key); + } + } + + @Override + public Optional removeMetadata(String key, Class clazz) { + synchronized (userMetaKeys) { + userMetaKeys.add(key); + return super.removeMetadata(key, clazz); } } }