Add the StorageOperation block

One operation - READ, WRITE, LIST or DELETE - on one storage of the
catalog. What its template holds is the storage type's business: an
object key or glob for S3, a SQL statement for PostgreSQL. The block
stays the same for any type added later.

Ports follow from the operation: every ${{name}} of the template that is
not the execution's own becomes an input, a WRITE takes 'content' of the
kinds the storage accepts, and 'result' is what came back - files, texts,
JSON rows, names, or where a write went and how much it touched.

The storage type and the catalog check the block through the flow's own
validator, each problem on the field it is about: a read-only storage
refusing a WRITE, a placeholder inside a SQL string literal, an instance
of another type. Writes and deletes are mocked in bias experiments, like
HTTP calls. The STORAGE_OPERATION event says which storage, what and how
much, never the values.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Lucio Lelii 2026-09-28 11:46:46 +02:00
parent c472c80674
commit 02e8b0f751
12 changed files with 822 additions and 2 deletions

View File

@ -0,0 +1,133 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.blocks.configurations;
import com.fasterxml.jackson.annotation.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.
*
* <p>What {@code template} holds depends on the storage, not on this block: an object key or glob
* for an object store, a SQL statement for a database. The storage type checks it - see
* {@link ValidStorageOperation} - so this block stays the same whichever kind of storage is added.
*/
@NoArgsConstructor(access = lombok.AccessLevel.PROTECTED)
@Getter
@EqualsAndHashCode(callSuper = true)
@ValidStorageOperation
public class StorageOperationBlockConfiguration extends BlockConfiguration<StorageOperationBlockType> {
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<StorageOperationBlockType> 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);
}
}

View File

