Make the storage node a resource that operations are linked to

The storage node no longer has targets and views, nor edges into and
out of it that other steps read and write through. It says where a
storage is - a catalog storage or the runner's own connection - what an
execution sees of it, and what is prepared at the start; it has no ports
and nothing connects to it.

Reading and writing is done by StorageOperation blocks, each a step of
its own, with its status, events and errors where they happen. Every
operation is linked to a Storage node, which it names by id; the editor
draws that as a link, so the field is hidden (a new @UiHidden,
x-ui-hidden). Its ports no longer depend on the storage: a read returns
what 'Read as' says, and the validation checks it, like the template,
against the linked node's type.

Gone with that: the hooks that wrote and read inside other steps, the
lazy and synthetic inputs, the checks on edge types and view parameters.
The ordering rule stays, between the operations that write to a node
and those that read from it.

context.iteration is now readable by every executor, through a view of
the execution's variables that writes through to the shared map.

The fuzz test wires random operations on an in-memory node into random
flows; every one the validation accepts still runs to an end.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Lucio Lelii 2026-09-28 17:34:39 +02:00
parent b94a21ac03
commit a06a00758e
28 changed files with 512 additions and 1650 deletions

View File

@ -50,6 +50,7 @@ import it.cnr.isti.workflow.manager.configurations.annotations.SchemaAllowedValu
import it.cnr.isti.workflow.manager.configurations.annotations.UiDescription;
import it.cnr.isti.workflow.manager.configurations.annotations.UiRequiredWhen;
import it.cnr.isti.workflow.manager.configurations.annotations.UiEnabledWhen;
import it.cnr.isti.workflow.manager.configurations.annotations.UiHidden;
import it.cnr.isti.workflow.manager.configurations.annotations.UiVisibleWhen;
import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel;
import it.cnr.isti.workflow.manager.configurations.annotations.UiOrder;
@ -103,6 +104,7 @@ public class JsonSchemaProducer {
Map<Class<?>, Map<String, LongText>> longTextMap = collectLongTextMetadata(type);
Map<Class<?>, Set<String>> acceptsPlaceholderMap = collectAcceptsVariablePlaceholderMetadata(type);
Map<Class<?>, Set<String>> defaultsWhenEmptyMap = collectDefaultsWhenEmptyMetadata(type);
Map<Class<?>, Set<String>> hiddenMap = collectFlagMetadata(type, UiHidden.class);
Map<Class<?>, Map<String, Structural>> structuralMap = collectStructuralMetadata(type);
Map<Class<?>, Map<String, UiOptionalGroup>> uiOptionalGroupMap = collectUiOptionalGroupMetadata(type);
Map<Class<?>, Map<String, UiEnabledWhen>> uiEnabledWhenMap = collectUiEnabledWhenMetadata(type);
@ -122,6 +124,7 @@ public class JsonSchemaProducer {
applyLongTextMetadata(root, getMergedMetadata(longTextMap, type));
applyAcceptsVariablePlaceholderMetadata(root, mergedNames(acceptsPlaceholderMap, type));
applyDefaultsWhenEmptyMetadata(root, mergedNames(defaultsWhenEmptyMap, type));
applyFlagMetadata(root, mergedNames(hiddenMap, type), "x-ui-hidden");
applyStructuralMetadata(root, getMergedMetadata(structuralMap, type));
applyUiOptionalGroupMetadata(root, getMergedMetadata(uiOptionalGroupMap, type));
applyUiEnabledWhenMetadata(root, getMergedMetadata(uiEnabledWhenMap, type));
@ -171,6 +174,7 @@ public class JsonSchemaProducer {
applyLongTextMetadata(classSchema, getMergedMetadata(longTextMap, matchedClass));
applyAcceptsVariablePlaceholderMetadata(classSchema, mergedNames(acceptsPlaceholderMap, matchedClass));
applyDefaultsWhenEmptyMetadata(classSchema, mergedNames(defaultsWhenEmptyMap, matchedClass));
applyFlagMetadata(classSchema, mergedNames(hiddenMap, matchedClass), "x-ui-hidden");
applyStructuralMetadata(classSchema, getMergedMetadata(structuralMap, matchedClass));
applyUiOptionalGroupMetadata(classSchema, getMergedMetadata(uiOptionalGroupMap, matchedClass));
applyUiEnabledWhenMetadata(classSchema, getMergedMetadata(uiEnabledWhenMap, matchedClass));
@ -880,6 +884,54 @@ public class JsonSchemaProducer {
return result;
}
/** The fields carrying a marker annotation, per class, the way the other collectors walk types. */
private Map<Class<?>, Set<String>> collectFlagMetadata(Class<?> rootClass,
Class<? extends java.lang.annotation.Annotation> annotation) {
Map<Class<?>, Set<String>> result = new HashMap<>();
Set<Class<?>> visited = new HashSet<>();
Queue<Class<?>> queue = new ArrayDeque<>();
queue.add(rootClass);
while (!queue.isEmpty()) {
Class<?> current = queue.poll();
if (current == null || !visited.add(current) || isTerminalType(current)) {
continue;
}
Set<String> names = new LinkedHashSet<>();
for (Field field : current.getDeclaredFields()) {
if (field.getAnnotation(annotation) != null) {
names.add(field.getName());
}
enqueueRelatedTypes(queue, field.getGenericType(), field.getType());
}
if (current.isRecord()) {
for (RecordComponent component : current.getRecordComponents()) {
if (component.getAnnotation(annotation) != null) {
names.add(component.getName());
}
enqueueRelatedTypes(queue, component.getGenericType(), component.getType());
}
}
if (current.getSuperclass() != null) {
queue.add(current.getSuperclass());
}
if (!names.isEmpty()) {
result.put(current, names);
}
}
return result;
}
private void applyFlagMetadata(ObjectNode classSchema, Set<String> names, String key) {
if (names == null || names.isEmpty() || !(classSchema.get("properties") instanceof ObjectNode properties)) {
return;
}
for (String name : names) {
if (properties.get(name) instanceof ObjectNode propertySchema) {
propertySchema.put(key, true);
}
}
}
private void applyDefaultsWhenEmptyMetadata(ObjectNode classSchema, Set<String> names) {
if (names == null || names.isEmpty()) {
return;

View File

@ -4,8 +4,6 @@
package it.cnr.isti.workflow.manager.blocks.configurations;
import java.util.List;
import java.util.Optional;
import com.fasterxml.jackson.annotation.JsonProperty;
@ -13,17 +11,13 @@ import it.cnr.isti.workflow.manager.blocks.types.StorageNodeBlockType;
import it.cnr.isti.workflow.manager.configurations.annotations.FieldRetriever;
import it.cnr.isti.workflow.manager.configurations.annotations.LongText;
import it.cnr.isti.workflow.manager.configurations.annotations.SchemaAllowedValues;
import it.cnr.isti.workflow.manager.configurations.annotations.Structural;
import it.cnr.isti.workflow.manager.configurations.annotations.UiDescription;
import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel;
import it.cnr.isti.workflow.manager.configurations.annotations.UiOrder;
import it.cnr.isti.workflow.manager.configurations.annotations.UiUniqueItemsBy;
import it.cnr.isti.workflow.manager.configurations.annotations.UiVisibleWhen;
import it.cnr.isti.workflow.manager.storage.StorageInit;
import it.cnr.isti.workflow.manager.storage.validation.ValidStorageNode;
import jakarta.validation.Valid;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.Size;
import lombok.Builder;
import lombok.EqualsAndHashCode;
import lombok.Getter;
@ -33,11 +27,12 @@ import lombok.NonNull;
/**
* A storage the flow keeps data in, as a node of its own rather than a step.
*
* <p>It is prepared once, when the execution starts: {@code createIfMissing} makes an object
* store's bucket, {@code initScript} runs a database's DDL. Its {@link #targets} are where outputs
* connected to it are written, each time the step producing them finishes; its {@link #views} are
* read afresh each time a step they feed is about to run. The order between a step that writes and
* one that reads is not the node's: it is set between those two steps with a dependency.
* <p>It has no ports and never runs. It says where a storage is and how it is reached - a catalog
* storage or the runner's own connection - what an execution sees of it, and what is prepared once,
* when the execution starts: {@code createIfMissing} makes an object store's bucket,
* {@code initScript} runs a database's DDL. What is read and written there is done by the storage
* operations linked to it, each a step of the flow; the order between one that writes and one that
* reads is set between those two, with a dependency.
*/
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
@Getter
@ -45,10 +40,12 @@ import lombok.NonNull;
@ValidStorageNode
public class StorageNodeConfiguration extends BlockConfiguration<StorageNodeBlockType> {
public static final int MAX_PORTS = 16;
public static final String SOURCE_CATALOG = "CATALOG";
public static final String SOURCE_PERSONAL = "PERSONAL";
public static final String SCOPE_SHARED = "SHARED";
public static final String SCOPE_PER_EXECUTION = "PER_EXECUTION";
@NotBlank
@Structural
@UiOrder(10)
@UiLabel("Storage type")
@FieldRetriever(name = "Storage", url = "/retriever/Storage/types")
@ -59,14 +56,14 @@ public class StorageNodeConfiguration extends BlockConfiguration<StorageNodeBloc
@UiLabel("Connection")
@UiDescription("A storage of the catalog, set up by whoever runs this service, or your own connection "
+ "saved in your vault and chosen when the execution starts.")
@SchemaAllowedValues(value = { StorageOperationBlockConfiguration.SOURCE_CATALOG,
StorageOperationBlockConfiguration.SOURCE_PERSONAL }, defaultValue = StorageOperationBlockConfiguration.SOURCE_CATALOG)
@SchemaAllowedValues(value = { SOURCE_CATALOG,
SOURCE_PERSONAL }, defaultValue = SOURCE_CATALOG)
@JsonProperty(required = false)
String source = StorageOperationBlockConfiguration.SOURCE_CATALOG;
String source = SOURCE_CATALOG;
@UiOrder(20)
@UiLabel("Storage")
@UiVisibleWhen(field = "source", equals = StorageOperationBlockConfiguration.SOURCE_CATALOG)
@UiVisibleWhen(field = "source", equals = SOURCE_CATALOG)
@FieldRetriever(name = "Storage", url = "/retriever/Storage/instances", dependsOn = { "storageType" })
@JsonProperty(required = false)
String instance;
@ -74,10 +71,10 @@ public class StorageNodeConfiguration extends BlockConfiguration<StorageNodeBloc
@UiOrder(30)
@UiLabel("Scope")
@UiDescription("Shared: the same data for every execution. Per execution: an object store prefix of this execution's own.")
@SchemaAllowedValues(value = { StorageOperationBlockConfiguration.SCOPE_SHARED,
StorageOperationBlockConfiguration.SCOPE_PER_EXECUTION }, defaultValue = StorageOperationBlockConfiguration.SCOPE_SHARED)
@SchemaAllowedValues(value = { SCOPE_SHARED,
SCOPE_PER_EXECUTION }, defaultValue = SCOPE_SHARED)
@JsonProperty(required = false)
String scope = StorageOperationBlockConfiguration.SCOPE_SHARED;
String scope = SCOPE_SHARED;
@UiOrder(40)
@UiLabel("Create if missing")
@ -93,38 +90,16 @@ public class StorageNodeConfiguration extends BlockConfiguration<StorageNodeBloc
@JsonProperty(required = false)
String initScript;
@Structural
@Valid
@Size(max = MAX_PORTS)
@UiUniqueItemsBy("name")
@UiOrder(60)
@UiLabel("Targets")
@UiDescription("Where outputs connected to this node are written.")
@JsonProperty(required = false)
List<StorageNodeTarget> targets = List.of();
@Structural
@Valid
@Size(max = MAX_PORTS)
@UiUniqueItemsBy("name")
@UiOrder(70)
@UiLabel("Views")
@UiDescription("What steps connected to this node read, afresh each time they run.")
@JsonProperty(required = false)
List<StorageNodeView> views = List.of();
@Builder
public StorageNodeConfiguration(@NonNull String name, String storageType, String source, String instance, String scope,
Boolean createIfMissing, String initScript, List<StorageNodeTarget> targets, List<StorageNodeView> views) {
Boolean createIfMissing, String initScript) {
super(name);
this.storageType = storageType;
this.source = source == null || source.isBlank() ? StorageOperationBlockConfiguration.SOURCE_CATALOG : source;
this.source = source == null || source.isBlank() ? SOURCE_CATALOG : source;
this.instance = instance;
this.scope = scope == null || scope.isBlank() ? StorageOperationBlockConfiguration.SCOPE_SHARED : scope;
this.scope = scope == null || scope.isBlank() ? SCOPE_SHARED : scope;
this.createIfMissing = createIfMissing != null && createIfMissing;
this.initScript = initScript;
this.targets = targets == null ? List.of() : List.copyOf(targets);
this.views = views == null ? List.of() : List.copyOf(views);
}
@Override
@ -135,36 +110,18 @@ public class StorageNodeConfiguration extends BlockConfiguration<StorageNodeBloc
public static StorageNodeConfiguration empty() {
StorageNodeConfiguration configuration = new StorageNodeConfiguration();
configuration.name = StorageNodeBlockType.TYPE;
configuration.targets = List.of();
configuration.views = List.of();
return configuration;
}
public boolean usesPersonalConnection() {
return StorageOperationBlockConfiguration.SOURCE_PERSONAL.equals(source);
return SOURCE_PERSONAL.equals(source);
}
public boolean perExecution() {
return StorageOperationBlockConfiguration.SCOPE_PER_EXECUTION.equals(scope);
return SCOPE_PER_EXECUTION.equals(scope);
}
public StorageInit init() {
return new StorageInit(Boolean.TRUE.equals(createIfMissing), initScript);
}
public List<StorageNodeTarget> targetList() {
return targets == null ? List.of() : targets;
}
public List<StorageNodeView> viewList() {
return views == null ? List.of() : views;
}
public Optional<StorageNodeTarget> target(String name) {
return targetList().stream().filter(target -> target != null && name.equals(target.name())).findFirst();
}
public Optional<StorageNodeView> view(String name) {
return viewList().stream().filter(view -> view != null && name.equals(view.name())).findFirst();
}
}

View File

@ -1,49 +0,0 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.blocks.configurations;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonProperty;
import it.cnr.isti.workflow.manager.configurations.annotations.LongText;
import it.cnr.isti.workflow.manager.configurations.annotations.UiDescription;
import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel;
import it.cnr.isti.workflow.manager.configurations.annotations.UiOrder;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.Pattern;
import jakarta.validation.constraints.Size;
/**
* A place in a storage node that outputs are written to. It becomes an input port of the node,
* which - unlike a step's - takes as many connections as there are steps writing there.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record StorageNodeTarget(
@NotBlank
@Size(max = 64)
@Pattern(regexp = "^[A-Za-z][A-Za-z0-9_-]*$",
message = "a target name must start with a letter and contain only letters, digits, '-' or '_'")
@JsonProperty(required = true)
@UiLabel("Name")
@UiDescription("Name of the input port outputs are connected to.")
@UiOrder(10)
String name,
@JsonProperty(required = true)
@UiLabel("Key or statement")
@UiDescription("Object store: the key, e.g. rounds/${{context.iteration}}.json. Database: an INSERT, UPDATE or "
+ "MERGE, e.g. INSERT INTO notes(body) VALUES (${{value}}). ${{value}} is what is written, "
+ "${{value.field}} a field of it.")
@LongText(placeholder = "rounds/${{context.iteration}}.json", acceptVariableAsPlaceholder = true)
@UiOrder(20)
String template,
@JsonProperty(required = false)
@UiLabel("Content type")
@UiDescription("For an object store; guessed from the value when empty.")
@UiOrder(30)
String contentType) {
}

View File

@ -1,52 +0,0 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.blocks.configurations;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonProperty;
import it.cnr.isti.workflow.manager.configurations.annotations.LongText;
import it.cnr.isti.workflow.manager.configurations.annotations.SchemaAllowedValues;
import it.cnr.isti.workflow.manager.configurations.annotations.UiDescription;
import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel;
import it.cnr.isti.workflow.manager.configurations.annotations.UiOrder;
import it.cnr.isti.workflow.manager.storage.ViewShape;
import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.Pattern;
import jakarta.validation.constraints.Size;
/**
* A way of reading a storage node. It becomes an output port, read afresh each time the step it
* feeds is about to run. A {@code ${{name}}} that is neither the execution's own nor a global
* becomes a parameter input of the node, {@code <view>.<name>}, fed by another step's output.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record StorageNodeView(
@NotBlank
@Size(max = 64)
@Pattern(regexp = "^[A-Za-z][A-Za-z0-9_-]*$",
message = "a view name must start with a letter and contain only letters, digits, '-' or '_'")
@JsonProperty(required = true)
@UiLabel("Name")
@UiDescription("Name of the output port it is read through.")
@UiOrder(10)
String name,
@JsonProperty(required = true)
@UiLabel("Key, pattern or query")
@UiDescription("Object store: a key or a glob, e.g. rounds/*.json. Database: a SELECT, e.g. SELECT body FROM "
+ "notes WHERE round < ${{round}}. Values are bound, never pasted in: write ${{name}} without quotes.")
@LongText(placeholder = "rounds/*.json", acceptVariableAsPlaceholder = true)
@UiOrder(20)
String query,
@JsonProperty(required = false)
@UiLabel("Read as")
@UiDescription("Files, texts or JSON for an object store; a database always returns its rows as JSON.")
@SchemaAllowedValues(value = { "FILES", "TEXTS", "JSON", "KEYS" }, defaultValue = "FILES")
@UiOrder(30)
ViewShape as) {
}

View File

@ -7,12 +7,11 @@ package it.cnr.isti.workflow.manager.blocks.configurations;
import com.fasterxml.jackson.annotation.JsonProperty;
import it.cnr.isti.workflow.manager.blocks.types.StorageOperationBlockType;
import it.cnr.isti.workflow.manager.configurations.annotations.FieldRetriever;
import it.cnr.isti.workflow.manager.configurations.annotations.LongText;
import it.cnr.isti.workflow.manager.configurations.annotations.SchemaAllowedValues;
import it.cnr.isti.workflow.manager.configurations.annotations.Structural;
import it.cnr.isti.workflow.manager.configurations.annotations.UiContextKeys;
import it.cnr.isti.workflow.manager.configurations.annotations.UiDescription;
import it.cnr.isti.workflow.manager.configurations.annotations.UiHidden;
import it.cnr.isti.workflow.manager.configurations.annotations.UiLabel;
import it.cnr.isti.workflow.manager.configurations.annotations.UiOrder;
import it.cnr.isti.workflow.manager.configurations.annotations.UiVisibleWhen;
@ -26,11 +25,14 @@ import lombok.NoArgsConstructor;
import lombok.NonNull;
/**
* One operation on one storage of the catalog.
* One operation on the storage of a Storage node of the flow.
*
* <p>What {@code template} holds depends on the storage, not on this block: an object key or glob
* for an object store, a SQL statement for a database. The storage type checks it - see
* {@link ValidStorageOperation} - so this block stays the same whichever kind of storage is added.
* <p>Always on a node: the node says where the storage is, how it is reached, what an execution
* sees of it and what is prepared at the start; the operation says what to do there. The editor
* draws which node as a link between the two, which is what {@link #storageNode} holds - so it is
* not a field of the form. What {@code template} holds depends on the node's storage: an object key
* or glob for an object store, a SQL statement for a database; the validation checks it against the
* linked node.
*/
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
@Getter
@ -38,50 +40,11 @@ import lombok.NonNull;
@ValidStorageOperation
public class StorageOperationBlockConfiguration extends BlockConfiguration<StorageOperationBlockType> {
public static final String SOURCE_CATALOG = "CATALOG";
public static final String SOURCE_PERSONAL = "PERSONAL";
public static final String SOURCE_NODE = "NODE";
public static final String SCOPE_SHARED = "SHARED";
public static final String SCOPE_PER_EXECUTION = "PER_EXECUTION";
@Structural
@UiOrder(10)
@UiLabel("Storage type")
@UiVisibleWhen(field = "source", equalsAny = { SOURCE_CATALOG, SOURCE_PERSONAL })
@FieldRetriever(name = "Storage", url = "/retriever/Storage/types")
@JsonProperty(required = false)
String storageType;
/**
* Where the connection comes from. Not a default to go back to: the two are alternatives, and
* a new block starts on the catalog.
*/
@UiOrder(15)
@UiLabel("Connection")
@UiDescription("A storage of the catalog, set up by whoever runs this service; your own connection, saved in "
+ "your vault and chosen when the execution starts; or a storage node of this flow, whose connection, "
+ "preparation and scope it shares.")
@SchemaAllowedValues(value = { SOURCE_CATALOG, SOURCE_PERSONAL, SOURCE_NODE }, defaultValue = SOURCE_CATALOG)
@JsonProperty(required = false)
String source = SOURCE_CATALOG;
@UiOrder(18)
@UiLabel("Storage node")
@UiDescription("A storage node of this flow.")
@UiVisibleWhen(field = "source", equals = SOURCE_NODE)
@FieldRetriever(name = "FlowNodes", url = "/secure-retriever/FlowNodes/storage/items",
dependsOn = { UiContextKeys.FLOW_ID }, requiresAuth = true)
/** The id of the Storage node it works on: set by linking the two in the editor. */
@UiHidden
@JsonProperty(required = false)
String storageNode;
@UiOrder(20)
@UiLabel("Storage")
@UiDescription("A storage of the catalog, set up by whoever runs this service.")
@UiVisibleWhen(field = "source", equals = SOURCE_CATALOG)
@FieldRetriever(name = "Storage", url = "/retriever/Storage/instances", dependsOn = { "storageType" })
@JsonProperty(required = false)
String instance;
@Structural
@UiOrder(30)
@UiLabel("Operation")
@ -92,7 +55,7 @@ public class StorageOperationBlockConfiguration extends BlockConfiguration<Stora
@Structural
@UiOrder(40)
@UiLabel("Key, pattern or statement")
@UiDescription("Object store: a key such as reports/${{title}}.json, or for READ and LIST a glob such as "
@UiDescription("Object store: a key such as reports/${{context.iteration}}.json, or for READ and LIST a glob such as "
+ "reports/*.json. Database: a SELECT for READ, an INSERT, UPDATE or MERGE for WRITE, a DELETE for DELETE, "
+ "a schema (or nothing) for LIST. Other ${{name}} placeholders become inputs; ${{value.field}} is a field of what is written.")
@LongText(placeholder = "reports/*.json or SELECT * FROM notes WHERE round = ${{round}}",
@ -105,7 +68,7 @@ public class StorageOperationBlockConfiguration extends BlockConfiguration<Stora
@UiOrder(50)
@UiLabel("Read as")
@UiVisibleWhen(field = "operation", equals = "READ")
@UiDescription("Files, texts or JSON for an object store; a database always returns its rows as JSON.")
@UiDescription("Files, texts, JSON or just the names for an object store; a database returns its rows as JSON.")
@SchemaAllowedValues(value = { "FILES", "TEXTS", "JSON", "KEYS" }, defaultValue = "FILES")
@JsonProperty(required = false)
ViewShape as;
@ -117,27 +80,15 @@ public class StorageOperationBlockConfiguration extends BlockConfiguration<Stora
@JsonProperty(required = false)
String contentType;
@UiOrder(70)
@UiLabel("Scope")
@UiVisibleWhen(field = "source", equalsAny = { SOURCE_CATALOG, SOURCE_PERSONAL })
@UiDescription("Shared: the same data for every execution. Per execution: an object store prefix of this execution's own.")
@SchemaAllowedValues(value = { SCOPE_SHARED, SCOPE_PER_EXECUTION }, defaultValue = SCOPE_SHARED)
@JsonProperty(required = false)
String scope = SCOPE_SHARED;
@Builder
public StorageOperationBlockConfiguration(@NonNull String name, String storageType, String source, String instance,
String storageNode, StorageOperation operation, String template, ViewShape as, String contentType, String scope) {
public StorageOperationBlockConfiguration(@NonNull String name, String storageNode, StorageOperation operation,
String template, ViewShape as, String contentType) {
super(name);
this.storageType = storageType;
this.source = source == null || source.isBlank() ? SOURCE_CATALOG : source;
this.instance = instance;
this.storageNode = storageNode;
this.operation = operation == null ? StorageOperation.READ : operation;
this.template = template;
this.as = as;
this.contentType = contentType;
this.scope = scope == null || scope.isBlank() ? SCOPE_SHARED : scope;
}
@Override
@ -155,17 +106,8 @@ public class StorageOperationBlockConfiguration extends BlockConfiguration<Stora
return operation == null ? StorageOperation.READ : operation;
}
/** On a storage node of the flow, sharing its connection, preparation and scope. */
public boolean usesStorageNode() {
return SOURCE_NODE.equals(source);
}
/** On the runner's own connection, from their vault, rather than a catalog storage. */
public boolean usesPersonalConnection() {
return SOURCE_PERSONAL.equals(source);
}
public boolean perExecution() {
return SCOPE_PER_EXECUTION.equals(scope);
/** What a read returns: as asked, files otherwise. The ports follow it, whatever the node is. */
public ViewShape effectiveShape() {
return as == null ? ViewShape.FILES : as;
}
}

View File

@ -4,74 +4,28 @@
package it.cnr.isti.workflow.manager.blocks.factories;
import java.util.ArrayList;
import java.util.List;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.IOCapability;
import it.cnr.isti.workflow.manager.blocks.IOCapabilityType;
import it.cnr.isti.workflow.manager.blocks.IOCapabilityTypes;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeTarget;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeView;
import it.cnr.isti.workflow.manager.blocks.types.StorageNodeBlockType;
import it.cnr.isti.workflow.manager.ios.IODescriptor;
import it.cnr.isti.workflow.manager.ios.IOType;
import it.cnr.isti.workflow.manager.storage.StorageFlows;
import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
import it.cnr.isti.workflow.manager.storage.StorageType;
import it.cnr.isti.workflow.manager.storage.StorageTypes;
import it.cnr.isti.workflow.manager.storage.ViewShape;
import it.cnr.isti.workflow.manager.storage.validation.StorageNodeValidator;
/**
* A storage node's ports: an input for each target, an output for each view, and an input for each
* parameter a view has, named {@code <view>.<parameter>}. Never refuses a draft - see the
* validation for what is wrong with one.
*/
/** A storage node has no ports: the storage operations linked to it do the reading and writing. */
@Component
public class StorageNodeFactory implements BlockFactory<StorageNodeBlockType, StorageNodeConfiguration> {
private static final List<IOCapability> PARAMETER_CAPABILITIES = List.of(
new IOCapability(IOCapabilityType.TEXT, false), new IOCapability(IOCapabilityType.JSON, false));
private static final List<IOCapability> DEFAULT_WRITABLE = List.of(new IOCapability(IOCapabilityType.FILE, false),
new IOCapability(IOCapabilityType.TEXT, false), new IOCapability(IOCapabilityType.JSON, false));
private final StorageNodeBlockType blockType;
private final StorageTypes types;
public StorageNodeFactory(StorageNodeBlockType blockType, StorageTypes types) {
public StorageNodeFactory(StorageNodeBlockType blockType) {
this.blockType = blockType;
this.types = types;
}
@Override
public Block<StorageNodeBlockType> create(StorageNodeConfiguration configuration) {
StorageType type = types.find(configuration.getStorageType()).orElse(null);
List<IODescriptor> inputs = new ArrayList<>();
List<IODescriptor> outputs = new ArrayList<>();
for (StorageNodeTarget target : configuration.targetList()) {
if (target != null && target.name() != null && !target.name().isBlank()) {
inputs.add(IODescriptor.input(target.name(), IOType.ANY, false, type == null ? DEFAULT_WRITABLE : type.writableKinds()));
}
}
for (StorageNodeView view : configuration.viewList()) {
if (view == null || view.name() == null || view.name().isBlank()) {
continue;
}
ViewShape shape = type == null ? (view.as() == null ? ViewShape.FILES : view.as()) : StorageNodeValidator.shape(view, type);
outputs.add(IODescriptor.output(view.name(), shape.portType(), shape.multiple(),
List.of(new IOCapability(IOCapabilityTypes.from(shape.portType()), shape.multiple()))));
for (String parameter : StoragePlaceholders.parameters(view.query())) {
inputs.add(IODescriptor.input(StorageFlows.parameterPort(view.name(), parameter), IOType.TEXT, false,
PARAMETER_CAPABILITIES));
}
}
return Block.<StorageNodeBlockType>builder()
.inputs(inputs)
.outputs(outputs)
.specificConfiguration(configuration)
.type(blockType)
.build();
@ -89,12 +43,11 @@ public class StorageNodeFactory implements BlockFactory<StorageNodeBlockType, St
@Override
public List<IOCapability> supportedInputCapabilities() {
return DEFAULT_WRITABLE;
return List.of();
}
@Override
public List<IOCapability> supportedOutputCapabilities() {
return List.of(new IOCapability(IOCapabilityType.FILE, true), new IOCapability(IOCapabilityType.TEXT, true),
new IOCapability(IOCapabilityType.JSON, false));
return List.of();
}
}

View File

@ -19,10 +19,7 @@ import it.cnr.isti.workflow.manager.ios.IODescriptor;
import it.cnr.isti.workflow.manager.ios.IOType;
import it.cnr.isti.workflow.manager.storage.StorageOperation;
import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
import it.cnr.isti.workflow.manager.storage.StorageType;
import it.cnr.isti.workflow.manager.storage.StorageTypes;
import it.cnr.isti.workflow.manager.storage.ViewShape;
import it.cnr.isti.workflow.manager.storage.validation.StorageOperationValidator;
/**
* The ports of a storage operation, which follow from what it does.
@ -30,7 +27,8 @@ import it.cnr.isti.workflow.manager.storage.validation.StorageOperationValidator
* <ul>
* <li>Every {@code ${{name}}} of the template that is not the execution's own ({@code context.*},
* {@code global.*}) nor the written value is an input.
* <li>A WRITE takes {@code content}, of the kinds the storage accepts.
* <li>A WRITE takes {@code content}: a file, a text or JSON - the linked node's storage says which it
* can keep, and the validation checks it.
* <li>{@code result} is what came back: the objects or rows read, the names listed, or - for a write
* or a delete - where it went and how much it touched.
* </ul>
@ -46,20 +44,17 @@ public class StorageOperationBlockFactory implements BlockFactory<StorageOperati
private static final List<IOCapability> PARAMETER_CAPABILITIES = List.of(
new IOCapability(IOCapabilityType.TEXT, false), new IOCapability(IOCapabilityType.JSON, false));
private static final List<IOCapability> DEFAULT_WRITABLE = List.of(new IOCapability(IOCapabilityType.FILE, false),
private static final List<IOCapability> WRITABLE = List.of(new IOCapability(IOCapabilityType.FILE, false),
new IOCapability(IOCapabilityType.TEXT, false), new IOCapability(IOCapabilityType.JSON, false));
private final StorageOperationBlockType blockType;
private final StorageTypes types;
public StorageOperationBlockFactory(StorageOperationBlockType blockType, StorageTypes types) {
public StorageOperationBlockFactory(StorageOperationBlockType blockType) {
this.blockType = blockType;
this.types = types;
}
@Override
public Block<StorageOperationBlockType> create(StorageOperationBlockConfiguration configuration) {
StorageType type = types.find(configuration.getStorageType()).orElse(null);
StorageOperation operation = configuration.effectiveOperation();
List<IODescriptor> inputs = new ArrayList<>();
for (String parameter : StoragePlaceholders.parameters(configuration.getTemplate())) {
@ -68,22 +63,21 @@ public class StorageOperationBlockFactory implements BlockFactory<StorageOperati
}
}
if (operation == StorageOperation.WRITE) {
inputs.add(IODescriptor.input(CONTENT_INPUT, IOType.ANY, false, type == null ? DEFAULT_WRITABLE : type.writableKinds()));
inputs.add(IODescriptor.input(CONTENT_INPUT, IOType.ANY, false, WRITABLE));
}
return Block.<StorageOperationBlockType>builder()
.inputs(inputs)
.output(result(operation, configuration, type))
.output(result(operation, configuration))
.specificConfiguration(configuration)
.type(blockType)
.build();
}
private static IODescriptor result(StorageOperation operation, StorageOperationBlockConfiguration configuration,
StorageType type) {
/** The same whatever node it is linked to, so linking it to another changes no port. */
private static IODescriptor result(StorageOperation operation, StorageOperationBlockConfiguration configuration) {
return switch (operation) {
case READ -> {
ViewShape shape = type == null ? (configuration.getAs() == null ? ViewShape.FILES : configuration.getAs())
: StorageOperationValidator.shape(configuration, type);
ViewShape shape = configuration.effectiveShape();
yield IODescriptor.output(RESULT_OUTPUT, shape.portType(), shape.multiple(),
List.of(new IOCapability(IOCapabilityTypes.from(shape.portType()), shape.multiple())));
}
@ -106,7 +100,7 @@ public class StorageOperationBlockFactory implements BlockFactory<StorageOperati
@Override
public List<IOCapability> supportedInputCapabilities() {
return DEFAULT_WRITABLE;
return WRITABLE;
}
@Override

View File

@ -0,0 +1,19 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.configurations.annotations;
import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
/**
* A field the editor does not show as a field, because something else on the canvas says it - a
* link drawn between two nodes, say. It is still part of the configuration, saved and validated.
*/
@Target({ ElementType.FIELD, ElementType.RECORD_COMPONENT })
@Retention(RetentionPolicy.RUNTIME)
public @interface UiHidden {
}

View File

@ -1,57 +0,0 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.configurations.retrievers;
import java.util.List;
import java.util.Map;
import org.springframework.http.HttpStatus;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ResponseStatusException;
import it.cnr.isti.workflow.manager.auth.repo.LoginEntity;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeConfiguration;
import it.cnr.isti.workflow.manager.executions.design.FlowSharedVariableCatalogService;
import it.cnr.isti.workflow.manager.storage.StorageFlows;
/**
* Nodes of the flow being edited that another node refers to: for now its storage nodes, which a
* storage operation can work on. The id is what is stored; the name is what is shown, so renaming
* the node does not break the reference.
*/
@Component
public class FlowNodesFieldRetriever implements SecureDynamicFieldRetriever {
public static final String CATEGORY = "FlowNodes";
private final FlowSharedVariableCatalogService blocks;
public FlowNodesFieldRetriever(FlowSharedVariableCatalogService blocks) {
this.blocks = blocks;
}
@Override
public String getCategory() {
return CATEGORY;
}
@Override
public List<RetrieverItem> retrieve(String parameter, Map<String, String> params, LoginEntity user) {
if (!"storage".equals(parameter)) {
throw new ResponseStatusException(HttpStatus.NOT_FOUND, "Unknown FlowNodes retriever parameter: " + parameter);
}
return blocks.blocksOf(params.get("flowId"), user.getUsername()).stream()
.filter(StorageFlows::isStorageNode)
.map(node -> {
StorageNodeConfiguration configuration = (StorageNodeConfiguration) node.getSpecificConfiguration();
String where = configuration.usesPersonalConnection() ? "your own connection" : configuration.getInstance();
return new RetrieverItem(
new RetrieverItemDescriptor(node.getName(), configuration.getStorageType() + " storage node on " + where,
Map.of("storageType", String.valueOf(configuration.getStorageType()))),
node.getId(), false, true, List.of());
})
.toList();
}
}

View File

@ -16,7 +16,6 @@ import it.cnr.isti.workflow.manager.blocks.configurations.ChatInteractionBlockCo
import it.cnr.isti.workflow.manager.blocks.configurations.ConditionalBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration;
import it.cnr.isti.workflow.manager.storage.StorageCredentials;
import it.cnr.isti.workflow.manager.blocks.configurations.HumanInteractiveBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration;
@ -131,13 +130,10 @@ final class AuthorizationRequirementResolver {
.addStepReference(block.getId(), block.getName());
}
/** A storage node or step working on the runner's own connection asks for it from their vault. */
/** A storage node on the runner's own connection asks for it from their vault. */
private static void collectStorageRequirement(Map<String, RequirementAccumulator> requirements, Block<?> block) {
String storageType = null;
if (block != null && block.getSpecificConfiguration() instanceof StorageOperationBlockConfiguration configuration
&& configuration.usesPersonalConnection()) {
storageType = configuration.getStorageType();
} else if (block != null && block.getSpecificConfiguration() instanceof StorageNodeConfiguration configuration
if (block != null && block.getSpecificConfiguration() instanceof StorageNodeConfiguration configuration
&& configuration.usesPersonalConnection()) {
storageType = configuration.getStorageType();
}

View File

@ -32,15 +32,10 @@ import it.cnr.isti.workflow.manager.flows.model.FlowNode;
import it.cnr.isti.workflow.manager.flows.loops.FlowLoops;
import it.cnr.isti.workflow.manager.flows.model.Connection;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.IOCapability;
import it.cnr.isti.workflow.manager.blocks.IOCapabilityType;
import it.cnr.isti.workflow.manager.ios.IODescriptor;
import it.cnr.isti.workflow.manager.ios.IOType;
import it.cnr.isti.workflow.manager.storage.StorageFlows;
import it.cnr.isti.workflow.manager.storage.StorageNodeRuntime;
import it.cnr.isti.workflow.manager.flows.model.Dependency;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.ios.IODescriptor;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import lombok.Builder;
import lombok.Getter;
@ -61,8 +56,6 @@ public class ExecutionObject {
@JsonIgnore
private transient StorageFlows.Compiled storage;
@JsonIgnore
private transient StorageNodeRuntime storageRuntime;
List<Dependency> stepDependencies = new ArrayList<>();
@ -155,7 +148,6 @@ public class ExecutionObject {
: List.copyOf(executionFlow.getConnections());
FlowLoops.Analysis loops = FlowLoops.analyze(executionFlow);
List<Step<?>> steps = getStepsFromFlow(executionFlow, loops);
wireStorage(steps);
this.executorService = createExecutorService(steps.size());
@ -184,11 +176,7 @@ public class ExecutionObject {
throw new IllegalStateException("No executor found for node " + node.getName() + ". Cannot create execution");
});
flow.getNodes().forEach(node -> steps.add(new Step<>(node)));
// The inputs a step waits on for its storage view parameters, before anything is wired to them.
this.storage.syntheticInputs().forEach((stepId, names) -> steps.stream()
.filter(step -> step.getId().equals(stepId)).findFirst()
.ifPresent(step -> names.forEach(name -> step.addSyntheticInput(
IODescriptor.input(name, IOType.ANY, false, List.of(new IOCapability(IOCapabilityType.ANY, false)))))));
List<Connection> connections = flow.getConnections() == null ? List.of() : flow.getConnections();
for (Connection connection: connections){
Step<?> sourceStep = steps.stream().filter(s -> s.getId().equals(connection.getSourceId())).findFirst().orElseThrow();
@ -215,37 +203,6 @@ public class ExecutionObject {
return steps;
}
/**
* Gives each step that reads from or writes to a storage node what it needs to: its view-fed
* inputs are there from the start, and its hooks do the reading and writing when it runs.
*/
private void wireStorage(List<Step<?>> steps) {
if (!this.storage.hasStorage()) {
return;
}
for (Step<?> step : steps) {
List<StorageFlows.Read> reads = this.storage.readsBy(step.getId());
List<StorageFlows.Write> writes = this.storage.writesFrom(step.getId());
if (reads.isEmpty() && writes.isEmpty()) {
continue;
}
reads.forEach(read -> step.provideInputLazily(read.inputName()));
step.setStorageHooks(new StorageStepHooksImpl(reads, writes, this.storage.storageNodes(), this::storageRuntime));
}
}
/** Given by the service that builds the execution: the storage beans are not the execution's. */
public void configureStorage(StorageNodeRuntime runtime) {
this.storageRuntime = runtime;
}
private StorageNodeRuntime storageRuntime() {
if (this.storageRuntime == null) {
throw new IllegalStateException("This execution reads or writes a storage, and has not been given one");
}
return this.storageRuntime;
}
/** One of the flow's storage nodes, by id, for a storage operation that works on it. */
public java.util.Optional<Block<?>> storageNode(String id) {
return java.util.Optional.ofNullable(id == null ? null : this.storage.storageNodes().get(id));

View File

@ -1613,7 +1613,6 @@ public class ExecutionsService {
}
private void attachPersistence(ExecutionObject executionObject) {
executionObject.configureStorage(storageNodeRuntime);
executionObject.setStateChangeListener(() -> {
persist(executionObject);
cleanupManagedResourcesIfFinal(executionObject);

View File

@ -1,132 +0,0 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.executions;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.function.Supplier;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeTarget;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeView;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.executions.steps.Step;
import it.cnr.isti.workflow.manager.executions.steps.StepStorageHooks;
import it.cnr.isti.workflow.manager.storage.StorageFlows;
import it.cnr.isti.workflow.manager.storage.StorageNodeRuntime;
import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
import it.cnr.isti.workflow.manager.storage.StorageReadResult;
import it.cnr.isti.workflow.manager.storage.StorageWriteResult;
/** One step's reads and writes, compiled from the flow's storage edges when the execution is built. */
final class StorageStepHooksImpl implements StepStorageHooks {
static final String READ_FAILED = "STORAGE_READ_FAILED";
static final String WRITE_FAILED = "STORAGE_WRITE_FAILED";
private final List<StorageFlows.Read> reads;
private final List<StorageFlows.Write> writes;
private final Map<String, Block<?>> storageNodes;
private final Supplier<StorageNodeRuntime> runtime;
StorageStepHooksImpl(List<StorageFlows.Read> reads, List<StorageFlows.Write> writes, Map<String, Block<?>> storageNodes,
Supplier<StorageNodeRuntime> runtime) {
this.reads = List.copyOf(reads);
this.writes = List.copyOf(writes);
this.storageNodes = storageNodes;
this.runtime = runtime;
}
@Override
public void read(Step<?> step) {
for (StorageFlows.Read read : reads) {
Block<?> node = storageNodes.get(read.storageId());
StorageNodeConfiguration configuration = (StorageNodeConfiguration) node.getSpecificConfiguration();
StorageNodeView view = configuration.view(read.viewName()).orElseThrow();
StorageReadResult result;
try {
result = runtime.get().read(node, view, name -> parameterValue(step, read, name), step.getAuthorizations(),
step.runtimeVariables());
input(step, read.inputName()).provideRead(result.value());
} catch (RuntimeException ex) {
throw new NodeExecutionException(READ_FAILED, "Could not read " + view.name() + " of "
+ StorageNodeRuntime.describe(node) + " for " + read.inputName() + ": " + ex.getMessage(), ex);
}
Map<String, Object> details = new LinkedHashMap<>();
details.put("storage", node.getId());
details.put("view", view.name());
details.put("count", result.count());
details.put("truncated", result.truncated());
log(step, "Read " + view.name() + " of " + StorageNodeRuntime.describe(node) + ": " + result.count()
+ (result.truncated() ? "+ (truncated)" : ""), details);
}
}
@Override
public void write(Step<?> step, Map<String, Object> outputs) {
for (StorageFlows.Write write : writes) {
if (!outputs.containsKey(write.outputName())) {
// A branch not taken, or an output this round did not produce: nothing to write.
continue;
}
Block<?> node = storageNodes.get(write.storageId());
StorageNodeConfiguration configuration = (StorageNodeConfiguration) node.getSpecificConfiguration();
StorageNodeTarget target = configuration.target(write.targetName()).orElseThrow();
if (step.getBiasExecutionContext() != null && step.getBiasExecutionContext().isExperiment()) {
log(step, "Not written to " + StorageNodeRuntime.describe(node) + " in a bias experiment", Map.of(
"storage", node.getId(), "target", target.name()), ExecutionEventType.BIAS_SIDE_EFFECT_MOCKED);
continue;
}
StorageWriteResult result;
try {
result = runtime.get().write(node, target, outputs.get(write.outputName()), step.getAuthorizations(),
step.runtimeVariables());
} catch (RuntimeException ex) {
throw new NodeExecutionException(WRITE_FAILED, "Could not write " + write.outputName() + " to "
+ StorageNodeRuntime.describe(node) + "." + target.name() + ": " + ex.getMessage(), ex);
}
Map<String, Object> details = new LinkedHashMap<>();
details.put("storage", node.getId());
details.put("target", target.name());
details.put("affected", result.affected());
if (result.location() != null) {
details.put("location", result.location());
}
log(step, "Wrote " + write.outputName() + " to " + StorageNodeRuntime.describe(node) + "." + target.name()
+ (result.location() == null ? " (" + result.affected() + " rows)" : ": " + result.location()), details);
}
}
/** A view parameter from the input this step got for it; the execution's own names otherwise. */
private static Object parameterValue(Step<?> step, StorageFlows.Read read, String name) {
String synthetic = StorageFlows.syntheticInput(read.storageId(), read.viewName(), name);
for (Input input : step.getInputs()) {
if (input.getDescriptor().getName().equals(synthetic)) {
return input.getValue();
}
}
if (StoragePlaceholders.isExecutionName(name)) {
return step.runtimeVariables().get(name);
}
return null;
}
private static Input input(Step<?> step, String name) {
return step.getInputs().stream().filter(input -> input.getDescriptor().getName().equals(name)).findFirst()
.orElseThrow(() -> new IllegalStateException("Step " + step.getNode().getName() + " has no input " + name));
}
private static void log(Step<?> step, String message, Map<String, Object> details) {
log(step, message, details, ExecutionEventType.STORAGE_OPERATION);
}
private static void log(Step<?> step, String message, Map<String, Object> details, ExecutionEventType type) {
if (step.getEventLogger() != null) {
step.getEventLogger().info(type, message, details);
}
}
}

View File

@ -22,19 +22,14 @@ import it.cnr.isti.workflow.manager.executions.ExecutionRuntimeContextSupport;
import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor;
import it.cnr.isti.workflow.manager.executions.ExecutionsService;
import it.cnr.isti.workflow.manager.executions.steps.Input;
import it.cnr.isti.workflow.manager.storage.StorageConnection;
import it.cnr.isti.workflow.manager.storage.StorageException;
import it.cnr.isti.workflow.manager.storage.StorageConnectionResolver;
import it.cnr.isti.workflow.manager.storage.StorageNodeRuntime;
import it.cnr.isti.workflow.manager.storage.StorageOperation;
import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
import it.cnr.isti.workflow.manager.storage.StorageReadResult;
import it.cnr.isti.workflow.manager.storage.StorageScope;
import it.cnr.isti.workflow.manager.storage.StorageSession;
import it.cnr.isti.workflow.manager.storage.StorageType;
import it.cnr.isti.workflow.manager.storage.StorageTypes;
import it.cnr.isti.workflow.manager.storage.StorageWriteResult;
import it.cnr.isti.workflow.manager.storage.validation.StorageOperationValidator;
import tools.jackson.databind.ObjectMapper;
import tools.jackson.databind.node.ObjectNode;
@ -47,16 +42,12 @@ import tools.jackson.databind.node.ObjectNode;
@Component
public class StorageOperationExecutor implements BlockExecutor<StorageOperationBlockType> {
private final StorageTypes types;
private final StorageConnectionResolver connections;
private final StorageNodeRuntime nodes;
private final ObjectProvider<ExecutionsService> executions;
private final ObjectMapper objectMapper;
public StorageOperationExecutor(StorageTypes types, StorageConnectionResolver connections, StorageNodeRuntime nodes,
ObjectProvider<ExecutionsService> executions, ObjectMapper objectMapper) {
this.types = types;
this.connections = connections;
public StorageOperationExecutor(StorageNodeRuntime nodes, ObjectProvider<ExecutionsService> executions,
ObjectMapper objectMapper) {
this.nodes = nodes;
this.executions = executions;
this.objectMapper = objectMapper;
@ -67,27 +58,21 @@ public class StorageOperationExecutor implements BlockExecutor<StorageOperationB
Map<String, Object> authorizations, Map<String, Object> executionVariables,
Map<String, ExecutionVariableDescriptor> executionVariableDescriptors, ExecutionEventLogger eventLogger) {
StorageOperationBlockConfiguration configuration = (StorageOperationBlockConfiguration) block.getSpecificConfiguration();
Block<?> node = configuration.usesStorageNode() ? storageNode(configuration, executionVariables) : null;
StorageType type = node != null ? nodes.typeOf(node) : types.require(configuration.getStorageType());
StorageConnection connection = node != null ? null
: configuration.usesPersonalConnection()
? connections.personal(type, authorizations, executionVariables)
: connections.catalog(type, configuration.getInstance());
String storageName = node != null ? StorageNodeRuntime.describe(node) : connection.name();
Block<?> node = storageNode(configuration, executionVariables);
StorageType type = nodes.typeOf(node);
String storageName = StorageNodeRuntime.describe(node);
StorageOperation operation = configuration.effectiveOperation();
Function<String, Object> values = values(inputs, executionVariables);
long started = System.nanoTime();
Map<String, Object> details = new LinkedHashMap<>();
details.put("storage", node != null ? node.getId() : configuration.usesPersonalConnection() ? "personal" : configuration.getInstance());
details.put("storage", node.getId());
details.put("type", type.getName());
details.put("operation", operation.name());
Object result;
try (StorageSession session = node != null ? nodes.open(node, authorizations, executionVariables)
: type.open(connection, scope(configuration, executionVariables))) {
try (StorageSession session = nodes.open(node, authorizations, executionVariables)) {
result = switch (operation) {
case READ -> {
StorageReadResult read = session.read(configuration.getTemplate(),
StorageOperationValidator.shape(configuration, type), values);
StorageReadResult read = session.read(configuration.getTemplate(), configuration.effectiveShape(), values);
details.put("count", read.count());
details.put("truncated", read.truncated());
yield read.value();
@ -119,6 +104,9 @@ public class StorageOperationExecutor implements BlockExecutor<StorageOperationB
*/
private Block<?> storageNode(StorageOperationBlockConfiguration configuration, Map<String, Object> executionVariables) {
Object root = executionVariables == null ? null : executionVariables.get(ExecutionRuntimeContextSupport.ROOT_EXECUTION_ID);
if (configuration.getStorageNode() == null || configuration.getStorageNode().isBlank()) {
throw new StorageException(StorageException.CONNECTION_INVALID, "This storage operation is not linked to a Storage node");
}
if (root == null) {
throw new StorageException(StorageException.CONNECTION_INVALID, "This step has no execution to find its storage node in");
}
@ -148,18 +136,6 @@ public class StorageOperationExecutor implements BlockExecutor<StorageOperationB
return details.containsKey("count") ? ": " + details.get("count") + (Boolean.TRUE.equals(details.get("truncated")) ? "+ (truncated)" : "") : "";
}
private static StorageScope scope(StorageOperationBlockConfiguration configuration, Map<String, Object> executionVariables) {
if (!configuration.perExecution()) {
return StorageScope.shared();
}
Object root = executionVariables == null ? null : executionVariables.get(ExecutionRuntimeContextSupport.ROOT_EXECUTION_ID);
if (root == null || String.valueOf(root).isBlank()) {
throw new StorageException(StorageException.CONNECTION_INVALID,
"A per-execution scope needs the execution's id, and this step has none");
}
return StorageScope.perExecution(String.valueOf(root));
}
/** An input by its port name; the execution's own names from its variables. */
private static Function<String, Object> values(List<Input> inputs, Map<String, Object> executionVariables) {
Map<String, Object> byName = new LinkedHashMap<>();

View File

@ -76,26 +76,6 @@ public class Input {
this.registered = true;
}
/**
* Fed by a storage view, which is read when the step starts rather than delivered by a
* connection: there from the start, nothing waits for it, and nobody is asked for it.
*/
void provideLazily() {
this.registered = true;
this.resolutionState = InputResolutionState.VALUE;
}
/**
* What a storage view held just as the step was about to run. The step is already on its way,
* so nothing is told: this is the value it runs with, not news that may make it ready.
*/
public void provideRead(Object value) {
Object normalizedValue = normalizeJsonInput(value);
validateValue(normalizedValue);
this.value = normalizedValue;
this.resolutionState = InputResolutionState.VALUE;
}
/**
* Forgets this input's value so the next round of a loop can deliver a new one. Only a loop
* going round does this, and only for inputs fed from inside the loop: an input's state

View File

@ -170,11 +170,6 @@ public class Step<N extends FlowNode> implements InputListener {
@JsonIgnore
private final Set<String> skippedDependencyIds = new LinkedHashSet<>();
/** Reads from and writes to the flow's storage nodes, when this step has any to do. */
@JsonIgnore
@Setter
private StepStorageHooks storageHooks;
@Builder
public Step(@NonNull N node) {
// Initialize the step with the provided block
@ -233,13 +228,6 @@ public class Step<N extends FlowNode> implements InputListener {
logger.info("Executing step " + this.id + " of node " + this.node.getName());
if (!isSimulated() && this.node.isUserInteractive()) {
try {
// Read now, so the person answering sees what the storage holds.
readFromStorage();
} catch (Throwable e) {
failWith(e);
return;
}
this.status = StepStatus.WAITING_FOR_INTERACTION;
logger.info("Step " + this.id + " is waiting for user interaction.");
if (this.node instanceof Block<?> block
@ -253,12 +241,11 @@ public class Step<N extends FlowNode> implements InputListener {
this.status = StepStatus.RUNNING;
listener.started(this.id);
try {
readFromStorage();
NodeExecutionResult executionResult = isSimulated() && this.node.isUserInteractive()
? NodeExecutors.simulateResult(this.node, executorInputs(), authorizations, executionVariables,
? NodeExecutors.simulateResult(this.node, this.inputs, authorizations, executorVariables(),
executionVariableDescriptors, this.interactionSimulationDescriptor, this.eventLogger,
this.biasExecutionContext)
: NodeExecutors.executeResult(this.node, executorInputs(), authorizations, executionVariables,
: NodeExecutors.executeResult(this.node, this.inputs, authorizations, executorVariables(),
executionVariableDescriptors, this.eventLogger, this.biasExecutionContext,
new ContainerExecutionContext(this.parentExecutionId, this.id,
this.executionSimulationEnabled, this.interactionSimulationDescriptor));
@ -268,8 +255,6 @@ public class Step<N extends FlowNode> implements InputListener {
executionResult.onSuspended().run();
return;
}
// Written before the step counts as done, so whoever depends on it finds it there.
writeToStorage(executionResult.outputs());
listener.commit(() -> settle(executionResult));
} catch (Throwable e) {
failWith(e);
@ -290,54 +275,13 @@ public class Step<N extends FlowNode> implements InputListener {
});
}
private void readFromStorage() {
if (this.storageHooks != null) {
this.storageHooks.read(this);
}
}
private void writeToStorage(Map<String, Object> outputs) {
if (this.storageHooks != null && outputs != null) {
this.storageHooks.write(this, outputs);
}
}
/**
* The inputs the node itself declares. A step may also have inputs of its own for the storage
* view parameters it waits on; they are the engine's business, not the node's.
* The execution's variables as this step's executor sees them: the shared map itself for
* everything written - executors register sessions in it - with this step's round readable as
* {@code context.iteration}, so a template can tell one round of a loop from the next.
*/
private List<Input> executorInputs() {
if (this.inputs.stream().noneMatch(input -> StepStorageHooks.isSynthetic(input.getDescriptor().getName()))) {
return this.inputs;
}
return this.inputs.stream().filter(input -> !StepStorageHooks.isSynthetic(input.getDescriptor().getName())).toList();
}
/**
* The execution's variables, and this step's round: {@code context.iteration}, so a storage
* template can tell one round from the next. A copy: executors are handed the shared map
* itself, which some of them write to, and this one is only for reading.
*/
public Map<String, Object> runtimeVariables() {
Map<String, Object> variables = new HashMap<>(this.executionVariables == null ? Map.of() : this.executionVariables);
variables.put(it.cnr.isti.workflow.manager.executions.ExecutionRuntimeContextSupport.ITERATION, this.iteration);
return variables;
}
/** An input the node does not declare, for a storage view parameter this step waits on. */
public synchronized void addSyntheticInput(IODescriptor descriptor) {
Input input = new Input(descriptor);
input.setListener(this);
this.inputs.add(input);
refreshStateFromInputs();
}
/** An input a storage view fills when the step starts: see {@link Input#provideLazily()}. */
public synchronized void provideInputLazily(String inputName) {
this.inputs.stream().filter(input -> input.getDescriptor().getName().equals(inputName)).findFirst()
.orElseThrow(() -> new IllegalArgumentException("Step " + this.node.getName() + " has no input " + inputName))
.provideLazily();
refreshStateFromInputs();
private Map<String, Object> executorVariables() {
return new StepVariables(this.executionVariables, this.iteration);
}
private void settle(NodeExecutionResult executionResult) {
@ -378,7 +322,7 @@ public class Step<N extends FlowNode> implements InputListener {
listener.resumed(this.id);
InteractionResult interactionResult;
try {
interactionResult = NodeExecutors.interact(this.node, executorInputs(), providedOutputs,
interactionResult = NodeExecutors.interact(this.node, this.inputs, providedOutputs,
Map.copyOf(this.partialResults), authorizations, executionVariables, executionVariableDescriptors,
this.eventLogger, this.biasExecutionContext);
} catch (RuntimeException exception) {
@ -402,14 +346,6 @@ public class Step<N extends FlowNode> implements InputListener {
throw new IllegalArgumentException("Missing output value for: " + name);
}
}
try {
writeToStorage(resolvedOutputs);
} catch (RuntimeException exception) {
// Nothing was taken: the answer can be sent again once the storage is back.
this.status = StepStatus.WAITING_FOR_INTERACTION;
listener.paused(this.id);
throw exception;
}
listener.commit(() -> {
boolean goingRound = goesRound(resolvedOutputs);
for (Output output : this.outputs) {
@ -592,7 +528,7 @@ public class Step<N extends FlowNode> implements InputListener {
return;
}
try {
NodeExecutors.cancel(this.node, executorInputs(), Map.copyOf(this.partialResults), this.authorizations,
NodeExecutors.cancel(this.node, this.inputs, Map.copyOf(this.partialResults), this.authorizations,
this.executionVariables, this.executionVariableDescriptors, this.eventLogger);
} catch (RuntimeException ex) {
logger.warn("Failed to cancel resources for step {}", this.id, ex);

View File

@ -1,26 +0,0 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.executions.steps;
import java.util.Map;
/**
* What a step does with the flow's storage nodes: read the views that feed its inputs just before it
* runs, and write its outputs to the targets they are connected to just before it counts as done.
* Both run on the step's own thread, outside the execution's lock, and a failure in either fails the
* step with the reason.
*/
public interface StepStorageHooks {
String SYNTHETIC_PREFIX = "@storage/";
void read(Step<?> step);
void write(Step<?> step, Map<String, Object> outputs);
static boolean isSynthetic(String inputName) {
return inputName != null && inputName.startsWith(SYNTHETIC_PREFIX);
}
}

View File

@ -0,0 +1,57 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.executions.steps;
import java.util.AbstractMap;
import java.util.LinkedHashSet;
import java.util.Map;
import java.util.Set;
import it.cnr.isti.workflow.manager.executions.ExecutionRuntimeContextSupport;
/**
* The execution's variables plus one of the step's own, {@code context.iteration}.
*
* <p>Not a copy: executors write to the variables they are handed - a shared MCP session is
* registered there - and those writes have to reach the execution. Every write goes through to the
* shared map; only the round is this step's, readable and never stored.
*/
final class StepVariables extends AbstractMap<String, Object> {
private final Map<String, Object> shared;
private final int iteration;
StepVariables(Map<String, Object> shared, int iteration) {
this.shared = shared == null ? new java.util.HashMap<>() : shared;
this.iteration = iteration;
}
@Override
public Object get(Object key) {
return ExecutionRuntimeContextSupport.ITERATION.equals(key) ? iteration : shared.get(key);
}
@Override
public boolean containsKey(Object key) {
return ExecutionRuntimeContextSupport.ITERATION.equals(key) || shared.containsKey(key);
}
@Override
public Object put(String key, Object value) {
return shared.put(key, value);
}
@Override
public Object remove(Object key) {
return shared.remove(key);
}
@Override
public Set<Entry<String, Object>> entrySet() {
Set<Entry<String, Object>> entries = new LinkedHashSet<>(shared.entrySet());
entries.add(new SimpleImmutableEntry<>(ExecutionRuntimeContextSupport.ITERATION, iteration));
return entries;
}
}

View File

@ -45,11 +45,11 @@ public record NodeTypeCapabilities(
}
/**
* A node that is not a step: it never runs, and the edges into and out of it are reads and
* writes the steps on their other end do. It takes no dependencies either way - the order
* between a step that writes and one that reads is set between those two steps.
* A node that is not a step: it never runs and nothing connects to it. What uses it - a storage
* operation using a storage node - names it, and the editor draws that as a link of its own.
* It takes no dependencies either way: the order is between the steps that use it.
*/
public static NodeTypeCapabilities resource() {
return new NodeTypeCapabilities(NodeVisualRole.RESOURCE, false, false, true, true, false, false, false);
return new NodeTypeCapabilities(NodeVisualRole.RESOURCE, false, false, false, false, false, false, false);
}
}

View File

@ -438,12 +438,11 @@ public class FlowExecutionValidator {
if (flowData == null || flowData.getConnections() == null) {
return List.of();
}
StorageFlows.Compiled storage = StorageFlows.compile(flowData);
FlowLoops.Analysis loops = FlowLoops.analyze(storage.executionFlow());
FlowLoops.Analysis loops = FlowLoops.analyze(StorageFlows.executionView(flowData));
Map<String, List<Connection>> byInput = new java.util.LinkedHashMap<>();
for (Connection connection : flowData.getConnections()) {
if (connection == null || connection.getTargetId() == null || connection.getTargetName() == null
|| loops.isBackEdge(connection) || isStorageTarget(storage, connection)) {
|| loops.isBackEdge(connection)) {
continue;
}
byInput.computeIfAbsent(connection.getTargetId() + "\u0000" + connection.getTargetName(),
@ -463,71 +462,45 @@ public class FlowExecutionValidator {
return errors;
}
/** A storage node's target takes as many writers as the flow has; its view parameters do not. */
private static boolean isStorageTarget(StorageFlows.Compiled storage, Connection connection) {
return storage.writes().stream().anyMatch(write -> write.storageId().equals(connection.getTargetId())
&& write.targetName().equals(connection.getTargetName()));
}
/**
* What storage nodes require of the flow around them.
* What storage nodes and the operations linked to them require of the flow.
*
* <ul>
* <li>They live in the top-level flow.
* <li>Their edges go between a step and the node: not from one storage node into another, and not
* from a container, whose outputs arrive where nothing may wait on a storage.
* <li>A step that writes to one and a step that reads from it are ordered, by connections or
* dependencies, one way or the other. Otherwise which comes first is chance: the read may or may
* not see the write. Reading before writing is a choice too - in a loop it reads the round
* before - so either direction satisfies it; only no direction at all is refused.
* <li>Storage nodes live in the top-level flow, and so do the operations linked to one.
* <li>An operation is linked to a node that is there, and what it does suits that node's storage
* and what the storage allows.
* <li>An operation that writes to a node and one that reads from it are ordered, by connections
* or dependencies, one way or the other. Otherwise which comes first is chance: the read may or
* may not see the write. Reading first is a choice too - in a loop it reads the round before -
* so either direction satisfies it; only no direction at all is refused.
* </ul>
*/
private List<ValidationError> validateStorage(FlowData flowData, boolean rootFlow) {
StorageFlows.Compiled storage = StorageFlows.compile(flowData);
List<Block<?>> operations = (flowData.getBlocks() == null ? List.<Block<?>>of() : flowData.getBlocks()).stream()
.filter(block -> block != null && block.getSpecificConfiguration() instanceof StorageOperationBlockConfiguration configuration
&& configuration.usesStorageNode())
.filter(block -> block != null && block.getSpecificConfiguration() instanceof StorageOperationBlockConfiguration)
.toList();
if (!storage.hasStorage() && operations.isEmpty()) {
return List.of();
}
List<ValidationError> errors = new ArrayList<>();
if (!rootFlow) {
operations.forEach(block -> errors.add(new ValidationError(ValidationErrorCode.STORAGE_NODE_NOT_TOP_LEVEL, "block",
block.getId(), "specificConfiguration.storageNode",
"A storage operation inside a container cannot use a storage node yet; use a catalog storage or your own connection")));
} else {
errors.addAll(storageOperationReferences(storage, operations));
}
if (!rootFlow) {
storage.storageNodes().values().forEach(node -> errors.add(new ValidationError(
ValidationErrorCode.STORAGE_NODE_NOT_TOP_LEVEL, "block", node.getId(), "type",
"A storage node belongs in the top-level flow; steps inside a container can still reach it from there")));
"A storage node belongs in the top-level flow")));
operations.forEach(block -> errors.add(new ValidationError(ValidationErrorCode.STORAGE_NODE_NOT_TOP_LEVEL, "block",
block.getId(), "specificConfiguration.storageNode",
"A storage operation inside a container cannot be linked to a storage node yet; put it in the top-level flow")));
return errors;
}
for (Connection stray : storage.strayConnections()) {
errors.add(new ValidationError(ValidationErrorCode.STORAGE_CONNECTION_INVALID, "connection", stray.getId(),
"targetName", "A storage node connects to steps: a target takes a step's output, a view feeds a step's input",
List.of(stray.getSourceId(), stray.getTargetId())));
}
errors.addAll(storageOperationReferences(storage, operations));
Map<String, FlowNode> nodes = new HashMap<>();
flowData.getNodes().forEach(node -> nodes.put(node.getId(), node));
for (StorageFlows.Write write : storage.writes()) {
if (nodes.get(write.sourceId()) instanceof Container<?> container) {
errors.add(new ValidationError(ValidationErrorCode.STORAGE_CONNECTION_INVALID, "container", container.getId(),
"outputs." + write.outputName(), "A container's output cannot be written to a storage node yet; "
+ "write it from a step inside the container, or from one after it",
List.of(container.getId(), write.storageId())));
}
}
errors.addAll(storageEdgeTypes(storage, nodes, flowData));
errors.addAll(unfedViewParameters(storage));
Map<String, List<String>> graph = buildOutgoingGraph(flowData);
for (Block<?> node : storage.storageNodes().values()) {
Set<String> writers = new LinkedHashSet<>();
storage.writes().stream().filter(write -> write.storageId().equals(node.getId())).forEach(write -> writers.add(write.sourceId()));
Set<String> readers = new LinkedHashSet<>();
storage.reads().stream().filter(read -> read.storageId().equals(node.getId())).forEach(read -> readers.add(read.consumerId()));
for (Block<?> operation : operations) {
StorageOperationBlockConfiguration configuration = (StorageOperationBlockConfiguration) operation.getSpecificConfiguration();
if (node.getId().equals(configuration.getStorageNode())) {
@ -537,8 +510,7 @@ public class FlowExecutionValidator {
}
for (String reader : readers) {
for (String writer : writers) {
if (reader.equals(writer) || isReachable(writer, reader, graph, new HashSet<>())
|| isReachable(reader, writer, graph, new HashSet<>())) {
if (isReachable(writer, reader, graph, new HashSet<>()) || isReachable(reader, writer, graph, new HashSet<>())) {
continue;
}
String readerName = name(nodes.get(reader), reader);
@ -555,49 +527,20 @@ public class FlowExecutionValidator {
}
/**
* What travels along a storage edge has to fit where it lands. For a step's own connections the
* editor sees to this; a storage edge is only used when a step runs, so a mismatch here would
* otherwise surface as that step failing - a list of texts into an input that takes one, a file
* into a database's target.
*/
private List<ValidationError> storageEdgeTypes(StorageFlows.Compiled storage, Map<String, FlowNode> nodes, FlowData flowData) {
List<ValidationError> errors = new ArrayList<>();
for (Connection connection : flowData.getConnections() == null ? List.<Connection>of() : flowData.getConnections()) {
if (connection == null) {
continue;
}
FlowNode source = nodes.get(connection.getSourceId());
FlowNode target = nodes.get(connection.getTargetId());
boolean fromStorage = storage.storageNodes().containsKey(connection.getSourceId());
boolean toStorage = storage.storageNodes().containsKey(connection.getTargetId());
if (source == null || target == null || fromStorage == toStorage) {
continue;
}
it.cnr.isti.workflow.manager.ios.IODescriptor output = descriptor(source.getOutputs(), connection.getSourceName());
it.cnr.isti.workflow.manager.ios.IODescriptor input = descriptor(target.getInputs(), connection.getTargetName());
if (output == null || input == null || fits(output, input)) {
continue;
}
errors.add(new ValidationError(ValidationErrorCode.STORAGE_CONNECTION_INVALID, "connection", connection.getId(),
"targetName", name(source, source.getId()) + "." + output.getName() + " gives " + describe(output) + ", and "
+ name(target, target.getId()) + "." + input.getName() + " takes " + describe(input),
List.of(source.getId(), target.getId())));
}
return errors;
}
/**
* A storage operation on a node of the flow: the node is there, what the block does suits the
* node's type, and the node's storage allows it.
* A storage operation and the node it is linked to: the node is there, what the operation does
* suits the node's type, and the node's storage allows it.
*/
private List<ValidationError> storageOperationReferences(StorageFlows.Compiled storage, List<Block<?>> operations) {
List<ValidationError> errors = new ArrayList<>();
for (Block<?> block : operations) {
StorageOperationBlockConfiguration configuration = (StorageOperationBlockConfiguration) block.getSpecificConfiguration();
if (configuration.getStorageNode() == null || configuration.getStorageNode().isBlank()) {
continue;
}
Block<?> node = storage.storageNodes().get(configuration.getStorageNode());
if (node == null) {
errors.add(new ValidationError(ValidationErrorCode.STORAGE_CONNECTION_INVALID, "block", block.getId(),
"specificConfiguration.storageNode", "There is no such storage node in this flow"));
"specificConfiguration.storageNode", "The storage node this operation is linked to is not in the flow"));
continue;
}
it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeConfiguration nodeConfiguration =
@ -609,14 +552,16 @@ public class FlowExecutionValidator {
}
StorageOperation op = configuration.effectiveOperation();
List<String> problems = switch (op) {
case READ -> type.get().validateQuery(configuration.getTemplate(),
configuration.getAs() != null ? configuration.getAs() : type.get().viewShapes().getFirst());
case READ -> type.get().validateQuery(configuration.getTemplate(), configuration.effectiveShape());
case WRITE -> type.get().validateTemplate(configuration.getTemplate());
case DELETE -> type.get().validateDelete(configuration.getTemplate());
case LIST -> List.of();
};
String field = op == StorageOperation.READ && !type.get().viewShapes().contains(configuration.effectiveShape())
? "specificConfiguration.as" : "specificConfiguration.template";
problems.forEach(problem -> errors.add(new ValidationError(ValidationErrorCode.VALIDATION_ERROR, "block", block.getId(),
"specificConfiguration.template", problem)));
problem.contains("cannot return") || problem.contains("returns rows as JSON") ? "specificConfiguration.as" : field,
problem)));
boolean changes = op == StorageOperation.WRITE || op == StorageOperation.DELETE;
if (changes && !nodeConfiguration.usesPersonalConnection() && storageInstances != null) {
storageInstances.find(nodeConfiguration.getInstance()).filter(instance -> instance.readOnly())
@ -627,69 +572,6 @@ public class FlowExecutionValidator {
return errors;
}
/**
* A view read by some step needs every one of its parameters fed: the node is not a step, so
* nobody can type a value into it when the execution starts.
*/
private List<ValidationError> unfedViewParameters(StorageFlows.Compiled storage) {
List<ValidationError> errors = new ArrayList<>();
Set<String> reported = new HashSet<>();
for (StorageFlows.Read read : storage.reads()) {
Block<?> node = storage.storageNodes().get(read.storageId());
if (node == null || !(node.getSpecificConfiguration()
instanceof it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeConfiguration configuration)) {
continue;
}
configuration.view(read.viewName()).ifPresent(view -> {
for (String parameter : it.cnr.isti.workflow.manager.storage.StoragePlaceholders.parameters(view.query())) {
boolean fed = storage.parameterFeeds().stream().anyMatch(feed -> feed.storageId().equals(node.getId())
&& feed.viewName().equals(view.name()) && feed.parameter().equals(parameter));
String port = StorageFlows.parameterPort(view.name(), parameter);
if (!fed && reported.add(node.getId() + "\u0000" + port)) {
errors.add(new ValidationError(ValidationErrorCode.STORAGE_CONNECTION_INVALID, "block", node.getId(),
"inputs." + port, "The view " + view.name() + " of " + name(node, node.getId()) + " reads ${{"
+ parameter + "}}: connect a step's output to " + port,
List.of(node.getId(), read.consumerId())));
}
}
});
}
return errors;
}
private static boolean fits(it.cnr.isti.workflow.manager.ios.IODescriptor output, it.cnr.isti.workflow.manager.ios.IODescriptor input) {
if (output.isMultiple() && !input.isMultiple()) {
// Not even into ANY: an input that takes one value refuses a list of them.
return false;
}
if (input.getType() == it.cnr.isti.workflow.manager.ios.IOType.ANY) {
return accepts(input, output);
}
return input.getType() == output.getType()
|| (input.getType() == it.cnr.isti.workflow.manager.ios.IOType.FILE && output.getType() == it.cnr.isti.workflow.manager.ios.IOType.CSV);
}
/** An ANY input lists the kinds it takes; one that lists none takes anything. */
private static boolean accepts(it.cnr.isti.workflow.manager.ios.IODescriptor input, it.cnr.isti.workflow.manager.ios.IODescriptor output) {
if (input.getValueKinds() == null || input.getValueKinds().isEmpty()
|| output.getType() == it.cnr.isti.workflow.manager.ios.IOType.ANY) {
return true;
}
it.cnr.isti.workflow.manager.blocks.IOCapabilityType kind = it.cnr.isti.workflow.manager.blocks.IOCapabilityTypes.from(output.getType());
return input.getValueKinds().stream().anyMatch(capability -> capability.type() == kind
|| capability.type() == it.cnr.isti.workflow.manager.blocks.IOCapabilityType.ANY);
}
private static it.cnr.isti.workflow.manager.ios.IODescriptor descriptor(List<it.cnr.isti.workflow.manager.ios.IODescriptor> ports, String name) {
return ports == null || name == null ? null
: ports.stream().filter(port -> name.equals(port.getName())).findFirst().orElse(null);
}
private static String describe(it.cnr.isti.workflow.manager.ios.IODescriptor port) {
return (port.isMultiple() ? "a list of " : "one ") + port.getType().name().toLowerCase(java.util.Locale.ROOT)
+ (port.isMultiple() ? " values" : " value");
}
private static String name(FlowNode node, String id) {
return node == null || node.getName() == null || node.getName().isBlank() ? id : node.getName();
}

View File

@ -4,79 +4,36 @@
package it.cnr.isti.workflow.manager.storage;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeView;
import it.cnr.isti.workflow.manager.flows.model.Connection;
import it.cnr.isti.workflow.manager.flows.model.Dependency;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.flows.model.FlowNode;
/**
* What storage nodes mean for the rest of the flow.
* What storage nodes mean for the rest of the flow: nothing the engine runs.
*
* <p>A storage node never runs, so it cannot be a node of the graph that decides what runs when -
* a step writing to it and one reading from it would otherwise look like a path through it, and a
* loop reading what it wrote the round before like a cycle with no way out. {@link #compile} takes
* it out, and says instead what its edges were:
*
* <ul>
* <li>a connection into one of its targets is a {@link Write} the source step does when it finishes;
* <li>a connection out of one of its views is a {@link Read} the consumer does before it starts;
* <li>a connection into a view's parameter is a real ordering - the consumer cannot read before the
* parameter exists - so it becomes a connection from the parameter's source to the consumer, into an
* input of its own ({@link #syntheticInput}).
* </ul>
*
* <p>The graph that is left is the one the engine runs and every check of it - loops, deadlocks,
* branches, one source per input - looks at, so none of them needs to know storage nodes exist. A
* flow without any comes back as the very same object.
* <p>A storage node is where a flow's storage operations are pointed - its connection, its scope,
* what is prepared at the start - and never a step. {@link #compile} takes it out of the flow the
* engine runs and every check of the graph looks at. Nothing connects to it: a storage operation
* names the node it works on, and the editor draws that as a link, which is not a connection.
*/
public final class StorageFlows {
/** Inputs a step gets for the view parameters it depends on: never a port of its block. */
public static final String SYNTHETIC_PREFIX = it.cnr.isti.workflow.manager.executions.steps.StepStorageHooks.SYNTHETIC_PREFIX;
private StorageFlows() {
}
public record Write(String sourceId, String outputName, String storageId, String targetName) {
}
public record Read(String consumerId, String inputName, String storageId, String viewName) {
}
public record ParameterFeed(String sourceId, String outputName, String storageId, String viewName, String parameter) {
}
/**
* @param executionFlow the flow without its storage nodes, with a connection for each view
* parameter a consumer waits for
* @param storageNodes by id
* @param strayConnections connections that mean nothing here - from one storage node straight
* into another, or into a port it does not have - for the validation to report
*/
public record Compiled(FlowData executionFlow, Map<String, Block<?>> storageNodes, List<Write> writes, List<Read> reads,
List<ParameterFeed> parameterFeeds, Map<String, Set<String>> syntheticInputs, List<Connection> strayConnections) {
/** @param executionFlow the flow without its storage nodes; @param storageNodes by id */
public record Compiled(FlowData executionFlow, Map<String, Block<?>> storageNodes) {
public boolean hasStorage() {
return !storageNodes.isEmpty();
}
public List<Write> writesFrom(String stepId) {
return writes.stream().filter(write -> write.sourceId().equals(stepId)).toList();
}
public List<Read> readsBy(String stepId) {
return reads.stream().filter(read -> read.consumerId().equals(stepId)).toList();
}
}
public static boolean isStorageNode(FlowNode node) {
@ -87,84 +44,21 @@ public final class StorageFlows {
return flow != null && flow.getBlocks() != null && flow.getBlocks().stream().anyMatch(StorageFlows::isStorageNode);
}
/** The node's input port for one parameter of one of its views. */
public static String parameterPort(String viewName, String parameter) {
return viewName + "." + parameter;
}
public static String syntheticInput(String storageId, String viewName, String parameter) {
return SYNTHETIC_PREFIX + storageId + "/" + viewName + "/" + parameter;
}
public static boolean isSynthetic(String inputName) {
return inputName != null && inputName.startsWith(SYNTHETIC_PREFIX);
}
public static FlowData executionView(FlowData flow) {
return compile(flow).executionFlow();
}
/** A flow without storage nodes comes back as the very same object. */
public static Compiled compile(FlowData flow) {
if (!hasStorage(flow)) {
return new Compiled(flow, Map.of(), List.of(), List.of(), List.of(), Map.of(), List.of());
return new Compiled(flow, Map.of());
}
Map<String, Block<?>> storageNodes = new LinkedHashMap<>();
flow.getBlocks().stream().filter(StorageFlows::isStorageNode).forEach(block -> storageNodes.put(block.getId(), block));
List<Connection> kept = new ArrayList<>();
List<Write> writes = new ArrayList<>();
List<Read> reads = new ArrayList<>();
List<ParameterFeed> feeds = new ArrayList<>();
List<Connection> stray = new ArrayList<>();
for (Connection connection : flow.getConnections() == null ? List.<Connection>of() : flow.getConnections()) {
if (connection == null) {
continue;
}
Block<?> source = storageNodes.get(connection.getSourceId());
Block<?> target = storageNodes.get(connection.getTargetId());
if (source == null && target == null) {
kept.add(connection);
} else if (source != null && target != null) {
stray.add(connection);
} else if (target != null) {
StorageNodeConfiguration configuration = (StorageNodeConfiguration) target.getSpecificConfiguration();
String port = connection.getTargetName();
if (port != null && configuration.target(port).isPresent()) {
writes.add(new Write(connection.getSourceId(), connection.getSourceName(), target.getId(), port));
} else {
ParameterFeed feed = parameterFeed(connection, target.getId(), configuration);
if (feed == null) {
stray.add(connection);
} else {
feeds.add(feed);
}
}
} else {
StorageNodeConfiguration configuration = (StorageNodeConfiguration) source.getSpecificConfiguration();
if (connection.getSourceName() != null && configuration.view(connection.getSourceName()).isPresent()) {
reads.add(new Read(connection.getTargetId(), connection.getTargetName(), source.getId(), connection.getSourceName()));
} else {
stray.add(connection);
}
}
}
Map<String, Set<String>> synthetic = new LinkedHashMap<>();
Set<String> seen = new LinkedHashSet<>();
for (ParameterFeed feed : feeds) {
for (Read read : reads) {
if (!read.storageId().equals(feed.storageId()) || !read.viewName().equals(feed.viewName())) {
continue;
}
String input = syntheticInput(feed.storageId(), feed.viewName(), feed.parameter());
if (seen.add(feed.sourceId() + "\u0000" + feed.outputName() + "\u0000" + read.consumerId() + "\u0000" + input)) {
kept.add(Connection.builder().sourceId(feed.sourceId()).sourceName(feed.outputName())
.targetId(read.consumerId()).targetName(input).build());
synthetic.computeIfAbsent(read.consumerId(), ignored -> new LinkedHashSet<>()).add(input);
}
}
}
List<Connection> connections = (flow.getConnections() == null ? List.<Connection>of() : flow.getConnections()).stream()
.filter(connection -> connection != null && !storageNodes.containsKey(connection.getSourceId())
&& !storageNodes.containsKey(connection.getTargetId()))
.toList();
List<Dependency> dependencies = (flow.getDependencies() == null ? List.<Dependency>of() : flow.getDependencies()).stream()
.filter(dependency -> dependency != null && !storageNodes.containsKey(dependency.getSourceId())
&& !storageNodes.containsKey(dependency.getTargetId()))
@ -172,29 +66,10 @@ public final class StorageFlows {
FlowData executionFlow = new FlowData(
flow.getBlocks().stream().filter(block -> !storageNodes.containsKey(block.getId())).toList(),
flow.getContainers() == null ? List.of() : flow.getContainers(),
List.copyOf(kept),
connections,
dependencies,
flow.getGlobalInputs() == null ? List.of() : flow.getGlobalInputs(),
flow.getLanes() == null ? List.of() : flow.getLanes());
return new Compiled(executionFlow, Map.copyOf(storageNodes), List.copyOf(writes), List.copyOf(reads),
List.copyOf(feeds), synthetic, List.copyOf(stray));
}
private static ParameterFeed parameterFeed(Connection connection, String storageId, StorageNodeConfiguration configuration) {
String port = connection.getTargetName();
if (port == null) {
return null;
}
for (StorageNodeView view : configuration.viewList()) {
if (view == null || view.name() == null) {
continue;
}
for (String parameter : StoragePlaceholders.parameters(view.query())) {
if (port.equals(parameterPort(view.name(), parameter))) {
return new ParameterFeed(connection.getSourceId(), connection.getSourceName(), storageId, view.name(), parameter);
}
}
}
return null;
return new Compiled(executionFlow, Map.copyOf(storageNodes));
}
}

View File

@ -5,21 +5,17 @@
package it.cnr.isti.workflow.manager.storage;
import java.util.Map;
import java.util.function.Function;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeTarget;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeView;
import it.cnr.isti.workflow.manager.executions.ExecutionRuntimeContextSupport;
import it.cnr.isti.workflow.manager.storage.validation.StorageNodeValidator;
/**
* Carries out what a storage node means, for whichever step is on the other end of its edges:
* prepare it, write a target, read a view. Each opens a session of its own and closes it after -
* a node's edges are used at different moments by different steps, on different threads.
* Reaches a storage node's storage, for the execution that prepares it and for the storage
* operations linked to it. Each use opens a session of its own and closes it after - the
* operations on a node run at different moments, on different threads.
*/
@Component
public class StorageNodeRuntime {
@ -46,27 +42,6 @@ public class StorageNodeRuntime {
}
}
public StorageReadResult read(Block<?> node, StorageNodeView view, Function<String, Object> values,
Map<String, Object> authorizations, Map<String, Object> executionVariables) {
StorageNodeConfiguration configuration = configuration(node);
StorageType type = types.require(configuration.getStorageType());
try (StorageSession session = type.open(connection(type, configuration, authorizations, executionVariables),
scope(configuration, executionVariables))) {
return session.read(view.query(), StorageNodeValidator.shape(view, type), values);
}
}
public StorageWriteResult write(Block<?> node, StorageNodeTarget target, Object value,
Map<String, Object> authorizations, Map<String, Object> executionVariables) {
StorageNodeConfiguration configuration = configuration(node);
StorageType type = types.require(configuration.getStorageType());
Function<String, Object> executionValues = name -> executionVariables == null ? null : executionVariables.get(name);
try (StorageSession session = type.open(connection(type, configuration, authorizations, executionVariables),
scope(configuration, executionVariables))) {
return session.write(target.template(), value, target.contentType(), StoragePlaceholders.withValue(value, executionValues));
}
}
/** A session on the node's storage, as the node reaches it, for an operation block to use. */
public StorageSession open(Block<?> node, Map<String, Object> authorizations, Map<String, Object> executionVariables) {
StorageNodeConfiguration configuration = configuration(node);

View File

@ -5,30 +5,23 @@
package it.cnr.isti.workflow.manager.storage.validation;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Optional;
import java.util.Set;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeTarget;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeView;
import it.cnr.isti.workflow.manager.storage.StorageFlows;
import it.cnr.isti.workflow.manager.storage.StorageInstancesProvider;
import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
import it.cnr.isti.workflow.manager.storage.StorageType;
import it.cnr.isti.workflow.manager.storage.StorageTypes;
import it.cnr.isti.workflow.manager.storage.ViewShape;
import jakarta.validation.ConstraintValidator;
import jakarta.validation.ConstraintValidatorContext;
/**
* Asks the storage type and the catalog whether a storage node holds together: its init, each
* target and view, the names its ports get, and what the storage allows. Like
* {@link StorageOperationValidator}, it only judges under Spring's validator.
* Asks the storage type and the catalog whether a storage node holds together: its type, where its
* storage is, its scope and what it prepares. Like {@link StorageOperationValidator}, it only judges
* under Spring's validator.
*/
@Component
public class StorageNodeValidator implements ConstraintValidator<ValidStorageNode, StorageNodeConfiguration> {
@ -67,7 +60,6 @@ public class StorageNodeValidator implements ConstraintValidator<ValidStorageNod
return problems;
}
StorageType type = found.get();
boolean writes = !configuration.targetList().isEmpty();
boolean prepares = configuration.init().createIfMissing() || configuration.init().hasScript();
if (!configuration.usesPersonalConnection()) {
if (configuration.getInstance() == null || configuration.getInstance().isBlank()) {
@ -78,9 +70,6 @@ public class StorageNodeValidator implements ConstraintValidator<ValidStorageNod
problems.add(new String[] { "instance", instance.displayName() + " is a " + instance.type()
+ " storage, not " + type.getName() });
}
if (instance.readOnly() && writes) {
problems.add(new String[] { "targets", instance.displayName() + " is read-only: it can have views, not targets" });
}
if (!instance.allowInit() && prepares) {
problems.add(new String[] { "initScript", instance.displayName() + " may not be prepared by flows: "
+ "leave Create if missing and the init script empty" });
@ -93,44 +82,9 @@ public class StorageNodeValidator implements ConstraintValidator<ValidStorageNod
+ "a per-execution scope is for object stores" });
}
type.validateInit(configuration.init()).forEach(problem -> problems.add(new String[] { "initScript", problem }));
Set<String> ports = new HashSet<>();
for (StorageNodeTarget target : configuration.targetList()) {
if (target == null || target.name() == null) {
continue;
}
if (!ports.add(target.name())) {
problems.add(new String[] { "targets", "Two ports are called " + target.name() });
}
type.validateTemplate(target.template()).forEach(problem -> problems.add(new String[] { "targets",
target.name() + ": " + problem }));
StoragePlaceholders.parameters(target.template()).forEach(parameter -> problems.add(new String[] { "targets",
target.name() + ": ${{" + parameter + "}} is not known when writing; a target has ${{value}}, "
+ "${{value.field}}, ${{context.*}} and ${{global.*}}" }));
}
for (StorageNodeView view : configuration.viewList()) {
if (view == null || view.name() == null) {
continue;
}
if (!ports.add(view.name())) {
problems.add(new String[] { "views", "Two ports are called " + view.name() });
}
ViewShape shape = shape(view, type);
type.validateQuery(view.query(), shape).forEach(problem -> problems.add(new String[] { "views",
view.name() + ": " + problem }));
for (String parameter : StoragePlaceholders.parameters(view.query())) {
if (!ports.add(StorageFlows.parameterPort(view.name(), parameter))) {
problems.add(new String[] { "views", "Two ports are called " + StorageFlows.parameterPort(view.name(), parameter) });
}
}
}
return problems;
}
public static ViewShape shape(StorageNodeView view, StorageType type) {
return view.as() != null ? view.as() : type.viewShapes().getFirst();
}
private static String escape(String message) {
return message.replace("\\", "\\\\").replace("{", "\\{").replace("}", "\\}").replace("$", "\\$");
}

View File

@ -6,129 +6,52 @@ package it.cnr.isti.workflow.manager.storage.validation;
import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration;
import it.cnr.isti.workflow.manager.storage.StorageInstancesProvider;
import it.cnr.isti.workflow.manager.storage.StorageOperation;
import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
import it.cnr.isti.workflow.manager.storage.StorageType;
import it.cnr.isti.workflow.manager.storage.StorageTypes;
import it.cnr.isti.workflow.manager.storage.ViewShape;
import jakarta.validation.ConstraintValidator;
import jakarta.validation.ConstraintValidatorContext;
/**
* Asks the storage type whether it can do what the block says, and the catalog whether the storage
* exists and allows it - each problem reported on the field it is about.
*
* <p>Needs the storage beans, so it only judges under Spring's validator. A plain one, with nothing
* to ask, lets the configuration through: the block is checked again when it runs.
* What can be said of a storage operation on its own: that it is linked to a node, and that its
* placeholders make sense for what it does. Whether its template suits the linked node's storage is
* known only with the flow, and is checked there (FlowExecutionValidator).
*/
@Component
public class StorageOperationValidator implements ConstraintValidator<ValidStorageOperation, StorageOperationBlockConfiguration> {
@Autowired(required = false)
StorageTypes types;
@Autowired(required = false)
StorageInstancesProvider instances;
@Override
public boolean isValid(StorageOperationBlockConfiguration configuration, ConstraintValidatorContext context) {
if (configuration == null || types == null || instances == null) {
if (configuration == null) {
return true;
}
List<Problem> problems = problems(configuration);
List<String[]> problems = new ArrayList<>();
if (configuration.getStorageNode() == null || configuration.getStorageNode().isBlank()) {
problems.add(new String[] { "storageNode", "Link this operation to a Storage node" });
}
StorageOperation operation = configuration.effectiveOperation();
String template = configuration.getTemplate();
if (template != null) {
StoragePlaceholders.invalidNames(template).forEach(name -> problems.add(new String[] { "template",
"${{" + name + "}} is not a usable name: start with a letter, then letters, digits, '-', '_' or '.'" }));
if (operation != StorageOperation.WRITE
&& StoragePlaceholders.names(template).stream().anyMatch(StoragePlaceholders::isValueName)) {
problems.add(new String[] { "template", "${{value}} is what a WRITE writes; a " + operation + " has none" });
}
}
if (problems.isEmpty()) {
return true;
}
context.disableDefaultConstraintViolation();
for (Problem problem : problems) {
context.buildConstraintViolationWithTemplate(escape(problem.message()))
.addPropertyNode(problem.field())
.addConstraintViolation();
for (String[] problem : problems) {
context.buildConstraintViolationWithTemplate(escape(problem[1])).addPropertyNode(problem[0]).addConstraintViolation();
}
return false;
}
private record Problem(String field, String message) {
}
List<Problem> problems(StorageOperationBlockConfiguration configuration) {
List<Problem> problems = new ArrayList<>();
if (configuration.usesStorageNode()) {
// What the node allows is known only with the flow: FlowExecutionValidator checks it there.
if (isBlank(configuration.getStorageNode())) {
problems.add(new Problem("storageNode", "Choose a storage node of this flow"));
}
return problems;
}
if (isBlank(configuration.getStorageType())) {
problems.add(new Problem("storageType", "Choose a storage type"));
return problems;
}
Optional<StorageType> found = types.find(configuration.getStorageType());
if (found.isEmpty()) {
problems.add(new Problem("storageType", "Unknown storage type " + configuration.getStorageType()
+ "; this service has " + types.all().stream().map(StorageType::getName).toList()));
return problems;
}
StorageType type = found.get();
StorageOperation operation = configuration.effectiveOperation();
if (configuration.usesPersonalConnection()) {
// Whose connection, and so what it allows, is only known when someone runs the flow.
} else if (isBlank(configuration.getInstance())) {
problems.add(new Problem("instance", "Choose a storage of the catalog, or use your own connection"));
} else {
instances.find(configuration.getInstance()).ifPresentOrElse(instance -> {
if (!instance.type().equalsIgnoreCase(type.getName())) {
problems.add(new Problem("instance", instance.displayName() + " is a " + instance.type()
+ " storage, not " + type.getName()));
}
if (instance.readOnly() && (operation == StorageOperation.WRITE || operation == StorageOperation.DELETE)) {
problems.add(new Problem("operation", instance.displayName() + " is read-only: it allows READ and LIST"));
}
}, () -> problems.add(new Problem("instance", "There is no storage " + configuration.getInstance() + " in the catalog")));
}
String template = configuration.getTemplate();
List<String> templateProblems = switch (operation) {
case READ -> type.validateQuery(template, shape(configuration, type));
case WRITE -> type.validateTemplate(template);
case DELETE -> type.validateDelete(template);
case LIST -> listProblems(template);
};
templateProblems.forEach(problem -> problems.add(new Problem("template", problem)));
if (operation != StorageOperation.WRITE && template != null
&& StoragePlaceholders.names(template).stream().anyMatch(StoragePlaceholders::isValueName)) {
problems.add(new Problem("template", "${{value}} is what a WRITE writes; a " + operation + " has none"));
}
if (configuration.perExecution() && !type.supportsExecutionScope()) {
problems.add(new Problem("scope", "A " + type.getName() + " storage is shared by every execution; "
+ "a per-execution scope is for object stores"));
}
return problems;
}
public static ViewShape shape(StorageOperationBlockConfiguration configuration, StorageType type) {
return configuration.getAs() != null ? configuration.getAs() : type.viewShapes().getFirst();
}
private static List<String> listProblems(String template) {
List<String> problems = new ArrayList<>();
StoragePlaceholders.invalidNames(template).forEach(name -> problems.add(
"${{" + name + "}} is not a usable name: start with a letter, then letters, digits, '-', '_' or '.'"));
return problems;
}
private static boolean isBlank(String value) {
return value == null || value.isBlank();
}
/** Bean Validation would read {..} and $ in a message as its own syntax; ours are literal. */
private static String escape(String message) {
return message.replace("\\", "\\\\").replace("{", "\\{").replace("}", "\\}").replace("$", "\\$");
}

View File

@ -0,0 +1,43 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.blocks.configurations;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.TestPropertySource;
import tools.jackson.databind.JsonNode;
/** What the editor is told about storage nodes and the operations linked to them. */
@SpringBootTest
@TestPropertySource(locations = "classpath:test.properties")
class StorageSchemaTest {
@Autowired
private JsonSchemaProducer schemaProducer;
@Test
void theNodeAnOperationWorksOnIsNotAFieldButTheLinkDrawnToIt() {
JsonNode properties = schemaProducer.generateSchemaNode(StorageOperationBlockConfiguration.class).get("properties");
assertTrue(properties.get("storageNode").path("x-ui-hidden").asBoolean(), "set by drawing the link, not typed in");
assertFalse(properties.get("template").has("x-ui-hidden"));
assertFalse(properties.has("source") || properties.has("instance") || properties.has("storageType"),
"where the storage is, is the node's to say");
}
@Test
void aNodeSaysWhereItsStorageIsAndHasNoTargetsOrViews() {
JsonNode properties = schemaProducer.generateSchemaNode(StorageNodeConfiguration.class).get("properties");
assertEquals("CATALOG", properties.get("source").get("default").asString());
assertFalse(properties.has("targets") || properties.has("views"));
}
}

View File

@ -36,12 +36,12 @@ import it.cnr.isti.workflow.manager.blocks.configurations.ConditionalBlockConfig
import it.cnr.isti.workflow.manager.blocks.configurations.EndBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeTarget;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeView;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.factories.ConditionalBlockFactory;
import it.cnr.isti.workflow.manager.blocks.factories.EndBlockFactory;
import it.cnr.isti.workflow.manager.blocks.factories.LLMBlockFactory;
import it.cnr.isti.workflow.manager.blocks.factories.StorageNodeFactory;
import it.cnr.isti.workflow.manager.blocks.factories.StorageOperationBlockFactory;
import it.cnr.isti.workflow.manager.executions.ExecutionContext;
import it.cnr.isti.workflow.manager.executions.ExecutionObject;
import it.cnr.isti.workflow.manager.executions.ExecutionStatus;
@ -58,16 +58,17 @@ import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator;
import it.cnr.isti.workflow.manager.flows.validation.ValidationError;
import it.cnr.isti.workflow.manager.flows.validation.ValidationErrorCode;
import it.cnr.isti.workflow.manager.ios.IODescriptor;
import it.cnr.isti.workflow.manager.ios.IOType;
import it.cnr.isti.workflow.manager.llms.ChatMessage;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import it.cnr.isti.workflow.manager.llms.providers.LLMProvider;
/**
* Storage nodes wired in at random - writers, readers, view parameters, loops, dependencies - and
* every flow the validation accepts run for real. Each must finish, in SUCCESS or stopped by a loop
* at its limit: a storage edge must never leave a step waiting, nor make a flow the engine cannot
* finish look valid. The storage is in memory, so this is about the engine and the validation, not
* about any store.
* Storage operations on a storage node wired into random flows - writers, readers, a read waiting on
* a parameter, loops, dependencies either way - and every flow the validation accepts run for real.
* Each must finish, in SUCCESS or stopped by a loop at its limit: a storage node must never leave a
* step waiting, nor make a flow the engine cannot finish look valid. The storage is in memory, so
* this is about the engine and the validation, not about any store.
*/
@SpringBootTest
@TestPropertySource(locations = "classpath:test.properties")
@ -243,6 +244,9 @@ class StorageFuzzTest {
@Autowired
StorageNodeFactory storageFactory;
@Autowired
StorageOperationBlockFactory operationFactory;
@Test
void everyFlowWithStorageTheValidationAcceptsRunsToAnEnd() {
Random random = new Random(SEED);
@ -266,7 +270,7 @@ class StorageFuzzTest {
}
accepted++;
StorageFlows.Compiled compiled = StorageFlows.compile(flow);
if (!compiled.writes().isEmpty() || !compiled.reads().isEmpty()) {
if (flow.getBlocks().stream().anyMatch(block -> block.getSpecificConfiguration() instanceof StorageOperationBlockConfiguration)) {
acceptedWithStorageEdges++;
}
if (!FlowLoops.analyze(compiled.executionFlow()).loops().isEmpty()) {
@ -277,9 +281,9 @@ class StorageFuzzTest {
failures.add("trial " + trial + ": " + problem + "\n" + describe(flow));
}
}
System.out.printf("storage fuzz: %d flows, %d accepted (%d with storage edges, %d with loops), refused %s%n",
System.out.printf("storage fuzz: %d flows, %d accepted (%d with storage operations, %d with loops), refused %s%n",
FLOWS, accepted, acceptedWithStorageEdges, acceptedWithLoops, refused);
assertTrue(acceptedWithStorageEdges > FLOWS / 20, "too few flows with storage edges accepted: " + acceptedWithStorageEdges);
assertTrue(acceptedWithStorageEdges > FLOWS / 20, "too few flows with storage operations accepted: " + acceptedWithStorageEdges);
assertTrue(refused.containsKey(ValidationErrorCode.STORAGE_ORDER_UNDEFINED.name()), "the ordering rule never fired: " + refused);
assertTrue(failures.isEmpty(), failures.size() + " accepted flows did not run to an end:\n\n"
+ String.join("\n\n", failures.subList(0, Math.min(20, failures.size()))));
@ -331,11 +335,13 @@ class StorageFuzzTest {
}
private FlowData randomFlow(Random random) {
int size = 2 + random.nextInt(4);
Block<?> storage = storageFactory.create(StorageNodeConfiguration.builder().name("store")
.storageType("Memory").instance("mem").scope(StorageNodeConfiguration.SCOPE_PER_EXECUTION).build());
int size = 3 + random.nextInt(4);
List<Block<?>> steps = new ArrayList<>();
for (int i = 0; i < size; i++) {
String name = "n" + i;
steps.add(switch (random.nextInt(4)) {
steps.add(switch (random.nextInt(7)) {
case 0 -> llmFactory.create(LLMBlockConfiguration.builder().name(name + "-llm").llmDescriptor(ECHO)
.prompt("Go <<${{a}}x>>").build());
case 1 -> llmFactory.create(LLMBlockConfiguration.builder().name(name + "-reader").llmDescriptor(ECHO)
@ -344,21 +350,19 @@ class StorageFuzzTest {
.condition(List.of("${{response}}.length() > 3", "${{response}}.length() >= 0",
"${{response}}.length() % 2 == 0").get(random.nextInt(3)))
.outputTemplate("${{response}}").build());
case 3, 4 -> operationFactory.create(StorageOperationBlockConfiguration.builder().name(name + "-write")
.storageNode(storage.getId()).operation(StorageOperation.WRITE)
.template("k/${{context.iteration}}/" + random.nextInt(1000)).build());
case 5 -> operationFactory.create(StorageOperationBlockConfiguration.builder().name(name + "-read")
.storageNode(storage.getId()).operation(StorageOperation.READ).as(ViewShape.TEXTS)
.template(random.nextBoolean() ? "k/" : "k/${{from}}").build());
default -> endFactory.create(EndBlockConfiguration.builder().name(name + "-end")
.outcomeCode("END" + i).outcomeLabel("End " + i).build());
});
}
boolean parameterized = random.nextBoolean();
Block<?> storage = storageFactory.create(StorageNodeConfiguration.builder().name("store")
.storageType("Memory").instance("mem").scope("PER_EXECUTION")
.targets(List.of(new StorageNodeTarget("put", "k/${{context.iteration}}/" + random.nextInt(1000), null)))
.views(List.of(new StorageNodeView("all", parameterized ? "k/${{from}}" : "k/", null)))
.build());
FlowData.FlowDataBuilder builder = FlowData.builder();
FlowData.FlowDataBuilder builder = FlowData.builder().block(storage);
steps.forEach(builder::block);
builder.block(storage);
int connections = 1 + random.nextInt(size);
int connections = 1 + random.nextInt(size * 2);
for (int i = 0; i < connections; i++) {
Block<?> source = steps.get(random.nextInt(steps.size()));
Block<?> target = steps.get(random.nextInt(steps.size()));
@ -366,46 +370,23 @@ class StorageFuzzTest {
continue;
}
IODescriptor output = source.getOutputs().get(random.nextInt(source.getOutputs().size()));
List<IODescriptor> open = target.getInputs().stream().filter(port -> !port.getName().equals("seen")).toList();
if (open.isEmpty()) {
List<IODescriptor> fitting = target.getInputs().stream().filter(input -> fits(output, input)).toList();
if (fitting.isEmpty()) {
continue;
}
IODescriptor input = open.get(random.nextInt(open.size()));
Connection connection = connect(source, output.getName(), target, input.getName());
Connection connection = connect(source, output.getName(), target, fitting.get(random.nextInt(fitting.size())).getName());
if (random.nextInt(4) == 0) {
connection.setLoop(new LoopEdgeSettings(1 + random.nextInt(3)));
}
builder.connection(connection);
}
// Storage edges wired the way the editor would let them be: text into the target, the
// view's list into an input that takes a list, or into an end that takes anything.
List<Block<?>> producers = steps.stream().filter(step -> !step.getOutputs().isEmpty()).toList();
List<Block<?>> readers = steps.stream().filter(step -> step.getInputs().stream()
.anyMatch(port -> port.getName().equals("seen") || port.getName().equals(EndBlockFactory.INPUT_NAME))).toList();
List<Block<?>> writers = new ArrayList<>();
for (int i = 0; i < 1 + random.nextInt(2) && !producers.isEmpty(); i++) {
Block<?> writer = producers.get(random.nextInt(producers.size()));
writers.add(writer);
builder.connection(connect(writer, writer.getOutputs().get(random.nextInt(writer.getOutputs().size())).getName(), storage, "put"));
}
List<Block<?>> reading = new ArrayList<>();
for (int i = 0; i < 1 + random.nextInt(2) && !readers.isEmpty(); i++) {
Block<?> reader = readers.get(random.nextInt(readers.size()));
reading.add(reader);
String input = reader.getInputs().stream().anyMatch(port -> port.getName().equals("seen")) ? "seen" : EndBlockFactory.INPUT_NAME;
builder.connection(connect(storage, "all", reader, input));
}
if (parameterized && !producers.isEmpty()) {
Block<?> source = producers.get(random.nextInt(producers.size()));
builder.connection(connect(source, source.getOutputs().getFirst().getName(), storage, "all.from"));
}
if (!writers.isEmpty() && !reading.isEmpty() && random.nextBoolean()) {
List<Block<?>> writers = steps.stream().filter(step -> step.getName().endsWith("-write")).toList();
List<Block<?>> readers = steps.stream().filter(step -> step.getName().endsWith("-read")).toList();
if (!writers.isEmpty() && !readers.isEmpty() && random.nextInt(3) > 0) {
Block<?> writer = writers.get(random.nextInt(writers.size()));
Block<?> reader = reading.get(random.nextInt(reading.size()));
if (writer != reader) {
builder.dependency(random.nextInt(4) == 0 ? new Dependency(reader.getId(), writer.getId())
: new Dependency(writer.getId(), reader.getId()));
}
Block<?> reader = readers.get(random.nextInt(readers.size()));
builder.dependency(random.nextInt(4) == 0 ? new Dependency(reader.getId(), writer.getId())
: new Dependency(writer.getId(), reader.getId()));
}
for (int i = 0; i < random.nextInt(2); i++) {
Block<?> source = steps.get(random.nextInt(steps.size()));
@ -417,6 +398,14 @@ class StorageFuzzTest {
return builder.build();
}
/** What the editor lets through: a list only into a list, JSON only where JSON or anything goes. */
private static boolean fits(IODescriptor output, IODescriptor input) {
if (output.isMultiple() != input.isMultiple()) {
return false;
}
return input.getType() == IOType.ANY || input.getType() == output.getType();
}
private static Connection connect(Block<?> source, String output, Block<?> target, String input) {
return Connection.builder().sourceId(source.getId()).sourceName(output).targetId(target.getId()).targetName(input).build();
}

View File

@ -5,7 +5,7 @@
package it.cnr.isti.workflow.manager.storage;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.nio.file.Files;
@ -13,6 +13,7 @@ import java.nio.file.Path;
import java.time.Duration;
import java.time.Instant;
import java.util.List;
import java.util.UUID;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
@ -22,6 +23,7 @@ import org.springframework.context.annotation.Bean;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.springframework.test.context.TestPropertySource;
import org.springframework.web.server.ResponseStatusException;
import org.testcontainers.containers.MinIOContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
@ -32,15 +34,12 @@ import it.cnr.isti.workflow.manager.blocks.configurations.ConditionalBlockConfig
import it.cnr.isti.workflow.manager.blocks.configurations.EndBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.LLMBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeConfiguration;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeTarget;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageNodeView;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.factories.ConditionalBlockFactory;
import it.cnr.isti.workflow.manager.blocks.factories.EndBlockFactory;
import it.cnr.isti.workflow.manager.blocks.factories.LLMBlockFactory;
import it.cnr.isti.workflow.manager.blocks.factories.StorageNodeFactory;
import it.cnr.isti.workflow.manager.blocks.factories.StorageOperationBlockFactory;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration;
import it.cnr.isti.workflow.manager.executions.ExecutionEvent;
import it.cnr.isti.workflow.manager.executions.ExecutionEventType;
import it.cnr.isti.workflow.manager.executions.ExecutionObject;
import it.cnr.isti.workflow.manager.executions.ExecutionStatus;
@ -49,18 +48,25 @@ import it.cnr.isti.workflow.manager.executions.FieldKey;
import it.cnr.isti.workflow.manager.flows.model.Connection;
import it.cnr.isti.workflow.manager.flows.model.Dependency;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.flows.model.LoopEdgeSettings;
import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator;
import it.cnr.isti.workflow.manager.flows.validation.ValidationError;
import it.cnr.isti.workflow.manager.flows.validation.ValidationErrorCode;
import it.cnr.isti.workflow.manager.ios.IOType;
import it.cnr.isti.workflow.manager.llms.ChatMessage;
import it.cnr.isti.workflow.manager.llms.LLMDescriptor;
import it.cnr.isti.workflow.manager.llms.providers.LLMProvider;
import it.cnr.isti.workflow.manager.storage.postgres.PostgresStorageType;
import it.cnr.isti.workflow.manager.storage.s3.S3StorageType;
import it.cnr.isti.workflow.manager.vault.UserSecretService;
import it.cnr.isti.workflow.manager.vault.model.VaultSecretCreateRequest;
import it.cnr.isti.workflow.manager.vault.model.VaultSecretView;
import tools.jackson.databind.JsonNode;
import tools.jackson.databind.ObjectMapper;
/**
* Storage nodes run for real: prepared at the start, written by the step connected to a target
* before it counts as done, read by the step a view feeds just before it runs.
* Storage nodes and the operations linked to them, run for real against a MinIO and a PostgreSQL:
* the node prepared at the start, each operation a step on the node's connection and scope.
*/
@SpringBootTest
@TestPropertySource(locations = "classpath:test.properties")
@ -92,6 +98,8 @@ class StorageNodeExecutionTest {
POSTGRES.getHost(), POSTGRES.getFirstMappedPort(), POSTGRES.getDatabaseName(), POSTGRES.getUsername(),
POSTGRES.getPassword()));
registry.add("app.storage.instances.file", file::toString);
// The containers are on this machine: a personal connection to them has to be let through.
registry.add("app.storage.endpoint.allowed-host-patterns", MINIO::getHost);
}
@TestConfiguration
@ -132,7 +140,10 @@ class StorageNodeExecutionTest {
FlowExecutionValidator validator;
@Autowired
StorageNodeFactory storageFactory;
StorageNodeFactory nodeFactory;
@Autowired
StorageOperationBlockFactory operationFactory;
@Autowired
LLMBlockFactory llmFactory;
@ -144,71 +155,74 @@ class StorageNodeExecutionTest {
EndBlockFactory endFactory;
@Autowired
StorageOperationBlockFactory operationFactory;
UserSecretService secrets;
@Test
void theStepThatDependsOnTheWriterReadsWhatItWroteFromABucketMadeAtTheStart() {
Block<?> notes = notesStorage("notes-" + System.nanoTime());
Block<?> writer = llm("write-note", "<<${{note}}>>");
Block<?> reader = llm("read-notes", "<<${{seen[]}}>>");
FlowData flow = FlowData.builder().block(notes).block(writer).block(reader)
.connection(connect(writer, LLMBlockFactory.OUTPUT_NAME, notes, "notes"))
.connection(connect(notes, "all", reader, "seen"))
.dependency(new Dependency(writer.getId(), reader.getId()))
.build();
void anOperationThatDependsOnTheWriterReadsWhatItWroteToABucketMadeAtTheStart() {
Block<?> notes = s3Node("notes", "node-files", true);
Block<?> write = operation("save-note", notes, StorageOperation.WRITE, "notes/${{title}}.txt", null);
Block<?> read = operation("read-notes", notes, StorageOperation.READ, "notes/*.txt", ViewShape.TEXTS);
assertEquals(List.of("title", "content"), write.getInputs().stream().map(input -> input.getName()).toList());
assertTrue(notes.getInputs().isEmpty() && notes.getOutputs().isEmpty(), "a storage node has no ports");
FlowData flow = FlowData.builder().block(notes).block(write).block(read)
.dependency(new Dependency(write.getId(), read.getId())).build();
assertTrue(validator.collectErrors(flow).isEmpty(), () -> validator.collectErrors(flow).toString());
ExecutionObject execution = executionsService.createExecution("Write and read", flow);
executionsService.prepareInput(execution.getId(), writer.getId(), "note", "the list forgets items");
ExecutionObject execution = executionsService.createExecution("Write then read", flow);
executionsService.prepareInput(execution.getId(), write.getId(), "title", "round-1");
executionsService.prepareInput(execution.getId(), write.getId(), "content", "The list forgets items on reload");
executionsService.startExecution(execution.getId());
ExecutionObject finished = awaitFinal(execution.getId());
assertEquals(ExecutionStatus.SUCCESS, finished.getContext().getStatus(), () -> finished.getContext().getErrors().toString());
assertEquals("the list forgets items", result(finished, reader));
assertTrue(events(finished).stream().anyMatch(event -> event.getMessage().startsWith("Wrote response to notes.notes: notes/1.txt")),
() -> events(finished).toString());
assertTrue(events(finished).stream().anyMatch(event -> event.getMessage().startsWith("Read all of notes: 1")));
assertEquals(List.of("The list forgets items on reload"), result(finished, read));
assertTrue(finished.getContext().getEvents().stream().anyMatch(event -> event.getType() == ExecutionEventType.STORAGE_OPERATION
&& event.getMessage().equals("WRITE on notes: notes/round-1.txt")), () -> finished.getContext().getEvents().toString());
}
@Test
void aWriterAndAReaderOfTheSameStorageMustBeOrderedOneWayOrTheOther() {
Block<?> notes = notesStorage("unordered");
Block<?> writer = llm("write-note", "<<${{note}}>>");
Block<?> reader = llm("read-notes", "<<${{seen[]}}>>");
FlowData.FlowDataBuilder unordered = FlowData.builder().block(notes).block(writer).block(reader)
.connection(connect(writer, LLMBlockFactory.OUTPUT_NAME, notes, "notes"))
.connection(connect(notes, "all", reader, "seen"));
void aSqlWriteBindsTheFieldsOfItsContentAndAReadReturnsTheRows() throws Exception {
Block<?> db = nodeFactory.create(StorageNodeConfiguration.builder().name("db").storageType(PostgresStorageType.NAME)
.instance("node-db").initScript("CREATE TABLE IF NOT EXISTS rounds (round int, verdict text)").build());
Block<?> write = operation("record-round", db, StorageOperation.WRITE,
"INSERT INTO rounds(round, verdict) VALUES (${{value.round}}, ${{value.verdict}})", null);
Block<?> read = operation("rejected-rounds", db, StorageOperation.READ,
"SELECT round FROM rounds WHERE verdict = ${{verdict}} ORDER BY round", ViewShape.JSON);
assertEquals(IOType.JSON, read.getOutputs().getFirst().getType());
FlowData flow = FlowData.builder().block(db).block(write).block(read)
.dependency(new Dependency(write.getId(), read.getId())).build();
List<ValidationError> errors = validator.collectErrors(unordered.build());
ValidationError order = errors.stream().filter(error -> error.code() == ValidationErrorCode.STORAGE_ORDER_UNDEFINED)
.findFirst().orElseThrow(() -> new AssertionError(errors.toString()));
assertEquals(reader.getId(), order.id());
assertTrue(order.message().contains("Add a dependency write-note → read-notes"), order.message());
ExecutionObject execution = executionsService.createExecution("Record a round", flow);
executionsService.prepareInput(execution.getId(), write.getId(), "content",
new ObjectMapper().readTree("{\"round\": 7, \"verdict\": \"rejected\"}"));
executionsService.prepareInput(execution.getId(), read.getId(), "verdict", "rejected");
executionsService.startExecution(execution.getId());
ExecutionObject finished = awaitFinal(execution.getId());
assertTrue(validator.collectErrors(unordered.dependency(new Dependency(writer.getId(), reader.getId())).build()).isEmpty());
assertTrue(validator.collectErrors(FlowData.builder().block(notes).block(writer).block(reader)
.connection(connect(writer, LLMBlockFactory.OUTPUT_NAME, notes, "notes"))
.connection(connect(notes, "all", reader, "seen"))
.dependency(new Dependency(reader.getId(), writer.getId())).build()).isEmpty(),
"reading first is a choice as well");
assertEquals(ExecutionStatus.SUCCESS, finished.getContext().getStatus(), () -> finished.getContext().getErrors().toString());
assertEquals(7, ((JsonNode) result(finished, read)).get(0).get("round").asInt(), "the table was made by the node's init script");
}
@Test
void inALoopEachRoundReadsWhatTheRoundsBeforeItWrote() {
Block<?> drafts = storage("drafts", StorageNodeConfiguration.builder()
.storageType(S3StorageType.NAME).instance("node-files").createIfMissing(true)
.scope("PER_EXECUTION")
.targets(List.of(new StorageNodeTarget("drafts", "drafts/${{context.iteration}}.txt", null)))
.views(List.of(new StorageNodeView("earlier", "drafts/*.txt", ViewShape.TEXTS))));
Block<?> draft = llm("draft", "<<${{topic}}x>> after ${{earlier[]}}");
void inALoopEachRoundWritesAKeyOfItsOwn() {
Block<?> drafts = s3Node("drafts", "node-files", true);
Block<?> draft = llm("draft", "<<${{topic}}x>>");
Block<?> write = operation("keep-draft", drafts, StorageOperation.WRITE, "drafts/${{context.iteration}}.txt", null);
Block<?> check = conditionalFactory.create(ConditionalBlockConfiguration.builder().name("long-enough")
.condition("${{response}}.length() >= 4").outputTemplate("${{response}}").build());
Block<?> done = endFactory.create(EndBlockConfiguration.builder().name("done").outcomeCode("DONE").outcomeLabel("Done").build());
FlowData flow = FlowData.builder().block(drafts).block(draft).block(check).block(done)
Block<?> folder = llm("folder", "<<drafts/>> after ${{accepted}}");
Block<?> read = operation("all-drafts", drafts, StorageOperation.READ, "${{folder}}*.txt", ViewShape.TEXTS);
Connection back = connect(check, ConditionalBlockFactory.FALSE_OUTPUT, draft, "topic");
back.setLoop(new LoopEdgeSettings(5));
FlowData flow = FlowData.builder().block(drafts).block(draft).block(write).block(check).block(folder).block(read)
.connection(connect(draft, LLMBlockFactory.OUTPUT_NAME, write, StorageOperationBlockFactory.CONTENT_INPUT))
.connection(connect(draft, LLMBlockFactory.OUTPUT_NAME, check, "response"))
.connection(connect(check, ConditionalBlockFactory.FALSE_OUTPUT, draft, "topic"))
.connection(connect(check, ConditionalBlockFactory.TRUE_OUTPUT, done, EndBlockFactory.INPUT_NAME))
.connection(connect(draft, LLMBlockFactory.OUTPUT_NAME, drafts, "drafts"))
.connection(connect(drafts, "earlier", draft, "earlier"))
// The check waits for the write, so each round's draft is kept before the next begins.
.dependency(new Dependency(write.getId(), check.getId()))
.connection(back)
// Out of the loop through the guard's way out, so the read comes after every round.
.connection(connect(check, ConditionalBlockFactory.TRUE_OUTPUT, folder, "accepted"))
.connection(connect(folder, LLMBlockFactory.OUTPUT_NAME, read, "folder"))
.build();
assertTrue(validator.collectErrors(flow).isEmpty(), () -> validator.collectErrors(flow).toString());
@ -218,131 +232,97 @@ class StorageNodeExecutionTest {
ExecutionObject finished = awaitFinal(execution.getId());
assertEquals(ExecutionStatus.SUCCESS, finished.getContext().getStatus(), () -> finished.getContext().getErrors().toString());
assertEquals(3, finished.getContext().getSteps().get(draft.getId()).getIteration());
List<Object> readCounts = events(finished).stream().filter(event -> event.getMessage().startsWith("Read earlier"))
.map(event -> event.getDetails().get("count")).toList();
assertEquals(List.of(0, 1, 2), readCounts, "round n reads the n - 1 drafts written before it");
assertEquals(List.of("ax", "axx", "axxx"), result(finished, read), "one object per round, keyed by its round");
}
@Test
void aViewParameterFedByAStepMakesTheReaderWaitForIt() {
Block<?> people = storage("people", StorageNodeConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("node-db")
.initScript("""
CREATE TABLE IF NOT EXISTS people (name text, note text);
DELETE FROM people;
INSERT INTO people VALUES ('ann', 'writes the tests'), ('bob', 'reviews');""")
.views(List.of(new StorageNodeView("byName", "SELECT note FROM people WHERE name = ${{name}}", null))));
assertEquals(List.of("byName.name"), people.getInputs().stream().map(input -> input.getName()).toList());
Block<?> who = llm("who", "<<${{person}}>>");
Block<?> found = endFactory.create(EndBlockConfiguration.builder().name("found").outcomeCode("FOUND").outcomeLabel("Found").build());
FlowData flow = FlowData.builder().block(people).block(who).block(found)
.connection(connect(who, LLMBlockFactory.OUTPUT_NAME, people, "byName.name"))
.connection(connect(people, "byName", found, EndBlockFactory.INPUT_NAME))
.build();
void aNodeOnAPersonalConnectionIsAskedForAtTheStartAndItsOperationsWorkOnIt() {
String owner = "storage-owner-" + UUID.randomUUID();
String connection = """
{"endpoint": "%s", "bucket": "node-files", "accessKey": "%s", "secretKey": "%s"}"""
.formatted(MINIO.getS3URL(), MINIO.getUserName(), MINIO.getPassword());
VaultSecretView saved = secrets.create(owner, new VaultSecretCreateRequest("my minio", S3StorageType.NAME, null, connection, null));
assertEquals(MINIO.getS3URL() + "/node-files", saved.endpoint(), "the address is shown, the keys are not");
ExecutionObject execution = executionsService.createExecution("Look up", flow);
executionsService.prepareInput(execution.getId(), who.getId(), "person", "ann");
Block<?> mine = nodeFactory.create(StorageNodeConfiguration.builder().name("mine").storageType(S3StorageType.NAME)
.source(StorageNodeConfiguration.SOURCE_PERSONAL).createIfMissing(true).build());
Block<?> write = operation("save-mine", mine, StorageOperation.WRITE, "mine/${{context.rootExecutionId}}.txt", null);
ExecutionObject execution = executionsService.createExecution("Mine", FlowData.builder().block(mine).block(write).build(), owner);
String key = StorageCredentials.authorizationKey(S3StorageType.NAME);
assertTrue(execution.getRequiredAuthorizations().stream().anyMatch(requirement -> requirement.key().equals(key)
&& requirement.isVaultBacked()));
assertEquals(List.of(key), execution.getMissingAuthorizationKeys());
executionsService.setAuthorizationValue(execution.getId(), key, saved.id());
executionsService.prepareInput(execution.getId(), write.getId(), "content", "written with my keys");
executionsService.startExecution(execution.getId());
ExecutionObject finished = awaitFinal(execution.getId());
assertEquals(ExecutionStatus.SUCCESS, finished.getContext().getStatus(), () -> finished.getContext().getErrors().toString());
assertTrue(String.valueOf(finished.getContext().getOutcomes().getFirst().payload()).contains("writes the tests"),
() -> String.valueOf(finished.getContext().getOutcomes()));
assertEquals("mine/" + execution.getId() + ".txt", ((JsonNode) result(finished, write)).get("location").asString());
}
@Test
void aWriteThatFailsFailsTheStepThatProducedTheValue() {
Block<?> notes = storage("keyed", StorageNodeConfiguration.builder()
.storageType(S3StorageType.NAME).instance("node-files").createIfMissing(true)
.targets(List.of(new StorageNodeTarget("byId", "notes/${{value.id}}.txt", null))));
Block<?> writer = llm("write-note", "<<${{note}}>>");
FlowData flow = FlowData.builder().block(notes).block(writer)
.connection(connect(writer, LLMBlockFactory.OUTPUT_NAME, notes, "byId")).build();
ExecutionObject execution = executionsService.createExecution("No id", flow);
executionsService.prepareInput(execution.getId(), writer.getId(), "note", "plain text has no id");
executionsService.startExecution(execution.getId());
ExecutionObject finished = awaitFinal(execution.getId());
assertEquals(ExecutionStatus.ERROR, finished.getContext().getStatus());
assertEquals("STORAGE_WRITE_FAILED", finished.getContext().getErrorCodes().get(writer.getId()));
String error = finished.getContext().getErrors().get(writer.getId());
assertTrue(error.contains("Could not write response to keyed.byId") && error.contains("${{value.id}}"), error);
void theVaultRefusesAPersonalConnectionThatDoesNotHoldTogether() {
ResponseStatusException notJson = assertThrows(ResponseStatusException.class,
() -> secrets.create("someone", new VaultSecretCreateRequest("bad", S3StorageType.NAME, null, "just a key", null)));
assertTrue(notJson.getReason().contains("not JSON"), notJson.getReason());
ResponseStatusException privateAddress = assertThrows(ResponseStatusException.class,
() -> secrets.create("someone", new VaultSecretCreateRequest("bad", PostgresStorageType.NAME, null,
"{\"host\": \"10.1.2.3\", \"database\": \"d\", \"user\": \"u\", \"password\": \"p\"}", null)));
assertTrue(privateAddress.getReason().contains("private, local or otherwise disallowed"), privateAddress.getReason());
}
@Test
void theNodeIsCheckedAgainstWhatItsStorageAllows() {
Block<?> readOnly = storage("archive", StorageNodeConfiguration.builder()
.storageType(S3StorageType.NAME).instance("archive").createIfMissing(true)
.targets(List.of(new StorageNodeTarget("keep", "kept.txt", null))));
Block<?> quoted = storage("quoted", StorageNodeConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("node-db")
.views(List.of(new StorageNodeView("byName", "SELECT * FROM people WHERE name = '${{name}}'", null))));
void aWriterAndAReaderOfTheSameNodeMustBeOrderedOneWayOrTheOther() {
Block<?> notes = s3Node("notes", "node-files", false);
Block<?> write = operation("write-note", notes, StorageOperation.WRITE, "notes/a.txt", null);
Block<?> read = operation("read-notes", notes, StorageOperation.READ, "notes/*.txt", ViewShape.TEXTS);
List<ValidationError> errors = validator.collectErrors(FlowData.builder().block(readOnly).block(quoted).build());
List<ValidationError> errors = validator.collectErrors(FlowData.builder().block(notes).block(write).block(read).build());
ValidationError order = errors.stream().filter(error -> error.code() == ValidationErrorCode.STORAGE_ORDER_UNDEFINED)
.findFirst().orElseThrow(() -> new AssertionError(errors.toString()));
assertEquals(read.getId(), order.id());
assertTrue(order.message().contains("Add a dependency write-note → read-notes"), order.message());
assertHas(errors, readOnly, "specificConfiguration.targets", "is read-only");
assertHas(errors, readOnly, "specificConfiguration.initScript", "may not be prepared by flows");
assertHas(errors, quoted, "specificConfiguration.views", "${{name}} is inside a string literal");
assertTrue(validator.collectErrors(FlowData.builder().block(notes).block(write).block(read)
.dependency(new Dependency(read.getId(), write.getId())).build()).isEmpty(), "reading first is a choice as well");
}
@Test
void aViewMustFitWhereItIsReadAndHaveEveryParameterFed() {
Block<?> notes = notesStorage("fit");
Block<?> end = endFactory.create(EndBlockConfiguration.builder().name("end").outcomeCode("E").outcomeLabel("E").build());
List<ValidationError> listIntoOne = validator.collectErrors(FlowData.builder().block(notes).block(end)
.connection(connect(notes, "all", end, EndBlockFactory.INPUT_NAME)).build());
assertTrue(listIntoOne.stream().anyMatch(error -> error.code() == ValidationErrorCode.STORAGE_CONNECTION_INVALID
&& error.message().contains("gives a list of text values, and end.input takes one any value")), listIntoOne::toString);
Block<?> people = storage("people", StorageNodeConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("node-db")
.views(List.of(new StorageNodeView("byName", "SELECT note FROM people WHERE name = ${{name}}", null))));
List<ValidationError> unfed = validator.collectErrors(FlowData.builder().block(people).block(end)
.connection(connect(people, "byName", end, EndBlockFactory.INPUT_NAME)).build());
assertHas(unfed, people, "inputs.byName.name", "connect a step's output to byName.name");
}
@Test
void aStorageOperationOnANodeSharesItsConnectionAndScope() {
Block<?> notes = notesStorage("shared");
Block<?> writer = llm("write-note", "<<${{note}}>>");
Block<?> list = operationFactory.create(StorageOperationBlockConfiguration.builder().name("list-notes")
.source(StorageOperationBlockConfiguration.SOURCE_NODE).storageNode(notes.getId())
.operation(StorageOperation.LIST).template("notes/*").build());
FlowData flow = FlowData.builder().block(notes).block(writer).block(list)
.connection(connect(writer, LLMBlockFactory.OUTPUT_NAME, notes, "notes"))
.dependency(new Dependency(writer.getId(), list.getId()))
.build();
assertTrue(validator.collectErrors(flow).isEmpty(), () -> validator.collectErrors(flow).toString());
ExecutionObject execution = executionsService.createExecution("List the node", flow);
executionsService.prepareInput(execution.getId(), writer.getId(), "note", "one");
executionsService.startExecution(execution.getId());
ExecutionObject finished = awaitFinal(execution.getId());
assertEquals(ExecutionStatus.SUCCESS, finished.getContext().getStatus(), () -> finished.getContext().getErrors().toString());
assertEquals(List.of("notes/1.txt"), finished.getContext().getResult()
.get(new FieldKey(list.getId(), StorageOperationBlockFactory.RESULT_OUTPUT)), "the node's per-execution scope is the operation's");
}
@Test
void aStorageOperationOnANodeIsOrderedAndCheckedLikeAnyReaderOrWriter() {
Block<?> notes = notesStorage("checked");
Block<?> writer = llm("write-note", "<<${{note}}>>");
Block<?> list = operationFactory.create(StorageOperationBlockConfiguration.builder().name("list-notes")
.source(StorageOperationBlockConfiguration.SOURCE_NODE).storageNode(notes.getId())
.operation(StorageOperation.LIST).template("notes/*").build());
void anOperationIsCheckedAgainstTheNodeItIsLinkedTo() {
Block<?> archive = s3Node("archive", "archive", false);
Block<?> db = nodeFactory.create(StorageNodeConfiguration.builder().name("db").storageType(PostgresStorageType.NAME)
.instance("node-db").build());
Block<?> intoArchive = operation("write-archive", archive, StorageOperation.WRITE, "x.txt", null);
Block<?> quoted = operation("quoted", db, StorageOperation.READ, "SELECT * FROM t WHERE a = '${{a}}'", ViewShape.JSON);
Block<?> asFiles = operation("as-files", db, StorageOperation.READ, "SELECT 1", ViewShape.FILES);
Block<?> unlinked = operationFactory.create(StorageOperationBlockConfiguration.builder().name("unlinked")
.operation(StorageOperation.LIST).build());
Block<?> missing = operationFactory.create(StorageOperationBlockConfiguration.builder().name("elsewhere")
.source(StorageOperationBlockConfiguration.SOURCE_NODE).storageNode("no-such-node")
.operation(StorageOperation.READ).template("x.txt").build());
List<ValidationError> errors = validator.collectErrors(FlowData.builder().block(notes).block(writer).block(list).block(missing)
.connection(connect(writer, LLMBlockFactory.OUTPUT_NAME, notes, "notes")).build());
.storageNode("no-such-node").operation(StorageOperation.LIST).build());
assertTrue(errors.stream().anyMatch(error -> error.code() == ValidationErrorCode.STORAGE_ORDER_UNDEFINED
&& error.id().equals(list.getId())), errors::toString);
assertHas(errors, missing, "specificConfiguration.storageNode", "no such storage node");
List<ValidationError> errors = validator.collectErrors(FlowData.builder().block(archive).block(db)
.block(intoArchive).block(quoted).block(asFiles).block(unlinked).block(missing).build());
assertHas(errors, intoArchive, "specificConfiguration.operation", "is read-only");
assertHas(errors, quoted, "specificConfiguration.template", "${{a}} is inside a string literal");
assertHas(errors, asFiles, "specificConfiguration.as", "cannot return FILES");
assertHas(errors, unlinked, "specificConfiguration.storageNode", "Link this operation to a Storage node");
assertHas(errors, missing, "specificConfiguration.storageNode", "is not in the flow");
}
@Test
void aNodeIsCheckedAgainstWhatItsStorageAllows() {
Block<?> preparing = nodeFactory.create(StorageNodeConfiguration.builder().name("archive").storageType(S3StorageType.NAME)
.instance("archive").createIfMissing(true).build());
Block<?> scoped = nodeFactory.create(StorageNodeConfiguration.builder().name("db").storageType(PostgresStorageType.NAME)
.instance("node-db").scope(StorageNodeConfiguration.SCOPE_PER_EXECUTION).build());
List<ValidationError> errors = validator.collectErrors(FlowData.builder().block(preparing).block(scoped).build());
assertHas(errors, preparing, "specificConfiguration.initScript", "may not be prepared by flows");
assertHas(errors, scoped, "specificConfiguration.scope", "shared by every execution");
}
private static void assertHas(List<ValidationError> errors, Block<?> block, String field, String text) {
@ -350,15 +330,16 @@ class StorageNodeExecutionTest {
&& error.message().contains(text)), () -> "No error on " + field + " containing '" + text + "': " + errors);
}
private Block<?> notesStorage(String bucketSuffix) {
return storage("notes", StorageNodeConfiguration.builder()
.storageType(S3StorageType.NAME).instance("node-files").createIfMissing(true).scope("PER_EXECUTION")
.targets(List.of(new StorageNodeTarget("notes", "notes/${{context.iteration}}.txt", null)))
.views(List.of(new StorageNodeView("all", "notes/*.txt", ViewShape.TEXTS))));
private Block<?> s3Node(String name, String instance, boolean perExecution) {
return nodeFactory.create(StorageNodeConfiguration.builder().name(name).storageType(S3StorageType.NAME).instance(instance)
.createIfMissing(instance.equals("node-files"))
.scope(perExecution ? StorageNodeConfiguration.SCOPE_PER_EXECUTION : StorageNodeConfiguration.SCOPE_SHARED)
.build());
}
private Block<?> storage(String name, StorageNodeConfiguration.StorageNodeConfigurationBuilder builder) {
return storageFactory.create(builder.name(name).build());
private Block<?> operation(String name, Block<?> node, StorageOperation operation, String template, ViewShape as) {
return operationFactory.create(StorageOperationBlockConfiguration.builder().name(name).storageNode(node.getId())
.operation(operation).template(template).as(as).build());
}
private Block<?> llm(String name, String prompt) {
@ -370,12 +351,7 @@ class StorageNodeExecutionTest {
}
private static Object result(ExecutionObject execution, Block<?> block) {
return execution.getContext().getResult().get(new FieldKey(block.getId(), LLMBlockFactory.OUTPUT_NAME));
}
private static List<ExecutionEvent> events(ExecutionObject execution) {
return execution.getContext().getEvents().stream()
.filter(event -> event.getType() == ExecutionEventType.STORAGE_OPERATION).toList();
return execution.getContext().getResult().get(new FieldKey(block.getId(), StorageOperationBlockFactory.RESULT_OUTPUT));
}
private ExecutionObject awaitFinal(String id) {
@ -385,7 +361,6 @@ class StorageNodeExecutionTest {
Thread.onSpinWait();
execution = executionsService.getExecution(id);
}
assertFalse(!execution.getContext().getStatus().isFinalState(), "the execution did not finish");
return execution;
}
}

View File

@ -1,256 +0,0 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.nio.file.Files;
import java.nio.file.Path;
import java.time.Duration;
import java.time.Instant;
import java.util.List;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.springframework.test.context.TestPropertySource;
import org.testcontainers.containers.MinIOContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.postgresql.PostgreSQLContainer;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.factories.StorageOperationBlockFactory;
import it.cnr.isti.workflow.manager.blocks.types.StorageOperationBlockType;
import it.cnr.isti.workflow.manager.executions.ExecutionEventType;
import it.cnr.isti.workflow.manager.executions.ExecutionObject;
import it.cnr.isti.workflow.manager.executions.FieldKey;
import it.cnr.isti.workflow.manager.executions.ExecutionStatus;
import it.cnr.isti.workflow.manager.executions.ExecutionsService;
import it.cnr.isti.workflow.manager.flows.model.Dependency;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator;
import it.cnr.isti.workflow.manager.flows.validation.ValidationError;
import it.cnr.isti.workflow.manager.storage.postgres.PostgresStorageType;
import it.cnr.isti.workflow.manager.storage.s3.S3StorageType;
import it.cnr.isti.workflow.manager.ios.IOType;
import it.cnr.isti.workflow.manager.vault.UserSecretService;
import it.cnr.isti.workflow.manager.vault.model.VaultSecretCreateRequest;
import it.cnr.isti.workflow.manager.vault.model.VaultSecretView;
import tools.jackson.databind.JsonNode;
/** The storage block run for real, against a MinIO and a PostgreSQL from the catalog. */
@SpringBootTest
@TestPropertySource(locations = "classpath:test.properties")
@Testcontainers(disabledWithoutDocker = true)
class StorageOperationExecutionTest {
@Container
static final MinIOContainer MINIO = new MinIOContainer("minio/minio:latest");
@Container
static final PostgreSQLContainer POSTGRES = new PostgreSQLContainer("postgres:17-alpine");
@DynamicPropertySource
static void catalog(DynamicPropertyRegistry registry) throws Exception {
Path file = Files.createTempFile("storages", ".json");
file.toFile().deleteOnExit();
Files.writeString(file, """
{ "storages": [
{ "id": "flow-files", "type": "S3", "allowInit": true,
"connection": { "endpoint": "%s", "bucket": "flow-files", "accessKey": "%s", "secretKey": "%s" } },
{ "id": "archive", "type": "S3", "readOnly": true,
"connection": { "endpoint": "%s", "bucket": "flow-files", "accessKey": "%s", "secretKey": "%s" } },
{ "id": "notes-db", "type": "PostgreSQL", "allowInit": true,
"connection": { "host": "%s", "port": "%d", "database": "%s", "user": "%s", "password": "%s", "sslMode": "disable" } }
] }""".formatted(MINIO.getS3URL(), MINIO.getUserName(), MINIO.getPassword(),
MINIO.getS3URL(), MINIO.getUserName(), MINIO.getPassword(),
POSTGRES.getHost(), POSTGRES.getFirstMappedPort(), POSTGRES.getDatabaseName(), POSTGRES.getUsername(),
POSTGRES.getPassword()));
registry.add("app.storage.instances.file", file::toString);
// The containers are on this machine: a personal connection to them has to be allowed through.
registry.add("app.storage.endpoint.allowed-host-patterns", MINIO::getHost);
}
@Autowired
ExecutionsService executionsService;
@Autowired
StorageOperationBlockFactory factory;
@Autowired
FlowExecutionValidator validator;
@Autowired
StorageTypes types;
@Autowired
StorageInstancesProvider instances;
@Autowired
UserSecretService secrets;
@Test
void aStepWritesAndTheOneThatDependsOnItReadsWhatWasWritten() throws Exception {
prepareBucket();
Block<StorageOperationBlockType> write = block("save-note", StorageOperationBlockConfiguration.builder()
.storageType(S3StorageType.NAME).instance("flow-files").operation(StorageOperation.WRITE)
.template("notes/${{title}}.txt").scope(StorageOperationBlockConfiguration.SCOPE_PER_EXECUTION));
Block<StorageOperationBlockType> read = block("read-notes", StorageOperationBlockConfiguration.builder()
.storageType(S3StorageType.NAME).instance("flow-files").operation(StorageOperation.READ)
.template("notes/*.txt").as(ViewShape.TEXTS).scope(StorageOperationBlockConfiguration.SCOPE_PER_EXECUTION));
assertEquals(List.of("title", "content"), write.getInputs().stream().map(input -> input.getName()).toList());
assertEquals(IOType.TEXT, read.getOutputs().getFirst().getType());
assertTrue(read.getOutputs().getFirst().isMultiple());
FlowData flow = FlowData.builder().block(write).block(read)
.dependency(new Dependency(write.getId(), read.getId())).build();
ExecutionObject execution = executionsService.createExecution("Write then read", flow);
executionsService.prepareInput(execution.getId(), write.getId(), "title", "round-1");
executionsService.prepareInput(execution.getId(), write.getId(), "content", "The list forgets items on reload");
executionsService.startExecution(execution.getId());
ExecutionObject finished = awaitFinal(execution.getId());
assertEquals(ExecutionStatus.SUCCESS, finished.getContext().getStatus(), () -> finished.getContext().getErrors().toString());
assertEquals(List.of("The list forgets items on reload"), result(finished, read));
assertTrue(finished.getContext().getEvents().stream().anyMatch(event -> event.getType() == ExecutionEventType.STORAGE_OPERATION
&& event.getMessage().contains("notes/round-1.txt")));
}
@Test
void aSqlWriteBindsTheFieldsOfItsContentAndAReadReturnsTheRows() throws Exception {
try (StorageSession session = types.require(PostgresStorageType.NAME).open(
instances.find("notes-db").orElseThrow().connection(types.require(PostgresStorageType.NAME)), StorageScope.shared())) {
session.initialize(new StorageInit(false, "CREATE TABLE IF NOT EXISTS rounds (round int, verdict text)"));
}
Block<StorageOperationBlockType> write = block("record-round", StorageOperationBlockConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("notes-db").operation(StorageOperation.WRITE)
.template("INSERT INTO rounds(round, verdict) VALUES (${{value.round}}, ${{value.verdict}})"));
Block<StorageOperationBlockType> read = block("rejected-rounds", StorageOperationBlockConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("notes-db").operation(StorageOperation.READ)
.template("SELECT round FROM rounds WHERE verdict = ${{verdict}} ORDER BY round"));
assertEquals(List.of("verdict"), read.getInputs().stream().map(input -> input.getName()).toList());
assertEquals(IOType.JSON, read.getOutputs().getFirst().getType());
FlowData flow = FlowData.builder().block(write).block(read)
.dependency(new Dependency(write.getId(), read.getId())).build();
ExecutionObject execution = executionsService.createExecution("Record a round", flow);
executionsService.prepareInput(execution.getId(), write.getId(), "content",
new tools.jackson.databind.ObjectMapper().readTree("{\"round\": 7, \"verdict\": \"rejected\"}"));
executionsService.prepareInput(execution.getId(), read.getId(), "verdict", "rejected");
executionsService.startExecution(execution.getId());
ExecutionObject finished = awaitFinal(execution.getId());
assertEquals(ExecutionStatus.SUCCESS, finished.getContext().getStatus(), () -> finished.getContext().getErrors().toString());
JsonNode rows = (JsonNode) result(finished, read);
assertEquals(7, rows.get(0).get("round").asInt());
}
@Test
void aPersonalConnectionIsAskedForAtTheStartAndTheStepWorksOnIt() {
prepareBucket();
String owner = "storage-owner-" + java.util.UUID.randomUUID();
String connection = """
{"endpoint": "%s", "bucket": "flow-files", "accessKey": "%s", "secretKey": "%s"}"""
.formatted(MINIO.getS3URL(), MINIO.getUserName(), MINIO.getPassword());
VaultSecretView saved = secrets.create(owner, new VaultSecretCreateRequest("my minio", S3StorageType.NAME, null, connection, null));
assertEquals(MINIO.getS3URL() + "/flow-files", saved.endpoint(), "the address is shown, the keys are not");
Block<StorageOperationBlockType> write = block("save-mine", StorageOperationBlockConfiguration.builder()
.storageType(S3StorageType.NAME).source(StorageOperationBlockConfiguration.SOURCE_PERSONAL)
.operation(StorageOperation.WRITE).template("mine/${{context.rootExecutionId}}.txt"));
ExecutionObject execution = executionsService.createExecution("Mine", FlowData.builder().block(write).build(), owner);
String key = StorageCredentials.authorizationKey(S3StorageType.NAME);
assertTrue(execution.getRequiredAuthorizations().stream().anyMatch(requirement -> requirement.key().equals(key)
&& requirement.isVaultBacked() && requirement.provider().equals(S3StorageType.NAME)));
assertEquals(List.of(key), execution.getMissingAuthorizationKeys());
executionsService.setAuthorizationValue(execution.getId(), key, saved.id());
executionsService.prepareInput(execution.getId(), write.getId(), "content", "written with my keys");
executionsService.startExecution(execution.getId());
ExecutionObject finished = awaitFinal(execution.getId());
assertEquals(ExecutionStatus.SUCCESS, finished.getContext().getStatus(), () -> finished.getContext().getErrors().toString());
JsonNode written = (JsonNode) result(finished, write);
assertEquals("mine/" + execution.getId() + ".txt", written.get("location").asString());
}
@Test
void theVaultRefusesAPersonalConnectionThatDoesNotHoldTogether() {
org.springframework.web.server.ResponseStatusException notJson = org.junit.jupiter.api.Assertions.assertThrows(
org.springframework.web.server.ResponseStatusException.class,
() -> secrets.create("someone", new VaultSecretCreateRequest("bad", S3StorageType.NAME, null, "just a key", null)));
assertTrue(notJson.getReason().contains("not JSON"), notJson.getReason());
org.springframework.web.server.ResponseStatusException noSecret = org.junit.jupiter.api.Assertions.assertThrows(
org.springframework.web.server.ResponseStatusException.class,
() -> secrets.create("someone", new VaultSecretCreateRequest("bad", S3StorageType.NAME, null,
"{\"endpoint\": \"https://minio.example.org\", \"bucket\": \"b-1\", \"accessKey\": \"a\"}", null)));
assertTrue(noSecret.getReason().contains("has no secretKey set"), noSecret.getReason());
org.springframework.web.server.ResponseStatusException privateAddress = org.junit.jupiter.api.Assertions.assertThrows(
org.springframework.web.server.ResponseStatusException.class,
() -> secrets.create("someone", new VaultSecretCreateRequest("bad", PostgresStorageType.NAME, null,
"{\"host\": \"10.1.2.3\", \"database\": \"d\", \"user\": \"u\", \"password\": \"p\"}", null)));
assertTrue(privateAddress.getReason().contains("private, local or otherwise disallowed"), privateAddress.getReason());
}
@Test
void theFlowIsRefusedWhereTheStorageCannotDoWhatTheBlockAsks() {
Block<StorageOperationBlockType> intoArchive = block("write-archive", StorageOperationBlockConfiguration.builder()
.storageType(S3StorageType.NAME).instance("archive").operation(StorageOperation.WRITE).template("x.txt"));
Block<StorageOperationBlockType> quoted = block("quoted", StorageOperationBlockConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("notes-db").operation(StorageOperation.READ)
.template("SELECT * FROM rounds WHERE verdict = '${{verdict}}'"));
Block<StorageOperationBlockType> wrongType = block("wrong-type", StorageOperationBlockConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("flow-files").operation(StorageOperation.LIST));
Block<StorageOperationBlockType> scoped = block("scoped-db", StorageOperationBlockConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("notes-db").operation(StorageOperation.LIST)
.scope(StorageOperationBlockConfiguration.SCOPE_PER_EXECUTION));
List<ValidationError> errors = validator.collectErrors(FlowData.builder()
.block(intoArchive).block(quoted).block(wrongType).block(scoped).build());
assertHasError(errors, intoArchive, "specificConfiguration.operation", "is read-only");
assertHasError(errors, quoted, "specificConfiguration.template", "${{verdict}} is inside a string literal");
assertHasError(errors, wrongType, "specificConfiguration.instance", "is a S3 storage, not PostgreSQL");
assertHasError(errors, scoped, "specificConfiguration.scope", "shared by every execution");
}
private static Object result(ExecutionObject execution, Block<?> block) {
return execution.getContext().getResult().get(new FieldKey(block.getId(), StorageOperationBlockFactory.RESULT_OUTPUT));
}
private static void assertHasError(List<ValidationError> errors, Block<?> block, String field, String text) {
assertTrue(errors.stream().anyMatch(error -> block.getId().equals(error.id()) && field.equals(error.field())
&& error.message().contains(text)), () -> "No error on " + field + " containing '" + text + "': " + errors);
}
private void prepareBucket() {
StorageType s3 = types.require(S3StorageType.NAME);
try (StorageSession session = s3.open(instances.find("flow-files").orElseThrow().connection(s3), StorageScope.shared())) {
session.initialize(new StorageInit(true, null));
}
}
private Block<StorageOperationBlockType> block(String name,
StorageOperationBlockConfiguration.StorageOperationBlockConfigurationBuilder builder) {
return factory.create(builder.name(name).build());
}
private ExecutionObject awaitFinal(String id) {
Instant deadline = Instant.now().plus(Duration.ofSeconds(20));
ExecutionObject execution = executionsService.getExecution(id);
while (!execution.getContext().getStatus().isFinalState() && Instant.now().isBefore(deadline)) {
Thread.onSpinWait();
execution = executionsService.getExecution(id);
}
return execution;
}
}