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 new file mode 100644 index 0000000..63036ec --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/configurations/StorageOperationBlockConfiguration.java @@ -0,0 +1,133 @@ +// 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.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.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.UiVisibleWhen; +import it.cnr.isti.workflow.manager.storage.StorageOperation; +import it.cnr.isti.workflow.manager.storage.ViewShape; +import it.cnr.isti.workflow.manager.storage.validation.ValidStorageOperation; +import jakarta.validation.constraints.NotBlank; +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.NonNull; + +/** + * One operation on one storage of the catalog. + * + *

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. + */ +@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED) +@Getter +@EqualsAndHashCode(callSuper = true) +@ValidStorageOperation +public class StorageOperationBlockConfiguration extends BlockConfiguration { + + 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") + @JsonProperty(required = true) + String storageType; + + @NotBlank + @UiOrder(20) + @UiLabel("Storage") + @UiDescription("A storage of the catalog, set up by whoever runs this service.") + @FieldRetriever(name = "Storage", url = "/retriever/Storage/instances", dependsOn = { "storageType" }) + @JsonProperty(required = true) + String instance; + + @Structural + @UiOrder(30) + @UiLabel("Operation") + @SchemaAllowedValues(value = { "READ", "WRITE", "LIST", "DELETE" }, defaultValue = "READ") + @JsonProperty(required = false) + StorageOperation operation = StorageOperation.READ; + + @Structural + @UiOrder(40) + @UiLabel("Key, pattern or statement") + @UiDescription("Object store: a key such as reports/${{context.iteration}}.json, or for READ and LIST a glob such as " + + "reports/*.json. Database: a SELECT for READ, an INSERT, UPDATE or MERGE for WRITE, a DELETE for DELETE, " + + "a schema (or nothing) for LIST. Other ${{name}} placeholders become inputs; ${{value.field}} is a field of what is written.") + @LongText(placeholder = "reports/*.json or SELECT * FROM notes WHERE round = ${{round}}", + tip = "Values are bound, never pasted into a statement: write ${{name}} without quotes.", + acceptVariableAsPlaceholder = true) + @JsonProperty(required = false) + String template; + + @Structural + @UiOrder(50) + @UiLabel("Read as") + @UiVisibleWhen(field = "operation", equals = "READ") + @UiDescription("Files, texts or JSON for an object store; a database always returns its rows as JSON.") + @SchemaAllowedValues(value = { "FILES", "TEXTS", "JSON", "KEYS" }, defaultValue = "FILES") + @JsonProperty(required = false) + ViewShape as; + + @UiOrder(60) + @UiLabel("Content type") + @UiVisibleWhen(field = "operation", equals = "WRITE") + @UiDescription("For an object store; guessed from the value when empty.") + @JsonProperty(required = false) + String contentType; + + @UiOrder(70) + @UiLabel("Scope") + @UiDescription("Shared: the same data for every execution. Per execution: an object store prefix of this execution's own.") + @SchemaAllowedValues(value = { SCOPE_SHARED, SCOPE_PER_EXECUTION }, defaultValue = SCOPE_SHARED) + @JsonProperty(required = false) + String scope = SCOPE_SHARED; + + @Builder + public StorageOperationBlockConfiguration(@NonNull String name, String storageType, String instance, + StorageOperation operation, String template, ViewShape as, String contentType, String scope) { + super(name); + this.storageType = storageType; + this.instance = instance; + this.operation = operation == null ? StorageOperation.READ : operation; + this.template = template; + this.as = as; + this.contentType = contentType; + this.scope = scope == null || scope.isBlank() ? SCOPE_SHARED : scope; + } + + @Override + public Class getBlockType() { + return StorageOperationBlockType.class; + } + + public static StorageOperationBlockConfiguration empty() { + StorageOperationBlockConfiguration configuration = new StorageOperationBlockConfiguration(); + configuration.name = StorageOperationBlockType.TYPE; + return configuration; + } + + public StorageOperation effectiveOperation() { + return operation == null ? StorageOperation.READ : operation; + } + + public boolean perExecution() { + return SCOPE_PER_EXECUTION.equals(scope); + } +} 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 new file mode 100644 index 0000000..fb3fa91 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/factories/StorageOperationBlockFactory.java @@ -0,0 +1,117 @@ +// 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.factories; + +import java.util.ArrayList; +import java.util.List; + +import org.springframework.stereotype.Component; + +import it.cnr.isti.workflow.manager.blocks.Block; +import it.cnr.isti.workflow.manager.blocks.IOCapability; +import it.cnr.isti.workflow.manager.blocks.IOCapabilityType; +import it.cnr.isti.workflow.manager.blocks.IOCapabilityTypes; +import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.types.StorageOperationBlockType; +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. + * + *

+ * + *

Built from a draft as much as from a finished block, so it never refuses: what is wrong is said + * by the validation, on the field it is about. + */ +@Component +public class StorageOperationBlockFactory implements BlockFactory { + + public static final String CONTENT_INPUT = "content"; + public static final String RESULT_OUTPUT = "result"; + + 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 StorageOperationBlockType blockType; + private final StorageTypes types; + + public StorageOperationBlockFactory(StorageOperationBlockType blockType, StorageTypes types) { + 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())) { + if (!parameter.equals(CONTENT_INPUT)) { + inputs.add(IODescriptor.input(parameter, IOType.TEXT, false, PARAMETER_CAPABILITIES)); + } + } + if (operation == StorageOperation.WRITE) { + inputs.add(IODescriptor.input(CONTENT_INPUT, IOType.ANY, false, type == null ? DEFAULT_WRITABLE : type.writableKinds())); + } + return Block.builder() + .inputs(inputs) + .output(result(operation, configuration, type)) + .specificConfiguration(configuration) + .type(blockType) + .build(); + } + + private static IODescriptor result(StorageOperation operation, StorageOperationBlockConfiguration configuration, + StorageType type) { + return switch (operation) { + case READ -> { + ViewShape shape = type == null ? (configuration.getAs() == null ? ViewShape.FILES : configuration.getAs()) + : StorageOperationValidator.shape(configuration, type); + yield IODescriptor.output(RESULT_OUTPUT, shape.portType(), shape.multiple(), + List.of(new IOCapability(IOCapabilityTypes.from(shape.portType()), shape.multiple()))); + } + case LIST -> IODescriptor.output(RESULT_OUTPUT, IOType.TEXT, true, + List.of(new IOCapability(IOCapabilityType.TEXT, true))); + case WRITE, DELETE -> IODescriptor.output(RESULT_OUTPUT, IOType.JSON, false, + List.of(new IOCapability(IOCapabilityType.JSON, false))); + }; + } + + @Override + public Block createEmpty() { + return create(StorageOperationBlockConfiguration.empty()); + } + + @Override + public Class getBlockType() { + return StorageOperationBlockType.class; + } + + @Override + public List supportedInputCapabilities() { + return DEFAULT_WRITABLE; + } + + @Override + public List supportedOutputCapabilities() { + return List.of(new IOCapability(IOCapabilityType.FILE, true), new IOCapability(IOCapabilityType.TEXT, true), + new IOCapability(IOCapabilityType.JSON, false)); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/blocks/types/StorageOperationBlockType.java b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/StorageOperationBlockType.java new file mode 100644 index 0000000..9cf27ec --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/blocks/types/StorageOperationBlockType.java @@ -0,0 +1,42 @@ +// 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.types; + +import org.springframework.stereotype.Component; + +import it.cnr.isti.workflow.manager.blocks.configurations.BlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration; + +@Component(StorageOperationBlockType.TYPE) +public class StorageOperationBlockType implements BlockType { + + public static final String TYPE = "StorageOperation"; + + @Override + public String getName() { + return TYPE; + } + + @Override + public String getDescription() { + return "Reads, writes, lists or deletes in a storage (an S3-compatible object store or a PostgreSQL " + + "database) chosen from the catalog: a key or glob for an object store, a SQL statement for a database."; + } + + @Override + public boolean validate() { + return true; + } + + @Override + public boolean isUserInteractive() { + return false; + } + + @Override + public Class> getBlockConfigurationClass() { + return StorageOperationBlockConfiguration.class; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java index d31ee26..2d3badc 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/ExecutionEventType.java @@ -41,6 +41,7 @@ public enum ExecutionEventType { LLM_REQUEST, LLM_ATTACHMENTS_SENT, HTTP_REQUEST, + STORAGE_OPERATION, MCP_SESSION_OPENED, MCP_SESSION_REUSED, MCP_SESSION_CLOSED, diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasMockedSideEffect.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasMockedSideEffect.java index 453515f..758b394 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasMockedSideEffect.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/BiasMockedSideEffect.java @@ -8,6 +8,7 @@ import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentChatBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration; public record BiasMockedSideEffect( String nodeId, @@ -20,7 +21,8 @@ public record BiasMockedSideEffect( ? "HTTP" : configuration instanceof MCPAgentChatBlockConfiguration ? "MCP_AGENT_CHAT" - : configuration instanceof MCPAgentBlockConfiguration ? "MCP_AGENT" : "EXTERNAL"; + : configuration instanceof MCPAgentBlockConfiguration ? "MCP_AGENT" + : configuration instanceof StorageOperationBlockConfiguration ? "STORAGE" : "EXTERNAL"; return new BiasMockedSideEffect(block.getId(), block.getName(), kind); } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapterRegistry.java b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapterRegistry.java index 57afc3d..b973df3 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapterRegistry.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/bias/runtime/BiasBehaviorAdapterRegistry.java @@ -18,6 +18,7 @@ import it.cnr.isti.workflow.manager.blocks.Block; import it.cnr.isti.workflow.manager.blocks.configurations.HTTPServerCallBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentBlockConfiguration; import it.cnr.isti.workflow.manager.blocks.configurations.MCPAgentChatBlockConfiguration; +import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration; import it.cnr.isti.workflow.manager.executions.ExecutionEventLogger; import it.cnr.isti.workflow.manager.executions.ExecutionEventType; import it.cnr.isti.workflow.manager.executions.bias.BiasExecutionContext; @@ -199,7 +200,8 @@ public class BiasBehaviorAdapterRegistry { Object configuration = block == null ? null : block.getSpecificConfiguration(); return configuration instanceof HTTPServerCallBlockConfiguration || configuration instanceof MCPAgentBlockConfiguration - || configuration instanceof MCPAgentChatBlockConfiguration; + || configuration instanceof MCPAgentChatBlockConfiguration + || configuration instanceof StorageOperationBlockConfiguration; } private static String entityType(FlowNode node) { 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 new file mode 100644 index 0000000..01af817 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/executors/blocks/StorageOperationExecutor.java @@ -0,0 +1,162 @@ +// 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.executors.blocks; + +import java.util.LinkedHashMap; +import java.util.List; +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.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.ExecutionEventLogger; +import it.cnr.isti.workflow.manager.executions.ExecutionEventType; +import it.cnr.isti.workflow.manager.executions.ExecutionRuntimeContextSupport; +import it.cnr.isti.workflow.manager.executions.ExecutionVariableDescriptor; +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.StorageInstancesProvider; +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; + +/** + * Runs one storage operation, in a session opened for it and closed after. + * + *

The event it leaves says which storage, which operation, where and how much - never what was + * read or written, which may be anything the flow handles. + */ +@Component +public class StorageOperationExecutor implements BlockExecutor { + + private final StorageTypes types; + private final StorageInstancesProvider instances; + private final ObjectMapper objectMapper; + + public StorageOperationExecutor(StorageTypes types, StorageInstancesProvider instances, ObjectMapper objectMapper) { + this.types = types; + this.instances = instances; + this.objectMapper = objectMapper; + } + + @Override + public Map execute(Block block, List inputs, + Map authorizations, Map executionVariables, + Map executionVariableDescriptors, ExecutionEventLogger eventLogger) { + StorageOperationBlockConfiguration configuration = (StorageOperationBlockConfiguration) block.getSpecificConfiguration(); + StorageType type = types.require(configuration.getStorageType()); + StorageInstancesProvider.StorageInstance instance = instances.find(configuration.getInstance()) + .orElseThrow(() -> new StorageException(StorageException.CONNECTION_INVALID, + "There is no storage " + configuration.getInstance() + " in the catalog")); + if (!instance.type().equalsIgnoreCase(type.getName())) { + throw new StorageException(StorageException.CONNECTION_INVALID, + instance.displayName() + " is a " + instance.type() + " storage, not " + type.getName()); + } + StorageConnection connection = instance.connection(type); + StorageOperation operation = configuration.effectiveOperation(); + Function values = values(inputs, executionVariables); + long started = System.nanoTime(); + Map details = new LinkedHashMap<>(); + details.put("storage", instance.id()); + details.put("type", type.getName()); + details.put("operation", operation.name()); + Object result; + try (StorageSession session = type.open(connection, scope(configuration, executionVariables))) { + result = switch (operation) { + case READ -> { + StorageReadResult read = session.read(configuration.getTemplate(), + StorageOperationValidator.shape(configuration, type), values); + details.put("count", read.count()); + details.put("truncated", read.truncated()); + yield read.value(); + } + case LIST -> { + StorageReadResult listed = session.list(configuration.getTemplate(), values); + details.put("count", listed.count()); + details.put("truncated", listed.truncated()); + yield listed.value(); + } + case WRITE -> written(session.write(configuration.getTemplate(), + StoragePlaceholders.require(values, StorageOperationBlockFactory.CONTENT_INPUT), + configuration.getContentType(), + StoragePlaceholders.withValue(values.apply(StorageOperationBlockFactory.CONTENT_INPUT), values)), details); + case DELETE -> written(session.delete(configuration.getTemplate(), values), details); + }; + } + details.put("durationMs", (System.nanoTime() - started) / 1_000_000); + if (eventLogger != null) { + eventLogger.info(ExecutionEventType.STORAGE_OPERATION, + operation + " on " + instance.displayName() + describe(details), details); + } + return Map.of(StorageOperationBlockFactory.RESULT_OUTPUT, result); + } + + private ObjectNode written(StorageWriteResult written, Map details) { + ObjectNode result = objectMapper.createObjectNode(); + result.put("affected", written.affected()); + if (written.location() != null) { + result.put("location", written.location()); + details.put("location", written.location()); + } + if (written.rows() != null) { + result.set("rows", written.rows()); + } + details.put("count", written.affected()); + return result; + } + + private static String describe(Map details) { + if (details.containsKey("location")) { + return ": " + details.get("location"); + } + return details.containsKey("count") ? ": " + details.get("count") + (Boolean.TRUE.equals(details.get("truncated")) ? "+ (truncated)" : "") : ""; + } + + private static StorageScope scope(StorageOperationBlockConfiguration configuration, Map 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<>(); + if (inputs != null) { + for (Input input : inputs) { + byName.put(input.getDescriptor().getName(), input.getValue()); + } + } + return name -> { + if (byName.containsKey(name)) { + return byName.get(name); + } + return executionVariables == null ? null : executionVariables.get(name); + }; + } + + @Override + public Class getBlockType() { + return StorageOperationBlockType.class; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageType.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageType.java index 85c0a40..c58c205 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageType.java +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageType.java @@ -48,6 +48,14 @@ public interface StorageType { List validateInit(StorageInit init); + /** + * Whether an execution can have a part of this storage to itself, under a prefix of its own. An + * object store can; a database's tables are the database's, and are shared. + */ + default boolean supportsExecutionScope() { + return false; + } + /** The placeholders of a template or query that bind as values rather than render as text. */ default Set parameters(String templateOrQuery) { return StoragePlaceholders.parameters(templateOrQuery); diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageType.java b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageType.java index c30de23..8da7d25 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageType.java +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageType.java @@ -87,6 +87,11 @@ public class S3StorageType implements StorageType { return List.of(ViewShape.FILES, ViewShape.TEXTS, ViewShape.JSON, ViewShape.KEYS); } + @Override + public boolean supportsExecutionScope() { + return true; + } + @Override public List writableKinds() { return List.of(new IOCapability(IOCapabilityType.FILE, false), new IOCapability(IOCapabilityType.TEXT, false), 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 new file mode 100644 index 0000000..29e847b --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/validation/StorageOperationValidator.java @@ -0,0 +1,123 @@ +// 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.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. + */ +@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) { + return true; + } + List problems = problems(configuration); + if (problems.isEmpty()) { + return true; + } + context.disableDefaultConstraintViolation(); + for (Problem problem : problems) { + context.buildConstraintViolationWithTemplate(escape(problem.message())) + .addPropertyNode(problem.field()) + .addConstraintViolation(); + } + return false; + } + + private record Problem(String field, String message) { + } + + List problems(StorageOperationBlockConfiguration configuration) { + List problems = new ArrayList<>(); + if (isBlank(configuration.getStorageType())) { + 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 (!isBlank(configuration.getInstance())) { + 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/main/java/it/cnr/isti/workflow/manager/storage/validation/ValidStorageOperation.java b/src/main/java/it/cnr/isti/workflow/manager/storage/validation/ValidStorageOperation.java new file mode 100644 index 0000000..ad8e7ac --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/validation/ValidStorageOperation.java @@ -0,0 +1,26 @@ +// 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.validation; + +import java.lang.annotation.ElementType; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +import jakarta.validation.Constraint; +import jakarta.validation.Payload; + +/** A storage operation its storage type and catalog instance can actually carry out. */ +@Target(ElementType.TYPE) +@Retention(RetentionPolicy.RUNTIME) +@Constraint(validatedBy = StorageOperationValidator.class) +public @interface ValidStorageOperation { + + String message() default "invalid storage operation"; + + Class[] groups() default {}; + + Class[] payload() default {}; +} 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 new file mode 100644 index 0000000..5fd988d --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/storage/StorageOperationExecutionTest.java @@ -0,0 +1,199 @@ +// 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 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); + } + + @Autowired + ExecutionsService executionsService; + + @Autowired + StorageOperationBlockFactory factory; + + @Autowired + FlowExecutionValidator validator; + + @Autowired + StorageTypes types; + + @Autowired + StorageInstancesProvider instances; + + @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 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; + } +}