@ -0,0 +1,117 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.blocks.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.
*
* <ul>
* <li>Every {@code ${{name}}} of the template that is not the execution's own ({@code context.*},
* {@code global.*}) nor the written value is an input.
* <li>A WRITE takes {@code content}, of the kinds the storage accepts.
* <li>{@code result} is what came back: the objects or rows read, the names listed, or - for a write
* or a delete - where it went and how much it touched.
* </ul>
*
* <p>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<StorageOperationBlockType, StorageOperationBlockConfiguration> {
public static final String CONTENT_INPUT = "content";
public static final String RESULT_OUTPUT = "result";
private static final List<IOCapability> PARAMETER_CAPABILITIES = List.of(
new IOCapability(IOCapabilityType.TEXT, false), new IOCapability(IOCapabilityType.JSON, false));
private static final List<IOCapability> DEFAULT_WRITABLE = List.of(new IOCapability(IOCapabilityType.FILE, false),
new IOCapability(IOCapabilityType.TEXT, false), new IOCapability(IOCapabilityType.JSON, false));
private final StorageOperationBlockType blockType;
private final StorageTypes types;
public StorageOperationBlockFactory(StorageOperationBlockType blockType, StorageTypes types) {
this.blockType = blockType;
this.types = types;
}
@Override
public Block<StorageOperationBlockType> create(StorageOperationBlockConfiguration configuration) {
StorageType type = types.find(configuration.getStorageType()).orElse(null);
StorageOperation operation = configuration.effectiveOperation();
List<IODescriptor> inputs = new ArrayList<>();
for (String parameter : StoragePlaceholders.parameters(configuration.getTemplate())) {
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.<StorageOperationBlockType>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<StorageOperationBlockType> createEmpty() {
return create(StorageOperationBlockConfiguration.empty());
}
@Override
public Class<StorageOperationBlockType> getBlockType() {
return StorageOperationBlockType.class;
}
@Override
public List<IOCapability> supportedInputCapabilities() {
return DEFAULT_WRITABLE;
}
@Override
public List<IOCapability> supportedOutputCapabilities() {
return List.of(new IOCapability(IOCapabilityType.FILE, true), new IOCapability(IOCapabilityType.TEXT, true),
new IOCapability(IOCapabilityType.JSON, false));
}
}

View File

@ -0,0 +1,42 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.blocks.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<? extends BlockConfiguration<?>> getBlockConfigurationClass() {
return StorageOperationBlockConfiguration.class;
}
}

View File

@ -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,

View File

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

View File

@ -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) {

View File

@ -0,0 +1,162 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.executions.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.
*
* <p>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<StorageOperationBlockType> {
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<String, Object> execute(Block<StorageOperationBlockType> block, List<Input> inputs,
Map<String, Object> authorizations, Map<String, Object> executionVariables,
Map<String, ExecutionVariableDescriptor> 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<String, Object> values = values(inputs, executionVariables);
long started = System.nanoTime();
Map<String, Object> 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<String, Object> 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<String, Object> 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<String, Object> executionVariables) {
if (!configuration.perExecution()) {
return StorageScope.shared();
}
Object root = executionVariables == null ? null : executionVariables.get(ExecutionRuntimeContextSupport.ROOT_EXECUTION_ID);
if (root == null || String.valueOf(root).isBlank()) {
throw new StorageException(StorageException.CONNECTION_INVALID,
"A per-execution scope needs the execution's id, and this step has none");
}
return StorageScope.perExecution(String.valueOf(root));
}
/** An input by its port name; the execution's own names from its variables. */
private static Function<String, Object> values(List<Input> inputs, Map<String, Object> executionVariables) {
Map<String, Object> byName = new LinkedHashMap<>();
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<StorageOperationBlockType> getBlockType() {
return StorageOperationBlockType.class;
}
}

View File

@ -48,6 +48,14 @@ public interface StorageType {
List<String> 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<String> parameters(String templateOrQuery) {
return StoragePlaceholders.parameters(templateOrQuery);

View File

@ -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<IOCapability> writableKinds() {
return List.of(new IOCapability(IOCapabilityType.FILE, false), new IOCapability(IOCapabilityType.TEXT, false),

View File

@ -0,0 +1,123 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage.validation;
import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration;
import it.cnr.isti.workflow.manager.storage.StorageInstancesProvider;
import it.cnr.isti.workflow.manager.storage.StorageOperation;
import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
import it.cnr.isti.workflow.manager.storage.StorageType;
import it.cnr.isti.workflow.manager.storage.StorageTypes;
import it.cnr.isti.workflow.manager.storage.ViewShape;
import jakarta.validation.ConstraintValidator;
import jakarta.validation.ConstraintValidatorContext;
/**
* Asks the storage type whether it can do what the block says, and the catalog whether the storage
* exists and allows it - each problem reported on the field it is about.
*
* <p>Needs the storage beans, so it only judges under Spring's validator. A plain one, with nothing
* to ask, lets the configuration through: the block is checked again when it runs.
*/
@Component
public class StorageOperationValidator implements ConstraintValidator<ValidStorageOperation, StorageOperationBlockConfiguration> {
@Autowired(required = false)
StorageTypes types;
@Autowired(required = false)
StorageInstancesProvider instances;
@Override
public boolean isValid(StorageOperationBlockConfiguration configuration, ConstraintValidatorContext context) {
if (configuration == null || types == null || instances == null) {
return true;
}
List<Problem> 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<Problem> problems(StorageOperationBlockConfiguration configuration) {
List<Problem> problems = new ArrayList<>();
if (isBlank(configuration.getStorageType())) {
return problems;
}
Optional<StorageType> found = types.find(configuration.getStorageType());
if (found.isEmpty()) {
problems.add(new Problem("storageType", "Unknown storage type " + configuration.getStorageType()
+ "; this service has " + types.all().stream().map(StorageType::getName).toList()));
return problems;
}
StorageType type = found.get();
StorageOperation operation = configuration.effectiveOperation();
if (!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<String> templateProblems = switch (operation) {
case READ -> type.validateQuery(template, shape(configuration, type));
case WRITE -> type.validateTemplate(template);
case DELETE -> type.validateDelete(template);
case LIST -> listProblems(template);
};
templateProblems.forEach(problem -> problems.add(new Problem("template", problem)));
if (operation != StorageOperation.WRITE && template != null
&& StoragePlaceholders.names(template).stream().anyMatch(StoragePlaceholders::isValueName)) {
problems.add(new Problem("template", "${{value}} is what a WRITE writes; a " + operation + " has none"));
}
if (configuration.perExecution() && !type.supportsExecutionScope()) {
problems.add(new Problem("scope", "A " + type.getName() + " storage is shared by every execution; "
+ "a per-execution scope is for object stores"));
}
return problems;
}
public static ViewShape shape(StorageOperationBlockConfiguration configuration, StorageType type) {
return configuration.getAs() != null ? configuration.getAs() : type.viewShapes().getFirst();
}
private static List<String> listProblems(String template) {
List<String> problems = new ArrayList<>();
StoragePlaceholders.invalidNames(template).forEach(name -> problems.add(
"${{" + name + "}} is not a usable name: start with a letter, then letters, digits, '-', '_' or '.'"));
return problems;
}
private static boolean isBlank(String value) {
return value == null || value.isBlank();
}
/** Bean Validation would read {..} and $ in a message as its own syntax; ours are literal. */
private static String escape(String message) {
return message.replace("\\", "\\\\").replace("{", "\\{").replace("}", "\\}").replace("$", "\\$");
}
}

View File

@ -0,0 +1,26 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage.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<? extends Payload>[] payload() default {};
}

View File

@ -0,0 +1,199 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
import java.nio.file.Files;
import java.nio.file.Path;
import java.time.Duration;
import java.time.Instant;
import java.util.List;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.test.context.DynamicPropertyRegistry;
import org.springframework.test.context.DynamicPropertySource;
import org.springframework.test.context.TestPropertySource;
import org.testcontainers.containers.MinIOContainer;
import org.testcontainers.junit.jupiter.Container;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.postgresql.PostgreSQLContainer;
import it.cnr.isti.workflow.manager.blocks.Block;
import it.cnr.isti.workflow.manager.blocks.configurations.StorageOperationBlockConfiguration;
import it.cnr.isti.workflow.manager.blocks.factories.StorageOperationBlockFactory;
import it.cnr.isti.workflow.manager.blocks.types.StorageOperationBlockType;
import it.cnr.isti.workflow.manager.executions.ExecutionEventType;
import it.cnr.isti.workflow.manager.executions.ExecutionObject;
import it.cnr.isti.workflow.manager.executions.FieldKey;
import it.cnr.isti.workflow.manager.executions.ExecutionStatus;
import it.cnr.isti.workflow.manager.executions.ExecutionsService;
import it.cnr.isti.workflow.manager.flows.model.Dependency;
import it.cnr.isti.workflow.manager.flows.model.FlowData;
import it.cnr.isti.workflow.manager.flows.validation.FlowExecutionValidator;
import it.cnr.isti.workflow.manager.flows.validation.ValidationError;
import it.cnr.isti.workflow.manager.storage.postgres.PostgresStorageType;
import it.cnr.isti.workflow.manager.storage.s3.S3StorageType;
import it.cnr.isti.workflow.manager.ios.IOType;
import 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<StorageOperationBlockType> write = block("save-note", StorageOperationBlockConfiguration.builder()
.storageType(S3StorageType.NAME).instance("flow-files").operation(StorageOperation.WRITE)
.template("notes/${{title}}.txt").scope(StorageOperationBlockConfiguration.SCOPE_PER_EXECUTION));
Block<StorageOperationBlockType> read = block("read-notes", StorageOperationBlockConfiguration.builder()
.storageType(S3StorageType.NAME).instance("flow-files").operation(StorageOperation.READ)
.template("notes/*.txt").as(ViewShape.TEXTS).scope(StorageOperationBlockConfiguration.SCOPE_PER_EXECUTION));
assertEquals(List.of("title", "content"), write.getInputs().stream().map(input -> input.getName()).toList());
assertEquals(IOType.TEXT, read.getOutputs().getFirst().getType());
assertTrue(read.getOutputs().getFirst().isMultiple());
FlowData flow = FlowData.builder().block(write).block(read)
.dependency(new Dependency(write.getId(), read.getId())).build();
ExecutionObject execution = executionsService.createExecution("Write then read", flow);
executionsService.prepareInput(execution.getId(), write.getId(), "title", "round-1");
executionsService.prepareInput(execution.getId(), write.getId(), "content", "The list forgets items on reload");
executionsService.startExecution(execution.getId());
ExecutionObject finished = awaitFinal(execution.getId());
assertEquals(ExecutionStatus.SUCCESS, finished.getContext().getStatus(), () -> finished.getContext().getErrors().toString());
assertEquals(List.of("The list forgets items on reload"), result(finished, read));
assertTrue(finished.getContext().getEvents().stream().anyMatch(event -> event.getType() == ExecutionEventType.STORAGE_OPERATION
&& event.getMessage().contains("notes/round-1.txt")));
}
@Test
void aSqlWriteBindsTheFieldsOfItsContentAndAReadReturnsTheRows() throws Exception {
try (StorageSession session = types.require(PostgresStorageType.NAME).open(
instances.find("notes-db").orElseThrow().connection(types.require(PostgresStorageType.NAME)), StorageScope.shared())) {
session.initialize(new StorageInit(false, "CREATE TABLE IF NOT EXISTS rounds (round int, verdict text)"));
}
Block<StorageOperationBlockType> write = block("record-round", StorageOperationBlockConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("notes-db").operation(StorageOperation.WRITE)
.template("INSERT INTO rounds(round, verdict) VALUES (${{value.round}}, ${{value.verdict}})"));
Block<StorageOperationBlockType> read = block("rejected-rounds", StorageOperationBlockConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("notes-db").operation(StorageOperation.READ)
.template("SELECT round FROM rounds WHERE verdict = ${{verdict}} ORDER BY round"));
assertEquals(List.of("verdict"), read.getInputs().stream().map(input -> input.getName()).toList());
assertEquals(IOType.JSON, read.getOutputs().getFirst().getType());
FlowData flow = FlowData.builder().block(write).block(read)
.dependency(new Dependency(write.getId(), read.getId())).build();
ExecutionObject execution = executionsService.createExecution("Record a round", flow);
executionsService.prepareInput(execution.getId(), write.getId(), "content",
new tools.jackson.databind.ObjectMapper().readTree("{\"round\": 7, \"verdict\": \"rejected\"}"));
executionsService.prepareInput(execution.getId(), read.getId(), "verdict", "rejected");
executionsService.startExecution(execution.getId());
ExecutionObject finished = awaitFinal(execution.getId());
assertEquals(ExecutionStatus.SUCCESS, finished.getContext().getStatus(), () -> finished.getContext().getErrors().toString());
JsonNode rows = (JsonNode) result(finished, read);
assertEquals(7, rows.get(0).get("round").asInt());
}
@Test
void theFlowIsRefusedWhereTheStorageCannotDoWhatTheBlockAsks() {
Block<StorageOperationBlockType> intoArchive = block("write-archive", StorageOperationBlockConfiguration.builder()
.storageType(S3StorageType.NAME).instance("archive").operation(StorageOperation.WRITE).template("x.txt"));
Block<StorageOperationBlockType> quoted = block("quoted", StorageOperationBlockConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("notes-db").operation(StorageOperation.READ)
.template("SELECT * FROM rounds WHERE verdict = '${{verdict}}'"));
Block<StorageOperationBlockType> wrongType = block("wrong-type", StorageOperationBlockConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("flow-files").operation(StorageOperation.LIST));
Block<StorageOperationBlockType> scoped = block("scoped-db", StorageOperationBlockConfiguration.builder()
.storageType(PostgresStorageType.NAME).instance("notes-db").operation(StorageOperation.LIST)
.scope(StorageOperationBlockConfiguration.SCOPE_PER_EXECUTION));
List<ValidationError> errors = validator.collectErrors(FlowData.builder()
.block(intoArchive).block(quoted).block(wrongType).block(scoped).build());
assertHasError(errors, intoArchive, "specificConfiguration.operation", "is read-only");
assertHasError(errors, quoted, "specificConfiguration.template", "${{verdict}} is inside a string literal");
assertHasError(errors, wrongType, "specificConfiguration.instance", "is a S3 storage, not PostgreSQL");
assertHasError(errors, scoped, "specificConfiguration.scope", "shared by every execution");
}
private static Object result(ExecutionObject execution, Block<?> block) {
return execution.getContext().getResult().get(new FieldKey(block.getId(), StorageOperationBlockFactory.RESULT_OUTPUT));
}
private static void assertHasError(List<ValidationError> errors, Block<?> block, String field, String text) {
assertTrue(errors.stream().anyMatch(error -> block.getId().equals(error.id()) && field.equals(error.field())
&& error.message().contains(text)), () -> "No error on " + field + " containing '" + text + "': " + errors);
}
private void prepareBucket() {
StorageType s3 = types.require(S3StorageType.NAME);
try (StorageSession session = s3.open(instances.find("flow-files").orElseThrow().connection(s3), StorageScope.shared())) {
session.initialize(new StorageInit(true, null));
}
}
private Block<StorageOperationBlockType> block(String name,
StorageOperationBlockConfiguration.StorageOperationBlockConfigurationBuilder builder) {
return factory.create(builder.name(name).build());
}
private ExecutionObject awaitFinal(String id) {
Instant deadline = Instant.now().plus(Duration.ofSeconds(20));
ExecutionObject execution = executionsService.getExecution(id);
while (!execution.getContext().getStatus().isFinalState() && Instant.now().isBefore(deadline)) {
Thread.onSpinWait();
execution = executionsService.getExecution(id);
}
return execution;
}
}