Skip to content

Commit a269eee

Browse files
authored
[Fix #1691] MVStore additional objects persistence (#1715)
* [Fix #1691] Add metadata to mvstore Signed-off-by: Francisco Javier Tirado Sarti <ftirados@ibm.com> * [Fix #1691] MapInMemory approach Signed-off-by: Francisco Javier Tirado Sarti <ftirados@ibm.com> * [Fix #1691] Moving from int index to string generated index Signed-off-by: Francisco Javier Tirado Sarti <ftirados@ibm.com> * [Fix #1691] Use Ulid for improving versatility Signed-off-by: Francisco Javier Tirado Sarti <ftirados@ibm.com> --------- Signed-off-by: Francisco Javier Tirado Sarti <ftirados@ibm.com>
1 parent 8db8273 commit a269eee

26 files changed

Lines changed: 895 additions & 49 deletions

File tree

‎impl/core/src/main/java/io/serverlessworkflow/impl/WorkflowMutableInstance.java‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@ public class WorkflowMutableInstance implements WorkflowInstance {
5353
protected AtomicReference<CompletableFuture<WorkflowModel>> futureRef = new AtomicReference<>();
5454
protected Instant completedAt;
5555

56-
protected final Map<String, Object> additionalObjects = new ConcurrentHashMap<>();
56+
private Map<String, Object> additionalObjects = new ConcurrentHashMap<>();
5757

5858
protected final Map<String, Integer> iterationsMap = new ConcurrentHashMap<>();
5959

@@ -448,6 +448,10 @@ public void addCancelable(CompletableFuture<?> cancelable) {
448448
}
449449
}
450450

451+
protected void setMetadata(Map<String, Object> metadata) {
452+
this.additionalObjects = new ConcurrentHashMap<>(metadata);
453+
}
454+
451455
@Override
452456
public <T> T addMetadataIfAbsent(String key, Supplier<T> supplier) {
453457
return (T) additionalObjects.computeIfAbsent(key, k -> supplier.get());

‎impl/core/src/main/java/io/serverlessworkflow/impl/marshaller/DefaultOutputBuffer.java‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -122,6 +122,16 @@ public WorkflowOutputBuffer writeBytes(byte[] bytes) {
122122
return this;
123123
}
124124

125+
@Override
126+
public WorkflowOutputBuffer writeRawBytes(byte[] bytes) {
127+
try {
128+
output.write(bytes);
129+
} catch (IOException e) {
130+
throw new UncheckedIOException(e);
131+
}
132+
return this;
133+
}
134+
125135
@Override
126136
public void close() {
127137
try {

‎impl/core/src/main/java/io/serverlessworkflow/impl/marshaller/MarshallingUtils.java‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -173,6 +173,10 @@ public static WorkflowModel readModel(WorkflowBufferFactory factory, byte[] valu
173173
return readValue(factory, value, b -> (WorkflowModel) b.readObject());
174174
}
175175

176+
public static Object readObject(WorkflowBufferFactory factory, byte[] value) {
177+
return readValue(factory, value, b -> b.readObject());
178+
}
179+
176180
public static Instant readInstant(WorkflowBufferFactory factory, byte[] value) {
177181
return readValue(factory, value, WorkflowInputBuffer::readInstant);
178182
}

‎impl/core/src/main/java/io/serverlessworkflow/impl/marshaller/WorkflowOutputBuffer.java‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,8 @@ public interface WorkflowOutputBuffer extends AutoCloseable {
4141

4242
WorkflowOutputBuffer writeBytes(byte[] bytes);
4343

44+
WorkflowOutputBuffer writeRawBytes(byte[] bytes);
45+
4446
WorkflowOutputBuffer writeInstant(Instant instant);
4547

4648
default WorkflowOutputBuffer writeOffsetDateTime(OffsetDateTime time) {

‎impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/CompletedTaskInfo.java‎

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,14 +17,16 @@
1717

1818
import io.serverlessworkflow.impl.WorkflowModel;
1919
import java.time.Instant;
20+
import java.util.Map;
2021

2122
public record CompletedTaskInfo(
2223
Instant instant,
2324
WorkflowModel model,
2425
WorkflowModel context,
2526
Boolean isEndNode,
2627
String nextPosition,
27-
int iteration)
28+
int iteration,
29+
Map<String, Object> additionalObjects)
2830
implements PersistenceTaskInfo {
2931

3032
public CompletedTaskInfo(
@@ -35,4 +37,14 @@ public CompletedTaskInfo(
3537
String nextPosition) {
3638
this(instant, model, context, isEndNode, nextPosition, 1);
3739
}
40+
41+
public CompletedTaskInfo(
42+
Instant instant,
43+
WorkflowModel model,
44+
WorkflowModel context,
45+
Boolean isEndNode,
46+
String nextPosition,
47+
int iteration) {
48+
this(instant, model, context, isEndNode, nextPosition, iteration, Map.of());
49+
}
3850
}

‎impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/PersistenceInstanceInfo.java‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,5 +17,12 @@
1717

1818
import io.serverlessworkflow.impl.WorkflowModel;
1919
import java.time.Instant;
20+
import java.util.Map;
2021

21-
public record PersistenceInstanceInfo(Instant startedAt, WorkflowModel input) {}
22+
public record PersistenceInstanceInfo(
23+
Instant startedAt, WorkflowModel input, Map<String, Object> metadata) {
24+
25+
public PersistenceInstanceInfo(Instant startedAt, WorkflowModel input) {
26+
this(startedAt, input, Map.of());
27+
}
28+
}

‎impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/PersistenceWorkflowInfo.java‎

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,4 +25,15 @@ public record PersistenceWorkflowInfo(
2525
Instant startedAt,
2626
WorkflowModel input,
2727
WorkflowStatus status,
28-
Map<String, PersistenceTaskInfo> tasks) {}
28+
Map<String, PersistenceTaskInfo> tasks,
29+
Map<String, Object> metadata) {
30+
31+
public PersistenceWorkflowInfo(
32+
String id,
33+
Instant startedAt,
34+
WorkflowModel input,
35+
WorkflowStatus status,
36+
Map<String, PersistenceTaskInfo> tasks) {
37+
this(id, startedAt, input, status, tasks, Map.of());
38+
}
39+
}

‎impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/RetriedTaskInfo.java‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,4 +15,11 @@
1515
*/
1616
package io.serverlessworkflow.impl.persistence;
1717

18-
public record RetriedTaskInfo(int retryAttempt) implements PersistenceTaskInfo {}
18+
import java.util.Map;
19+
20+
public record RetriedTaskInfo(int retryAttempt, Map<String, Object> metadata)
21+
implements PersistenceTaskInfo {
22+
public RetriedTaskInfo(int retryAttempt) {
23+
this(retryAttempt, Map.of());
24+
}
25+
}

‎impl/persistence/api/src/main/java/io/serverlessworkflow/impl/persistence/WorkflowPersistenceInstance.java‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@ private WorkflowPersistenceInstance(WorkflowDefinition definition, PersistenceWo
4848
}
4949
});
5050
this.startedAt = info.startedAt();
51+
setMetadata(info.metadata());
5152
}
5253

5354
@Override
@@ -83,6 +84,7 @@ public void restoreContext(WorkflowContext workflow, TaskContext context) {
8384
: workflow.definition().taskExecutor(completedTaskInfo.nextPosition()),
8485
completedTaskInfo.isEndNode()));
8586
workflow.context(completedTaskInfo.context());
87+
setMetadata(completedTaskInfo.additionalObjects());
8688
} else if (taskInfo instanceof RetriedTaskInfo retriedTaskInfo) {
8789
if (context.retryAttempt() == 0) {
8890
context.retryAttempt(retriedTaskInfo.retryAttempt());
@@ -98,6 +100,7 @@ public void restoreContext(WorkflowContext workflow, TaskContext context) {
98100
}
99101
searchContext = tryContext.parent();
100102
}
103+
setMetadata(retriedTaskInfo.metadata());
101104
}
102105
}
103106
}
Lines changed: 71 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,71 @@
1+
/*
2+
* Copyright 2020-Present The Serverless Workflow Specification Authors
3+
*
4+
* Licensed under the Apache License, Version 2.0 (the "License");
5+
* you may not use this file except in compliance with the License.
6+
* You may obtain a copy of the License at
7+
*
8+
* http://www.apache.org/licenses/LICENSE-2.0
9+
*
10+
* Unless required by applicable law or agreed to in writing, software
11+
* distributed under the License is distributed on an "AS IS" BASIS,
12+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
* See the License for the specific language governing permissions and
14+
* limitations under the License.
15+
*/
16+
package io.serverlessworkflow.impl.persistence.hashing;
17+
18+
import com.github.f4b6a3.ulid.Ulid;
19+
import com.github.f4b6a3.ulid.UlidFactory;
20+
import io.serverlessworkflow.impl.marshaller.WorkflowInputBuffer;
21+
import java.util.Optional;
22+
23+
public class DefaultHashFactory implements HashFactory {
24+
25+
private final UlidFactory idFactory = UlidFactory.newMonotonicInstance();
26+
27+
@Override
28+
public Optional<HashItem> fromBuffer(byte id, WorkflowInputBuffer buffer) {
29+
return Optional.ofNullable(
30+
switch (id) {
31+
case MD5HashItem.ID -> new MD5HashItem(buffer);
32+
case IntegerHashItem.ID -> new IntegerHashItem(buffer);
33+
case HashItem.HASHING_DISABLED -> null;
34+
default -> throw new UnsupportedOperationException("Unsupported id " + id);
35+
});
36+
}
37+
38+
@Override
39+
public Optional<HashItem> fromData(byte[] data) {
40+
if (md5Condition(data)) {
41+
return Optional.of(new MD5HashItem(data));
42+
} else if (intCondition(data)) {
43+
return Optional.of(new IntegerHashItem(data));
44+
} else {
45+
return Optional.empty();
46+
}
47+
}
48+
49+
protected boolean intCondition(byte[] data) {
50+
return data.length > IntegerHashItem.SIZE_THRESHOLD;
51+
}
52+
53+
protected boolean md5Condition(byte[] data) {
54+
return data.length > MD5HashItem.SIZE_THRESHOLD;
55+
}
56+
57+
@Override
58+
public HashIndex indexFromBytes(byte[] bytes) {
59+
return new DefaultHashIndex(Ulid.from(bytes));
60+
}
61+
62+
@Override
63+
public HashIndex indexFromString(String str) {
64+
return new DefaultHashIndex(Ulid.from(str));
65+
}
66+
67+
@Override
68+
public HashIndex newIndex() {
69+
return new DefaultHashIndex(idFactory.create());
70+
}
71+
}

0 commit comments

Comments
 (0)