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; - } -}