From a06a00758e701bfb89d64405d75965b7d3706192 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Mon, 28 Sep 2026 17:34:39 +0200 Subject: [PATCH] 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) --- .../configurations/JsonSchemaProducer.java | 52 +++ .../StorageNodeConfiguration.java | 87 ++--- .../configurations/StorageNodeTarget.java | 49 --- .../configurations/StorageNodeView.java | 52 --- .../StorageOperationBlockConfiguration.java | 92 +---- .../blocks/factories/StorageNodeFactory.java | 55 +-- .../StorageOperationBlockFactory.java | 26 +- .../configurations/annotations/UiHidden.java | 19 ++ .../retrievers/FlowNodesFieldRetriever.java | 57 ---- .../AuthorizationRequirementResolver.java | 8 +- .../manager/executions/ExecutionObject.java | 45 +-- .../manager/executions/ExecutionsService.java | 1 - .../executions/StorageStepHooksImpl.java | 132 ------- .../blocks/StorageOperationExecutor.java | 46 +-- .../manager/executions/steps/Input.java | 20 -- .../manager/executions/steps/Step.java | 82 +---- .../executions/steps/StepStorageHooks.java | 26 -- .../executions/steps/StepVariables.java | 57 ++++ .../capabilities/NodeTypeCapabilities.java | 8 +- .../validation/FlowExecutionValidator.java | 176 ++-------- .../manager/storage/StorageFlows.java | 155 +-------- .../manager/storage/StorageNodeRuntime.java | 31 +- .../validation/StorageNodeValidator.java | 52 +-- .../validation/StorageOperationValidator.java | 117 ++----- .../configurations/StorageSchemaTest.java | 43 +++ .../manager/storage/StorageFuzzTest.java | 97 +++--- .../storage/StorageNodeExecutionTest.java | 321 ++++++++---------- .../StorageOperationExecutionTest.java | 256 -------------- 28 files changed, 512 insertions(+), 1650 deletions(-) delete mode 100644 src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeTarget.java delete mode 100644 src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeView.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiHidden.java delete mode 100644 src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/FlowNodesFieldRetriever.java delete mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/StorageStepHooksImpl.java delete mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/steps/StepStorageHooks.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/executions/steps/StepVariables.java create mode 100644 src/test/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageSchemaTest.java delete mode 100644 src/test/java/it/cnr/isti/workflow/manager/storage/StorageOperationExecutionTest.java diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java index 7c75a87..7e8c9e6 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/JsonSchemaProducer.java @@ -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, Map> longTextMap = collectLongTextMetadata(type); Map, Set> acceptsPlaceholderMap = collectAcceptsVariablePlaceholderMetadata(type); Map, Set> defaultsWhenEmptyMap = collectDefaultsWhenEmptyMetadata(type); + Map, Set> hiddenMap = collectFlagMetadata(type, UiHidden.class); Map, Map> structuralMap = collectStructuralMetadata(type); Map, Map> uiOptionalGroupMap = collectUiOptionalGroupMetadata(type); Map, Map> 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, Set> collectFlagMetadata(Class rootClass, + Class annotation) { + Map, Set> result = new HashMap<>(); + Set> visited = new HashSet<>(); + Queue> queue = new ArrayDeque<>(); + queue.add(rootClass); + while (!queue.isEmpty()) { + Class current = queue.poll(); + if (current == null || !visited.add(current) || isTerminalType(current)) { + continue; + } + Set 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 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 names) { if (names == null || names.isEmpty()) { return; diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeConfiguration.java index 252e271..aa44f9a 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeConfiguration.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeConfiguration.java @@ -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. * - *

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. + *

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 { - 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 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 views = List.of(); - @Builder public StorageNodeConfiguration(@NonNull String name, String storageType, String source, String instance, String scope, - Boolean createIfMissing, String initScript, List targets, List 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 targetList() { - return targets == null ? List.of() : targets; - } - - public List viewList() { - return views == null ? List.of() : views; - } - - public Optional target(String name) { - return targetList().stream().filter(target -> target != null && name.equals(target.name())).findFirst(); - } - - public Optional view(String name) { - return viewList().stream().filter(view -> view != null && name.equals(view.name())).findFirst(); - } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeTarget.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeTarget.java deleted file mode 100644 index df990d7..0000000 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeTarget.java +++ /dev/null @@ -1,49 +0,0 @@ -// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - 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) { -} diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeView.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeView.java deleted file mode 100644 index 623c4bc..0000000 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageNodeView.java +++ /dev/null @@ -1,52 +0,0 @@ -// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - 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 .}, 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) { -} diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageOperationBlockConfiguration.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageOperationBlockConfiguration.java index 61eb0ce..a75bc89 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageOperationBlockConfiguration.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageOperationBlockConfiguration.java @@ -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. * - *

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. + *

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 { - 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.}. 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 { - private static final List PARAMETER_CAPABILITIES = List.of( - new IOCapability(IOCapabilityType.TEXT, false), new IOCapability(IOCapabilityType.JSON, false)); - private static final List 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 create(StorageNodeConfiguration configuration) { - StorageType type = types.find(configuration.getStorageType()).orElse(null); - List inputs = new ArrayList<>(); - List 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.builder() - .inputs(inputs) - .outputs(outputs) .specificConfiguration(configuration) .type(blockType) .build(); @@ -89,12 +43,11 @@ public class StorageNodeFactory implements BlockFactory supportedInputCapabilities() { - return DEFAULT_WRITABLE; + return List.of(); } @Override public List supportedOutputCapabilities() { - return List.of(new IOCapability(IOCapabilityType.FILE, true), new IOCapability(IOCapabilityType.TEXT, true), - new IOCapability(IOCapabilityType.JSON, false)); + return List.of(); } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/StorageOperationBlockFactory.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/StorageOperationBlockFactory.java index fb3fa91..62129a0 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/StorageOperationBlockFactory.java +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/StorageOperationBlockFactory.java @@ -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 *

    *
  • Every {@code ${{name}}} of the template that is not the execution's own ({@code context.*}, * {@code global.*}) nor the written value is an input. - *
  • A WRITE takes {@code content}, of the kinds the storage accepts. + *
  • 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. *
  • {@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. *
@@ -46,20 +44,17 @@ public class StorageOperationBlockFactory implements BlockFactory PARAMETER_CAPABILITIES = List.of( new IOCapability(IOCapabilityType.TEXT, false), new IOCapability(IOCapabilityType.JSON, false)); - private static final List DEFAULT_WRITABLE = List.of(new IOCapability(IOCapabilityType.FILE, false), + private static final List 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 create(StorageOperationBlockConfiguration configuration) { - StorageType type = types.find(configuration.getStorageType()).orElse(null); StorageOperation operation = configuration.effectiveOperation(); List inputs = new ArrayList<>(); for (String parameter : StoragePlaceholders.parameters(configuration.getTemplate())) { @@ -68,22 +63,21 @@ public class StorageOperationBlockFactory implements BlockFactorybuilder() .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 supportedInputCapabilities() { - return DEFAULT_WRITABLE; + return WRITABLE; } @Override diff --git a/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiHidden.java b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiHidden.java new file mode 100644 index 0000000..815b293 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/configurations/annotations/UiHidden.java @@ -0,0 +1,19 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - 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 { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/FlowNodesFieldRetriever.java b/src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/FlowNodesFieldRetriever.java deleted file mode 100644 index 6eed0f8..0000000 --- a/src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/FlowNodesFieldRetriever.java +++ /dev/null @@ -1,57 +0,0 @@ -// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - 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 retrieve(String parameter, Map 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(); - } -} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/AuthorizationRequirementResolver.java b/src/main/java/it/cnr/isti/workflow/manager/executions/AuthorizationRequirementResolver.java index 3f07559..b06998f 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/AuthorizationRequirementResolver.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/AuthorizationRequirementResolver.java @@ -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 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(); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java index 2c4b0da..0cce072 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionObject.java @@ -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 stepDependencies = new ArrayList<>(); @@ -155,7 +148,6 @@ public class ExecutionObject { : List.copyOf(executionFlow.getConnections()); FlowLoops.Analysis loops = FlowLoops.analyze(executionFlow); List> 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 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> steps) { - if (!this.storage.hasStorage()) { - return; - } - for (Step step : steps) { - List reads = this.storage.readsBy(step.getId()); - List 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> storageNode(String id) { return java.util.Optional.ofNullable(id == null ? null : this.storage.storageNodes().get(id)); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java index 3dba3c6..b9e254d 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionsService.java @@ -1613,7 +1613,6 @@ public class ExecutionsService { } private void attachPersistence(ExecutionObject executionObject) { - executionObject.configureStorage(storageNodeRuntime); executionObject.setStateChangeListener(() -> { persist(executionObject); cleanupManagedResourcesIfFinal(executionObject); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/StorageStepHooksImpl.java b/src/main/java/it/cnr/isti/workflow/manager/executions/StorageStepHooksImpl.java deleted file mode 100644 index 0bab550..0000000 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/StorageStepHooksImpl.java +++ /dev/null @@ -1,132 +0,0 @@ -// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - 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 reads; - private final List writes; - private final Map> storageNodes; - private final Supplier runtime; - - StorageStepHooksImpl(List reads, List writes, Map> storageNodes, - Supplier 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 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 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 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 details) { - log(step, message, details, ExecutionEventType.STORAGE_OPERATION); - } - - private static void log(Step step, String message, Map details, ExecutionEventType type) { - if (step.getEventLogger() != null) { - step.getEventLogger().info(type, message, details); - } - } -} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/StorageOperationExecutor.java b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/StorageOperationExecutor.java index 4f18b89..33ad966 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/StorageOperationExecutor.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/StorageOperationExecutor.java @@ -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 { - private final StorageTypes types; - private final StorageConnectionResolver connections; private final StorageNodeRuntime nodes; private final ObjectProvider executions; private final ObjectMapper objectMapper; - public StorageOperationExecutor(StorageTypes types, StorageConnectionResolver connections, StorageNodeRuntime nodes, - ObjectProvider executions, ObjectMapper objectMapper) { - this.types = types; - this.connections = connections; + public StorageOperationExecutor(StorageNodeRuntime nodes, ObjectProvider executions, + ObjectMapper objectMapper) { this.nodes = nodes; this.executions = executions; this.objectMapper = objectMapper; @@ -67,27 +58,21 @@ public class StorageOperationExecutor implements BlockExecutor authorizations, Map executionVariables, Map 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 values = values(inputs, executionVariables); long started = System.nanoTime(); Map 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 storageNode(StorageOperationBlockConfiguration configuration, Map 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 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 values(List inputs, Map executionVariables) { Map byName = new LinkedHashMap<>(); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java index 83db2b5..1e33e38 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Input.java @@ -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 diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java index 0796751..543439e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/Step.java @@ -170,11 +170,6 @@ public class Step implements InputListener { @JsonIgnore private final Set 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 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 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 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 implements InputListener { }); } - private void readFromStorage() { - if (this.storageHooks != null) { - this.storageHooks.read(this); - } - } - - private void writeToStorage(Map 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 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 runtimeVariables() { - Map 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 executorVariables() { + return new StepVariables(this.executionVariables, this.iteration); } private void settle(NodeExecutionResult executionResult) { @@ -378,7 +322,7 @@ public class Step 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 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 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); diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/StepStorageHooks.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/StepStorageHooks.java deleted file mode 100644 index 085515a..0000000 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/StepStorageHooks.java +++ /dev/null @@ -1,26 +0,0 @@ -// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - 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 outputs); - - static boolean isSynthetic(String inputName) { - return inputName != null && inputName.startsWith(SYNTHETIC_PREFIX); - } -} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/steps/StepVariables.java b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/StepVariables.java new file mode 100644 index 0000000..f75cd48 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/steps/StepVariables.java @@ -0,0 +1,57 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - 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}. + * + *

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 { + + private final Map shared; + private final int iteration; + + StepVariables(Map 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> entrySet() { + Set> entries = new LinkedHashSet<>(shared.entrySet()); + entries.add(new SimpleImmutableEntry<>(ExecutionRuntimeContextSupport.ITERATION, iteration)); + return entries; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/model/capabilities/NodeTypeCapabilities.java b/src/main/java/it/cnr/isti/workflow/manager/flows/model/capabilities/NodeTypeCapabilities.java index 77b3cd1..11f1c20 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/model/capabilities/NodeTypeCapabilities.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/model/capabilities/NodeTypeCapabilities.java @@ -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); } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java index 8a6587d..22cf730 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/flows/validation/FlowExecutionValidator.java @@ -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> 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. * *

    - *
  • They live in the top-level flow. - *
  • 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. - *
  • 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. + *
  • Storage nodes live in the top-level flow, and so do the operations linked to one. + *
  • An operation is linked to a node that is there, and what it does suits that node's storage + * and what the storage allows. + *
  • 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. *
*/ private List validateStorage(FlowData flowData, boolean rootFlow) { StorageFlows.Compiled storage = StorageFlows.compile(flowData); List> operations = (flowData.getBlocks() == null ? List.>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 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 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> graph = buildOutgoingGraph(flowData); for (Block node : storage.storageNodes().values()) { Set writers = new LinkedHashSet<>(); - storage.writes().stream().filter(write -> write.storageId().equals(node.getId())).forEach(write -> writers.add(write.sourceId())); Set 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 storageEdgeTypes(StorageFlows.Compiled storage, Map nodes, FlowData flowData) { - List errors = new ArrayList<>(); - for (Connection connection : flowData.getConnections() == null ? List.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 storageOperationReferences(StorageFlows.Compiled storage, List> operations) { List 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 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 unfedViewParameters(StorageFlows.Compiled storage) { - List errors = new ArrayList<>(); - Set 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 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(); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageFlows.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageFlows.java index 87eb3b4..028436c 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageFlows.java +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageFlows.java @@ -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. * - *

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: - * - *

    - *
  • a connection into one of its targets is a {@link Write} the source step does when it finishes; - *
  • a connection out of one of its views is a {@link Read} the consumer does before it starts; - *
  • 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}). - *
- * - *

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. + *

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> storageNodes, List writes, List reads, - List parameterFeeds, Map> syntheticInputs, List strayConnections) { + /** @param executionFlow the flow without its storage nodes; @param storageNodes by id */ + public record Compiled(FlowData executionFlow, Map> storageNodes) { public boolean hasStorage() { return !storageNodes.isEmpty(); } - - public List writesFrom(String stepId) { - return writes.stream().filter(write -> write.sourceId().equals(stepId)).toList(); - } - - public List 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> storageNodes = new LinkedHashMap<>(); flow.getBlocks().stream().filter(StorageFlows::isStorageNode).forEach(block -> storageNodes.put(block.getId(), block)); - - List kept = new ArrayList<>(); - List writes = new ArrayList<>(); - List reads = new ArrayList<>(); - List feeds = new ArrayList<>(); - List stray = new ArrayList<>(); - for (Connection connection : flow.getConnections() == null ? List.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> synthetic = new LinkedHashMap<>(); - Set 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 connections = (flow.getConnections() == null ? List.of() : flow.getConnections()).stream() + .filter(connection -> connection != null && !storageNodes.containsKey(connection.getSourceId()) + && !storageNodes.containsKey(connection.getTargetId())) + .toList(); List dependencies = (flow.getDependencies() == null ? List.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)); } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageNodeRuntime.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageNodeRuntime.java index e37d67e..42daa0e 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageNodeRuntime.java +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageNodeRuntime.java @@ -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 values, - Map authorizations, Map 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 authorizations, Map executionVariables) { - StorageNodeConfiguration configuration = configuration(node); - StorageType type = types.require(configuration.getStorageType()); - Function 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 authorizations, Map executionVariables) { StorageNodeConfiguration configuration = configuration(node); diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/validation/StorageNodeValidator.java b/src/main/java/it/cnr/isti/workflow/manager/storage/validation/StorageNodeValidator.java index e972baf..6d40ea5 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/storage/validation/StorageNodeValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/validation/StorageNodeValidator.java @@ -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 { @@ -67,7 +60,6 @@ public class StorageNodeValidator implements ConstraintValidator problems.add(new String[] { "initScript", problem })); - - Set 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("$", "\\$"); } diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/validation/StorageOperationValidator.java b/src/main/java/it/cnr/isti/workflow/manager/storage/validation/StorageOperationValidator.java index ade7875..63d6094 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/storage/validation/StorageOperationValidator.java +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/validation/StorageOperationValidator.java @@ -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. - * - *

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 { - @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 problems = problems(configuration); + List 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 problems(StorageOperationBlockConfiguration configuration) { - List 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 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 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 listProblems(String template) { - List 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("$", "\\$"); } diff --git a/src/test/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageSchemaTest.java b/src/test/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageSchemaTest.java new file mode 100644 index 0000000..895a910 --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageSchemaTest.java @@ -0,0 +1,43 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - 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")); + } +} diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/StorageFuzzTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/StorageFuzzTest.java index a224cd9..4f3d12f 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/storage/StorageFuzzTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/storage/StorageFuzzTest.java @@ -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> 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 open = target.getInputs().stream().filter(port -> !port.getName().equals("seen")).toList(); - if (open.isEmpty()) { + List 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> producers = steps.stream().filter(step -> !step.getOutputs().isEmpty()).toList(); - List> readers = steps.stream().filter(step -> step.getInputs().stream() - .anyMatch(port -> port.getName().equals("seen") || port.getName().equals(EndBlockFactory.INPUT_NAME))).toList(); - List> 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> 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> writers = steps.stream().filter(step -> step.getName().endsWith("-write")).toList(); + List> 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(); } diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/StorageNodeExecutionTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/StorageNodeExecutionTest.java index f273d45..80da4ce 100644 --- a/src/test/java/it/cnr/isti/workflow/manager/storage/StorageNodeExecutionTest.java +++ b/src/test/java/it/cnr/isti/workflow/manager/storage/StorageNodeExecutionTest.java @@ -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 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", "<> 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 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 errors = validator.collectErrors(FlowData.builder().block(readOnly).block(quoted).build()); + List 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 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 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 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 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 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 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 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; } } diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/StorageOperationExecutionTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/StorageOperationExecutionTest.java deleted file mode 100644 index 0c047f1..0000000 --- a/src/test/java/it/cnr/isti/workflow/manager/storage/StorageOperationExecutionTest.java +++ /dev/null @@ -1,256 +0,0 @@ -// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - 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 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 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 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 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 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 intoArchive = block("write-archive", StorageOperationBlockConfiguration.builder() - .storageType(S3StorageType.NAME).instance("archive").operation(StorageOperation.WRITE).template("x.txt")); - Block quoted = block("quoted", StorageOperationBlockConfiguration.builder() - .storageType(PostgresStorageType.NAME).instance("notes-db").operation(StorageOperation.READ) - .template("SELECT * FROM rounds WHERE verdict = '${{verdict}}'")); - Block wrongType = block("wrong-type", StorageOperationBlockConfiguration.builder() - .storageType(PostgresStorageType.NAME).instance("flow-files").operation(StorageOperation.LIST)); - Block scoped = block("scoped-db", StorageOperationBlockConfiguration.builder() - .storageType(PostgresStorageType.NAME).instance("notes-db").operation(StorageOperation.LIST) - .scope(StorageOperationBlockConfiguration.SCOPE_PER_EXECUTION)); - - List 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 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 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; - } -}