- 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 org.springframework.beans.factory.annotation.Value;
+import org.springframework.stereotype.Component;
+
+import it.cnr.isti.workflow.manager.llms.providers.OutboundEndpointGuard;
+
+/**
+ * Where a person's own storage connection may point.
+ *
+ * The same rules as an LLM endpoint - see {@link OutboundEndpointGuard} - behind a separate
+ * setting, {@code app.storage.endpoint.*}: an operator may well open one to a private network and
+ * not the other. Applied when the connection is saved and again right before every connection,
+ * which is the check that protects; a catalog instance was chosen by the operator and does not
+ * pass through here.
+ */
+@Component
+public class StorageEndpointGuard {
+
+ private final OutboundEndpointGuard guard;
+
+ public StorageEndpointGuard(
+ @Value("${app.storage.endpoint.allow-private-network:false}") boolean allowPrivateNetwork,
+ @Value("${app.storage.endpoint.allowed-host-patterns:}") String allowedHostPatterns) {
+ this.guard = new OutboundEndpointGuard(allowPrivateNetwork, allowedHostPatterns);
+ }
+
+ /** An http(s) endpoint, such as an object store's. */
+ public void validateUrl(String url) {
+ try {
+ guard.validate(url);
+ } catch (IllegalArgumentException ex) {
+ throw new StorageException(StorageException.CONNECTION_INVALID, ex.getMessage());
+ }
+ }
+
+ /** A bare host, such as a database's. */
+ public void validateHost(String host) {
+ try {
+ guard.validateHost(host);
+ } catch (IllegalArgumentException ex) {
+ throw new StorageException(StorageException.CONNECTION_INVALID, ex.getMessage());
+ }
+ }
+
+ /** Checks a person's own connection; lets an operator's through. */
+ public void check(StorageConnection connection, String urlKey, String hostKey) {
+ if (connection.operatorChosen()) {
+ return;
+ }
+ if (urlKey != null) {
+ validateUrl(connection.requireText(urlKey));
+ }
+ if (hostKey != null) {
+ validateHost(connection.requireText(hostKey));
+ }
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageException.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageException.java
new file mode 100644
index 0000000..e4cb0b6
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageException.java
@@ -0,0 +1,30 @@
+// 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 it.cnr.isti.workflow.manager.executions.NodeExecutionException;
+
+/** A storage refusing, or failing, what a step asked of it - always said in the step's terms. */
+public class StorageException extends NodeExecutionException {
+
+ public static final String CONNECTION_INVALID = "STORAGE_CONNECTION_INVALID";
+ public static final String UNREACHABLE = "STORAGE_UNREACHABLE";
+ public static final String NOT_FOUND = "STORAGE_OBJECT_NOT_FOUND";
+ public static final String TOO_LARGE = "STORAGE_OBJECT_TOO_LARGE";
+ public static final String OPERATION_NOT_ALLOWED = "STORAGE_OPERATION_NOT_ALLOWED";
+ public static final String PLACEHOLDER_MISSING = "STORAGE_PLACEHOLDER_MISSING";
+ public static final String INVALID_LOCATION = "STORAGE_INVALID_LOCATION";
+ public static final String SQL_REJECTED = "STORAGE_SQL_REJECTED";
+ public static final String INIT_FAILED = "STORAGE_INIT_FAILED";
+ public static final String CONTENT_INVALID = "STORAGE_CONTENT_INVALID";
+
+ public StorageException(String errorCode, String message) {
+ super(errorCode, message);
+ }
+
+ public StorageException(String errorCode, String message, Throwable cause) {
+ super(errorCode, message, cause);
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInit.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInit.java
new file mode 100644
index 0000000..3036999
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInit.java
@@ -0,0 +1,25 @@
+// 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;
+
+/**
+ * What a storage node prepares before the first step of an execution runs. Both are idempotent by
+ * contract, because every execution runs them again.
+ *
+ * @param createIfMissing create the container the connection names (an object store's bucket) when
+ * it does not exist yet
+ * @param script statements that prepare a database - {@code CREATE TABLE IF NOT EXISTS ...} - run
+ * in one transaction; blank for none
+ */
+public record StorageInit(boolean createIfMissing, String script) {
+
+ public static StorageInit none() {
+ return new StorageInit(false, null);
+ }
+
+ public boolean hasScript() {
+ return script != null && !script.isBlank();
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstanceView.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstanceView.java
new file mode 100644
index 0000000..59e979a
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstanceView.java
@@ -0,0 +1,15 @@
+// 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;
+
+/** A catalog storage as a flow author sees it: what it is and what it allows, not how to reach it. */
+public record StorageInstanceView(String id, String name, String description, String type, boolean readOnly,
+ boolean allowInit) {
+
+ public static StorageInstanceView of(StorageInstancesProvider.StorageInstance instance) {
+ return new StorageInstanceView(instance.id(), instance.displayName(), instance.description(), instance.type(),
+ instance.readOnly(), instance.allowInit());
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProvider.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProvider.java
new file mode 100644
index 0000000..9c5b8ce
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProvider.java
@@ -0,0 +1,146 @@
+// 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 java.io.InputStream;
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.Set;
+
+import org.springframework.beans.factory.annotation.Value;
+import org.springframework.core.io.ClassPathResource;
+import org.springframework.core.io.Resource;
+import org.springframework.core.io.ResourceLoader;
+import org.springframework.stereotype.Component;
+import org.springframework.util.StringUtils;
+
+import com.fasterxml.jackson.annotation.JsonCreator;
+import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
+import com.fasterxml.jackson.annotation.JsonProperty;
+import tools.jackson.databind.JsonNode;
+import tools.jackson.databind.ObjectMapper;
+
+/**
+ * The storages whoever runs this service offers to every flow, from {@code storages.json}.
+ *
+ * Loaded the way the MCP server catalog is: {@code app.storage.instances.file} names the file,
+ * and without it the (empty) one on the classpath is used. Every entry is checked by its own type
+ * as it loads, and a file that does not hold together stops the service from starting: an instance
+ * that looked fine in the editor and failed at every run would be found much later and by someone
+ * else.
+ *
+ *
The file holds the credentials of the instances it lists, like {@code mcp-servers.json} does,
+ * so it belongs outside the image, readable by the service alone.
+ */
+@Component
+public class StorageInstancesProvider {
+
+ private final List instances;
+
+ public StorageInstancesProvider(ObjectMapper objectMapper, ResourceLoader resourceLoader, StorageTypes types,
+ @Value("${app.storage.instances.file:}") String location) {
+ this.instances = load(Objects.requireNonNull(objectMapper), Objects.requireNonNull(resourceLoader),
+ Objects.requireNonNull(types), location);
+ }
+
+ public List all() {
+ return instances;
+ }
+
+ public List ofType(String type) {
+ return instances.stream().filter(instance -> instance.type().equalsIgnoreCase(type)).toList();
+ }
+
+ public Optional find(String id) {
+ return instances.stream().filter(instance -> instance.id().equals(id)).findFirst();
+ }
+
+ private static List load(ObjectMapper objectMapper, ResourceLoader resourceLoader,
+ StorageTypes types, String location) {
+ Catalog catalog;
+ try (InputStream input = resolve(resourceLoader, location).getInputStream()) {
+ catalog = objectMapper.readValue(input, Catalog.class);
+ } catch (Exception ex) {
+ throw new IllegalStateException("Unable to load the storage catalog from " + describe(location), ex);
+ }
+ List loaded = new ArrayList<>();
+ Set ids = new HashSet<>();
+ for (StorageInstance instance : catalog.storages() == null ? List.of() : catalog.storages()) {
+ String where = "Storage catalog " + describe(location) + ", entry " + (loaded.size() + 1);
+ if (instance == null || !StringUtils.hasText(instance.id())) {
+ throw new IllegalStateException(where + " has no id");
+ }
+ if (!ids.add(instance.id())) {
+ throw new IllegalStateException(where + ": id " + instance.id() + " is used twice");
+ }
+ StorageType type = types.find(instance.type()).orElseThrow(() -> new IllegalStateException(
+ where + " (" + instance.id() + ") has unknown type " + instance.type()
+ + "; known types: " + types.all().stream().map(StorageType::getName).toList()));
+ try {
+ type.validateConnection(instance.connection(type));
+ } catch (StorageException ex) {
+ throw new IllegalStateException(where + " (" + instance.id() + "): " + ex.getMessage(), ex);
+ }
+ loaded.add(instance);
+ }
+ return List.copyOf(loaded);
+ }
+
+ private static Resource resolve(ResourceLoader resourceLoader, String location) {
+ if (!StringUtils.hasText(location)) {
+ return new ClassPathResource("storages.json");
+ }
+ String resolved = location.startsWith("classpath:") || location.startsWith("file:") ? location : "file:" + location;
+ Resource resource = resourceLoader.getResource(resolved);
+ if (!resource.exists()) {
+ throw new IllegalStateException("Storage catalog not found at " + resolved);
+ }
+ return resource;
+ }
+
+ private static String describe(String location) {
+ return StringUtils.hasText(location) ? location : "classpath:storages.json";
+ }
+
+ @JsonIgnoreProperties(ignoreUnknown = true)
+ private record Catalog(List storages) {
+ }
+
+ /**
+ * One storage of the catalog.
+ *
+ * @param readOnly only reads: a node writing to it, or a block deleting from it, is refused
+ * while the flow is edited
+ * @param allowInit whether a flow may prepare it - create its bucket, run its DDL script
+ */
+ @JsonIgnoreProperties(ignoreUnknown = true)
+ public record StorageInstance(String id, String name, String description, String type, boolean readOnly,
+ boolean allowInit, JsonNode connection) {
+
+ @JsonCreator
+ static StorageInstance of(@JsonProperty("id") String id,
+ @JsonProperty("name") String name,
+ @JsonProperty("description") String description,
+ @JsonProperty("type") String type,
+ @JsonProperty("readOnly") Boolean readOnly,
+ @JsonProperty("allowInit") Boolean allowInit,
+ @JsonProperty("connection") JsonNode connection) {
+ // Left out means no: an instance is writable-but-not-preparable only when the file says so.
+ return new StorageInstance(id, name, description, type, Boolean.TRUE.equals(readOnly),
+ Boolean.TRUE.equals(allowInit), connection);
+ }
+
+ public String displayName() {
+ return StringUtils.hasText(name) ? name : id;
+ }
+
+ public StorageConnection connection(StorageType type) {
+ return new StorageConnection(displayName(), type.getName(), connection, readOnly, allowInit, true);
+ }
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageOperation.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageOperation.java
new file mode 100644
index 0000000..c3973a1
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageOperation.java
@@ -0,0 +1,13 @@
+// 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;
+
+/** What a step can ask of a storage, one operation at a time. */
+public enum StorageOperation {
+ READ,
+ WRITE,
+ LIST,
+ DELETE
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StoragePlaceholders.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StoragePlaceholders.java
new file mode 100644
index 0000000..9feed59
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StoragePlaceholders.java
@@ -0,0 +1,157 @@
+// 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 java.util.LinkedHashSet;
+import java.util.Set;
+import java.util.function.Function;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
+
+import it.cnr.isti.workflow.manager.executions.ExecutionTemplateResolver;
+
+/**
+ * The ${{name}} placeholders of a storage template - an object key, a glob, a statement.
+ *
+ * Names come in three kinds, and the storage node treats them differently: {@code context.*}
+ * and {@code global.*} are the execution's own and always known; {@code value} and
+ * {@code value.} are the value being written; anything else is a parameter the node takes
+ * from another step through a port.
+ */
+public final class StoragePlaceholders {
+
+ private static final Pattern PLACEHOLDER = Pattern.compile("\\$\\{\\{\\s*([^}]*?)\\s*}}");
+ private static final Pattern VALID_NAME = Pattern.compile("^[A-Za-z][A-Za-z0-9_.-]*$");
+
+ public static final String VALUE = "value";
+
+ private StoragePlaceholders() {
+ }
+
+ /** Every name, in the order it first appears. */
+ public static Set names(String template) {
+ Set names = new LinkedHashSet<>();
+ if (template == null) {
+ return names;
+ }
+ Matcher matcher = PLACEHOLDER.matcher(template);
+ while (matcher.find()) {
+ names.add(matcher.group(1));
+ }
+ return names;
+ }
+
+ /** The names another step has to supply: neither the execution's own nor the written value. */
+ public static Set parameters(String template) {
+ Set parameters = new LinkedHashSet<>();
+ for (String name : names(template)) {
+ if (!isExecutionName(name) && !isValueName(name)) {
+ parameters.add(name);
+ }
+ }
+ return parameters;
+ }
+
+ public static boolean isExecutionName(String name) {
+ return name.startsWith("context.") || name.startsWith("global.") || name.startsWith("project.");
+ }
+
+ public static boolean isValueName(String name) {
+ return name.equals(VALUE) || name.startsWith(VALUE + ".");
+ }
+
+ /** Names that would not survive as a port or a lookup, for the node's validation. */
+ public static Set invalidNames(String template) {
+ Set invalid = new LinkedHashSet<>();
+ for (String name : names(template)) {
+ if (!VALID_NAME.matcher(name).matches()) {
+ invalid.add(name);
+ }
+ }
+ return invalid;
+ }
+
+ /**
+ * The template with every placeholder replaced by its value as text - for an object key or a
+ * glob, never for a statement, whose values are bound instead.
+ */
+ public static String render(String template, Function values) {
+ if (template == null) {
+ return null;
+ }
+ Matcher matcher = PLACEHOLDER.matcher(template);
+ StringBuilder rendered = new StringBuilder();
+ while (matcher.find()) {
+ Object value = require(values, matcher.group(1));
+ matcher.appendReplacement(rendered, Matcher.quoteReplacement(ExecutionTemplateResolver.formatValue(value)));
+ }
+ matcher.appendTail(rendered);
+ return rendered.toString();
+ }
+
+ public static Object require(Function values, String name) {
+ Object value = values == null ? null : values.apply(name);
+ if (value == null) {
+ throw new StorageException(StorageException.PLACEHOLDER_MISSING,
+ "No value for ${{" + name + "}}");
+ }
+ return value;
+ }
+
+ /**
+ * Lookups for a write: {@code value} is the value itself, {@code value.a.b} walks into a JSON
+ * value (or a map) field by field, and every other name goes to {@code values}.
+ */
+ public static Function withValue(Object value, Function values) {
+ return name -> {
+ if (name.equals(VALUE)) {
+ return value;
+ }
+ if (name.startsWith(VALUE + ".")) {
+ return field(value, name.substring(VALUE.length() + 1));
+ }
+ return values == null ? null : values.apply(name);
+ };
+ }
+
+ private static Object field(Object value, String path) {
+ Object current = value;
+ for (String part : path.split("\\.")) {
+ if (current instanceof tools.jackson.databind.JsonNode node) {
+ tools.jackson.databind.JsonNode next = node.get(part);
+ if (next == null || next.isNull()) {
+ return null;
+ }
+ current = next.isValueNode() ? scalar(next) : next;
+ } else if (current instanceof java.util.Map, ?> map) {
+ current = map.get(part);
+ } else {
+ return null;
+ }
+ if (current == null) {
+ return null;
+ }
+ }
+ return current;
+ }
+
+ private static Object scalar(tools.jackson.databind.JsonNode node) {
+ if (node.isString()) {
+ return node.asString();
+ }
+ if (node.isBoolean()) {
+ return node.asBoolean();
+ }
+ if (node.isNumber()) {
+ return node.numberValue();
+ }
+ return node.toString();
+ }
+
+ /** The pattern, for a type that has to find placeholders itself - inside a statement, say. */
+ public static Pattern pattern() {
+ return PLACEHOLDER;
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageReadResult.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageReadResult.java
new file mode 100644
index 0000000..307dc8a
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageReadResult.java
@@ -0,0 +1,15 @@
+// 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;
+
+/**
+ * What a read found.
+ *
+ * @param value in the shape asked for: a list of files, a list of texts or names, or one JSON value
+ * @param count how many objects or rows it holds
+ * @param truncated true when more matched than a read may return, and the rest were left out
+ */
+public record StorageReadResult(Object value, int count, boolean truncated) {
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageScope.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageScope.java
new file mode 100644
index 0000000..5da6313
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageScope.java
@@ -0,0 +1,25 @@
+// 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;
+
+/**
+ * Where inside the connection an execution works. Shared means where the connection already
+ * points; per execution adds a prefix of that execution's own, so every run starts empty and none
+ * sees another's data.
+ */
+public record StorageScope(String prefix) {
+
+ public static StorageScope shared() {
+ return new StorageScope("");
+ }
+
+ public static StorageScope perExecution(String rootExecutionId) {
+ return new StorageScope(rootExecutionId + "/");
+ }
+
+ public StorageScope {
+ prefix = prefix == null ? "" : prefix;
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSession.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSession.java
new file mode 100644
index 0000000..335069d
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSession.java
@@ -0,0 +1,32 @@
+// 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 java.util.function.Function;
+
+/**
+ * One conversation with a storage, opened for an operation and closed after it.
+ *
+ * The same four things for every type, each read the way the type reads it: a {@code template}
+ * is an object key or a statement, a {@code query} is a glob or a {@code SELECT}. Placeholders are
+ * looked up through {@code values}, and a type decides whether it renders them into the text (a
+ * key) or binds them (a statement).
+ */
+public interface StorageSession extends AutoCloseable {
+
+ void initialize(StorageInit init);
+
+ StorageWriteResult write(String template, Object value, String contentType, Function values);
+
+ StorageReadResult read(String query, ViewShape shape, Function values);
+
+ /** What there is: the keys under a prefix or glob, or the tables of a database. */
+ StorageReadResult list(String query, Function values);
+
+ StorageWriteResult delete(String template, Function values);
+
+ @Override
+ void close();
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSettings.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSettings.java
new file mode 100644
index 0000000..12db891
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSettings.java
@@ -0,0 +1,56 @@
+// 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 org.springframework.beans.factory.annotation.Value;
+import org.springframework.stereotype.Component;
+
+/** The limits every storage works within, whoever's connection it is. */
+@Component
+public class StorageSettings {
+
+ private final long maxObjectBytes;
+ private final int maxKeys;
+ private final int sqlStatementTimeoutSeconds;
+ private final int sqlMaxRows;
+ private final int connectTimeoutSeconds;
+
+ public StorageSettings(
+ @Value("${app.storage.max-object-bytes:33554432}") long maxObjectBytes,
+ @Value("${app.storage.max-keys:1000}") int maxKeys,
+ @Value("${app.storage.sql.statement-timeout-seconds:30}") int sqlStatementTimeoutSeconds,
+ @Value("${app.storage.sql.max-rows:1000}") int sqlMaxRows,
+ @Value("${app.storage.connect-timeout-seconds:10}") int connectTimeoutSeconds) {
+ this.maxObjectBytes = maxObjectBytes;
+ this.maxKeys = maxKeys;
+ this.sqlStatementTimeoutSeconds = sqlStatementTimeoutSeconds;
+ this.sqlMaxRows = sqlMaxRows;
+ this.connectTimeoutSeconds = connectTimeoutSeconds;
+ }
+
+ public static StorageSettings defaults() {
+ return new StorageSettings(32L * 1024 * 1024, 1000, 30, 1000, 10);
+ }
+
+ public long maxObjectBytes() {
+ return maxObjectBytes;
+ }
+
+ public int maxKeys() {
+ return maxKeys;
+ }
+
+ public int sqlStatementTimeoutSeconds() {
+ return sqlStatementTimeoutSeconds;
+ }
+
+ public int sqlMaxRows() {
+ return sqlMaxRows;
+ }
+
+ public int connectTimeoutSeconds() {
+ return connectTimeoutSeconds;
+ }
+}
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
new file mode 100644
index 0000000..85c0a40
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageType.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.storage;
+
+import java.util.List;
+import java.util.Set;
+
+import it.cnr.isti.workflow.manager.blocks.IOCapability;
+
+/**
+ * A kind of storage a flow can keep data in - an S3-compatible object store, a PostgreSQL database.
+ *
+ * Pluggable the way LLM providers are: each is a Spring bean, collected by {@link StorageTypes}
+ * and listed on {@code /storage/types}. The storage node and the storage block speak only to this
+ * interface, in four generic terms - a template to write to, a query to read with, a script to
+ * prepare, a list of shapes to read in - and each type says what they mean for it and checks them
+ * while the flow is edited, so a mistake shows up in the editor rather than halfway through a run.
+ */
+public interface StorageType {
+
+ /** How flows, catalog entries and vault credentials name this type, e.g. "S3". */
+ String getName();
+
+ String getDescription();
+
+ /** The fields of a person's own connection, as the vault form asks for them. */
+ List credentialFields();
+
+ /** The shapes a read can come back in. The first is the default. */
+ List viewShapes();
+
+ /** What a write accepts, for the kinds of value the target's port takes. */
+ List writableKinds();
+
+ /** A connection's settings, checked when the catalog loads or a person saves their own. */
+ void validateConnection(StorageConnection connection);
+
+ /** Problems with a write template, in words for the flow author; empty when there are none. */
+ List validateTemplate(String template);
+
+ /** Problems with a read query for the shape asked. */
+ List validateQuery(String query, ViewShape shape);
+
+ /** Problems with a delete template. */
+ List validateDelete(String template);
+
+ List validateInit(StorageInit init);
+
+ /** 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);
+ }
+
+ StorageSession open(StorageConnection connection, StorageScope scope);
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypeMetadata.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypeMetadata.java
new file mode 100644
index 0000000..5ce69a9
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypeMetadata.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;
+
+import java.util.List;
+
+/**
+ * What the editor and the vault form need to know about a storage type - never anything about a
+ * connection to one.
+ *
+ * @param viewShapes what a read can return, the default first
+ * @param writableKinds the kinds of value a write accepts (FILE, TEXT, JSON)
+ * @param credentialFields the fields of a person's own connection, for the vault form
+ */
+public record StorageTypeMetadata(String name, String description, List viewShapes,
+ List writableKinds, List credentialFields) {
+
+ public static StorageTypeMetadata of(StorageType type) {
+ return new StorageTypeMetadata(type.getName(), type.getDescription(),
+ type.viewShapes().stream().map(Enum::name).toList(),
+ type.writableKinds().stream().map(kind -> kind.type().name()).distinct().toList(),
+ type.credentialFields());
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypes.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypes.java
new file mode 100644
index 0000000..027818e
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypes.java
@@ -0,0 +1,40 @@
+// 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 java.util.Comparator;
+import java.util.List;
+import java.util.Optional;
+
+import org.springframework.stereotype.Component;
+
+/** Every {@link StorageType} this service has, found by the name flows use for it. */
+@Component
+public class StorageTypes {
+
+ private final List types;
+
+ public StorageTypes(List types) {
+ this.types = types.stream()
+ .sorted(Comparator.comparing(StorageType::getName, String.CASE_INSENSITIVE_ORDER))
+ .toList();
+ }
+
+ public List all() {
+ return types;
+ }
+
+ public Optional find(String name) {
+ if (name == null || name.isBlank()) {
+ return Optional.empty();
+ }
+ return types.stream().filter(type -> type.getName().equalsIgnoreCase(name.trim())).findFirst();
+ }
+
+ public StorageType require(String name) {
+ return find(name).orElseThrow(() -> new StorageException(StorageException.CONNECTION_INVALID,
+ "Unknown storage type: " + name));
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageWriteResult.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageWriteResult.java
new file mode 100644
index 0000000..0815ca8
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageWriteResult.java
@@ -0,0 +1,17 @@
+// 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 tools.jackson.databind.JsonNode;
+
+/**
+ * What a write or a delete did.
+ *
+ * @param location where it went: the object key, or null for a statement
+ * @param affected objects or rows touched
+ * @param rows what a statement's {@code RETURNING} gave back, null when it had none
+ */
+public record StorageWriteResult(String location, long affected, JsonNode rows) {
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/ViewShape.java b/src/main/java/it/cnr/isti/workflow/manager/storage/ViewShape.java
new file mode 100644
index 0000000..3fdaa7b
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/ViewShape.java
@@ -0,0 +1,39 @@
+// 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 it.cnr.isti.workflow.manager.ios.IOType;
+
+/**
+ * What a read hands to the step that asked for it. A type offers the shapes that mean something for
+ * it - an object store has files, a database has rows - and a node asking for another is refused
+ * while the flow is edited.
+ */
+public enum ViewShape {
+ /** Each matching object as a file, for a node that takes files. */
+ FILES(IOType.FILE, true),
+ /** Each matching object's content as text. */
+ TEXTS(IOType.TEXT, true),
+ /** Only the names of what matched. */
+ KEYS(IOType.TEXT, true),
+ /** One JSON value: the rows of a query, or the parsed content of the matching objects. */
+ JSON(IOType.JSON, false);
+
+ private final IOType portType;
+ private final boolean multiple;
+
+ ViewShape(IOType portType, boolean multiple) {
+ this.portType = portType;
+ this.multiple = multiple;
+ }
+
+ public IOType portType() {
+ return portType;
+ }
+
+ public boolean multiple() {
+ return multiple;
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresSession.java b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresSession.java
new file mode 100644
index 0000000..bdfb200
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresSession.java
@@ -0,0 +1,320 @@
+// 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.postgres;
+
+import java.io.File;
+import java.math.BigDecimal;
+import java.math.BigInteger;
+import java.sql.Array;
+import java.sql.Connection;
+import java.sql.PreparedStatement;
+import java.sql.ResultSet;
+import java.sql.ResultSetMetaData;
+import java.sql.SQLException;
+import java.sql.Statement;
+import java.sql.Types;
+import java.util.ArrayList;
+import java.util.Base64;
+import java.util.Collection;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.function.Function;
+
+import it.cnr.isti.workflow.manager.storage.StorageConnection;
+import it.cnr.isti.workflow.manager.storage.StorageException;
+import it.cnr.isti.workflow.manager.storage.StorageInit;
+import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
+import it.cnr.isti.workflow.manager.storage.StorageReadResult;
+import it.cnr.isti.workflow.manager.storage.StorageSession;
+import it.cnr.isti.workflow.manager.storage.StorageSettings;
+import it.cnr.isti.workflow.manager.storage.StorageWriteResult;
+import it.cnr.isti.workflow.manager.storage.ViewShape;
+import tools.jackson.databind.JsonNode;
+import tools.jackson.databind.ObjectMapper;
+import tools.jackson.databind.node.ArrayNode;
+import tools.jackson.databind.node.ObjectNode;
+
+/** One connection, for one piece of work: every statement in a transaction of its own. */
+final class PostgresSession implements StorageSession {
+
+ private final Connection jdbc;
+ private final StorageConnection connection;
+ private final StorageSettings settings;
+ private final ObjectMapper objectMapper;
+
+ PostgresSession(Connection jdbc, StorageConnection connection, StorageSettings settings, ObjectMapper objectMapper) {
+ this.jdbc = jdbc;
+ this.connection = connection;
+ this.settings = settings;
+ this.objectMapper = objectMapper;
+ }
+
+ @Override
+ public void initialize(StorageInit init) {
+ if (init == null || !init.hasScript()) {
+ return;
+ }
+ if (!connection.allowInit()) {
+ throw new StorageException(StorageException.INIT_FAILED,
+ connection.name() + " may not be prepared by flows: its init script cannot run");
+ }
+ inTransaction(false, "run the init script", () -> {
+ try (Statement statement = jdbc.createStatement()) {
+ statement.execute(init.script());
+ }
+ return null;
+ });
+ }
+
+ @Override
+ public StorageWriteResult write(String template, Object value, String contentType, Function values) {
+ requireWritable("write to");
+ return change(template, PostgresStorageType.WRITE_VERBS, "write", StoragePlaceholders.withValue(value, values));
+ }
+
+ @Override
+ public StorageReadResult read(String query, ViewShape shape, Function values) {
+ SqlTemplate template = checked(query, PostgresStorageType.READ_VERBS, "read");
+ return inTransaction(true, "read", () -> {
+ try (PreparedStatement statement = prepare(template, values)) {
+ statement.setMaxRows(settings.sqlMaxRows() + 1);
+ try (ResultSet rows = statement.executeQuery()) {
+ ArrayNode array = rowsOf(rows);
+ boolean truncated = array.size() > settings.sqlMaxRows();
+ if (truncated) {
+ array.remove(array.size() - 1);
+ }
+ return new StorageReadResult(array, array.size(), truncated);
+ }
+ }
+ });
+ }
+
+ @Override
+ public StorageReadResult list(String query, Function values) {
+ String schema = query == null || query.isBlank() ? null : StoragePlaceholders.render(query, values).trim();
+ return inTransaction(true, "list the tables", () -> {
+ try (PreparedStatement statement = jdbc.prepareStatement(
+ "SELECT table_name FROM information_schema.tables WHERE table_schema = coalesce(?, current_schema())"
+ + " AND table_type IN ('BASE TABLE', 'VIEW') ORDER BY table_name")) {
+ statement.setString(1, schema);
+ statement.setMaxRows(settings.maxKeys() + 1);
+ List tables = new ArrayList<>();
+ try (ResultSet rows = statement.executeQuery()) {
+ while (rows.next()) {
+ tables.add(rows.getString(1));
+ }
+ }
+ boolean truncated = tables.size() > settings.maxKeys();
+ List kept = truncated ? tables.subList(0, settings.maxKeys()) : tables;
+ return new StorageReadResult(List.copyOf(kept), kept.size(), truncated);
+ }
+ });
+ }
+
+ @Override
+ public StorageWriteResult delete(String template, Function values) {
+ requireWritable("delete from");
+ return change(template, PostgresStorageType.DELETE_VERBS, "delete", values);
+ }
+
+ @Override
+ public void close() {
+ try {
+ jdbc.close();
+ } catch (SQLException ignored) {
+ // Returned to the pool, or dropped: either way nothing more can be done with it.
+ }
+ }
+
+ private StorageWriteResult change(String sql, Set verbs, String what, Function values) {
+ SqlTemplate template = checked(sql, verbs, what);
+ return inTransaction(false, what, () -> {
+ try (PreparedStatement statement = prepare(template, values)) {
+ if (statement.execute()) {
+ try (ResultSet rows = statement.getResultSet()) {
+ ArrayNode returned = rowsOf(rows);
+ return new StorageWriteResult(null, returned.size(), returned);
+ }
+ }
+ return new StorageWriteResult(null, Math.max(0, statement.getUpdateCount()), null);
+ }
+ });
+ }
+
+ private void requireWritable(String verb) {
+ if (connection.readOnly()) {
+ throw new StorageException(StorageException.OPERATION_NOT_ALLOWED,
+ connection.name() + " is read-only: nothing can " + verb + " it");
+ }
+ }
+
+ /** Checked again here, not only in the editor: a flow saved before a rule existed still meets it. */
+ private SqlTemplate checked(String sql, Set verbs, String what) {
+ SqlTemplate template = SqlTemplate.parse(sql);
+ if (!template.problems().isEmpty()) {
+ throw new StorageException(StorageException.SQL_REJECTED, String.join("; ", template.problems()));
+ }
+ if (!verbs.contains(template.verb())) {
+ throw new StorageException(StorageException.SQL_REJECTED,
+ "A " + what + " must start with " + String.join(" or ", verbs.stream().sorted().toList())
+ + ", not " + (template.verb().isEmpty() ? "nothing" : template.verb()));
+ }
+ return template;
+ }
+
+ private PreparedStatement prepare(SqlTemplate template, Function values) throws SQLException {
+ PreparedStatement statement = jdbc.prepareStatement(template.sql());
+ List parameters = template.parameters();
+ for (int i = 0; i < parameters.size(); i++) {
+ bind(statement, i + 1, parameters.get(i), values == null ? null : values.apply(parameters.get(i)));
+ }
+ return statement;
+ }
+
+ private void bind(PreparedStatement statement, int index, String name, Object value) throws SQLException {
+ if (value instanceof JsonNode node) {
+ if (node.isNull() || node.isMissingNode()) {
+ value = null;
+ } else if (node.isString()) {
+ value = node.asString();
+ } else if (node.isBoolean()) {
+ value = node.asBoolean();
+ } else if (node.isNumber()) {
+ value = node.numberValue();
+ } else {
+ statement.setObject(index, node.toString(), Types.OTHER);
+ return;
+ }
+ }
+ if (value == null) {
+ if (!StoragePlaceholders.isValueName(name)) {
+ throw new StorageException(StorageException.PLACEHOLDER_MISSING, "No value for ${{" + name + "}}");
+ }
+ statement.setNull(index, Types.OTHER);
+ } else if (value instanceof File) {
+ throw new StorageException(StorageException.SQL_REJECTED,
+ "${{" + name + "}} is a file, and a file cannot be bound into a statement");
+ } else if (value instanceof Boolean flag) {
+ statement.setBoolean(index, flag);
+ } else if (value instanceof Integer || value instanceof Long || value instanceof Short) {
+ statement.setLong(index, ((Number) value).longValue());
+ } else if (value instanceof BigInteger || value instanceof BigDecimal || value instanceof Double
+ || value instanceof Float) {
+ statement.setBigDecimal(index, new BigDecimal(value.toString()));
+ } else if (value instanceof Map, ?> || value instanceof Collection>) {
+ statement.setObject(index, objectMapper.writeValueAsString(value), Types.OTHER);
+ } else {
+ statement.setString(index, value.toString());
+ }
+ }
+
+ private ArrayNode rowsOf(ResultSet rows) throws SQLException {
+ ArrayNode array = objectMapper.createArrayNode();
+ ResultSetMetaData meta = rows.getMetaData();
+ while (rows.next()) {
+ ObjectNode row = array.addObject();
+ for (int column = 1; column <= meta.getColumnCount(); column++) {
+ row.set(meta.getColumnLabel(column), cell(rows, column, meta.getColumnTypeName(column)));
+ }
+ }
+ return array;
+ }
+
+ private JsonNode cell(ResultSet rows, int column, String typeName) throws SQLException {
+ Object value = rows.getObject(column);
+ if (value == null) {
+ return objectMapper.nullNode();
+ }
+ if ("json".equalsIgnoreCase(typeName) || "jsonb".equalsIgnoreCase(typeName)) {
+ try {
+ return objectMapper.readTree(rows.getString(column));
+ } catch (Exception ex) {
+ return objectMapper.getNodeFactory().stringNode(rows.getString(column));
+ }
+ }
+ if (value instanceof Array array) {
+ ArrayNode items = objectMapper.createArrayNode();
+ for (Object item : (Object[]) array.getArray()) {
+ items.add(item == null ? objectMapper.nullNode() : objectMapper.valueToTree(item instanceof Number || item instanceof Boolean ? item : item.toString()));
+ }
+ return items;
+ }
+ if (value instanceof byte[] bytes) {
+ return objectMapper.getNodeFactory().stringNode(Base64.getEncoder().encodeToString(bytes));
+ }
+ if (value instanceof Number || value instanceof Boolean) {
+ return objectMapper.valueToTree(value);
+ }
+ return objectMapper.getNodeFactory().stringNode(value.toString());
+ }
+
+ @FunctionalInterface
+ private interface Work {
+ T run() throws SQLException;
+ }
+
+ /** One statement's transaction: read-only when asked, with the statement timeout set inside it. */
+ private T inTransaction(boolean readOnly, String what, Work work) {
+ try {
+ jdbc.setAutoCommit(false);
+ jdbc.setReadOnly(readOnly);
+ try (Statement timeout = jdbc.createStatement()) {
+ timeout.execute("SET LOCAL statement_timeout = " + (settings.sqlStatementTimeoutSeconds() * 1000L));
+ }
+ T result = work.run();
+ if (readOnly) {
+ jdbc.rollback();
+ } else {
+ jdbc.commit();
+ }
+ return result;
+ } catch (SQLException ex) {
+ rollbackQuietly();
+ throw failure(what, ex);
+ } catch (RuntimeException ex) {
+ rollbackQuietly();
+ throw ex;
+ } finally {
+ try {
+ jdbc.setReadOnly(false);
+ jdbc.setAutoCommit(true);
+ } catch (SQLException ignored) {
+ // The connection is closed right after; a pooled one is reset by the pool.
+ }
+ }
+ }
+
+ private void rollbackQuietly() {
+ try {
+ jdbc.rollback();
+ } catch (SQLException ignored) {
+ // Already failed; the original error is the one worth reporting.
+ }
+ }
+
+ private StorageException failure(String what, SQLException ex) {
+ String state = ex.getSQLState() == null ? "" : ex.getSQLState();
+ if (state.equals("25006")) {
+ return new StorageException(StorageException.OPERATION_NOT_ALLOWED,
+ "A read cannot change data in " + connection.name() + ": " + ex.getMessage(), ex);
+ }
+ if (state.startsWith("08") || state.equals("57P01") || state.equals("53300")) {
+ return new StorageException(StorageException.UNREACHABLE,
+ "Could not " + what + " on " + connection.name() + ": " + ex.getMessage(), ex);
+ }
+ if (state.equals("57014")) {
+ return new StorageException(StorageException.SQL_REJECTED, "The " + what + " on " + connection.name()
+ + " took longer than " + settings.sqlStatementTimeoutSeconds() + " seconds and was stopped", ex);
+ }
+ if (state.startsWith("28") || state.equals("42501")) {
+ return new StorageException(StorageException.OPERATION_NOT_ALLOWED,
+ connection.name() + " refused to " + what + ": " + ex.getMessage(), ex);
+ }
+ return new StorageException(what.equals("run the init script") ? StorageException.INIT_FAILED : StorageException.SQL_REJECTED,
+ connection.name() + " refused the " + what + ": " + ex.getMessage(), ex);
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageType.java b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageType.java
new file mode 100644
index 0000000..98634f1
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageType.java
@@ -0,0 +1,274 @@
+// 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.postgres;
+
+import java.net.URLEncoder;
+import java.nio.charset.StandardCharsets;
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.SQLException;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.regex.Pattern;
+
+import org.springframework.beans.factory.DisposableBean;
+import org.springframework.stereotype.Component;
+
+import com.zaxxer.hikari.HikariConfig;
+import com.zaxxer.hikari.HikariDataSource;
+
+import it.cnr.isti.workflow.manager.blocks.IOCapability;
+import it.cnr.isti.workflow.manager.blocks.IOCapabilityType;
+import it.cnr.isti.workflow.manager.storage.CredentialField;
+import it.cnr.isti.workflow.manager.storage.StorageConnection;
+import it.cnr.isti.workflow.manager.storage.StorageEndpointGuard;
+import it.cnr.isti.workflow.manager.storage.StorageException;
+import it.cnr.isti.workflow.manager.storage.StorageInit;
+import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
+import it.cnr.isti.workflow.manager.storage.StorageScope;
+import it.cnr.isti.workflow.manager.storage.StorageSession;
+import it.cnr.isti.workflow.manager.storage.StorageSettings;
+import it.cnr.isti.workflow.manager.storage.StorageType;
+import it.cnr.isti.workflow.manager.storage.ViewShape;
+import tools.jackson.databind.ObjectMapper;
+
+/**
+ * A PostgreSQL database, spoken to in SQL.
+ *
+ *
+ * - A query is a {@code SELECT} (or {@code WITH}), run in a read-only transaction, returning its
+ * rows as JSON.
+ *
- A template is an {@code INSERT}, {@code UPDATE} or {@code MERGE}; a delete is a
+ * {@code DELETE}. Both may end in {@code RETURNING}.
+ *
- Init is a script - {@code CREATE TABLE IF NOT EXISTS ...} - run in one transaction.
+ *
+ *
+ * Placeholders are bound, never spliced in: see {@link SqlTemplate}. The JDBC URL is built here
+ * from a host, a port and a database, and the driver's properties are set here too. Nothing a
+ * person types reaches either, because pgjdbc's URL parameters are themselves a way in: a
+ * {@code socketFactory} or {@code loggerFile} of the caller's choosing ran code or wrote files on
+ * the server (CVE-2022-21724).
+ */
+@Component
+public class PostgresStorageType implements StorageType, DisposableBean {
+
+ public static final String NAME = "PostgreSQL";
+
+ private static final Pattern HOST = Pattern.compile("^(?:[A-Za-z0-9](?:[A-Za-z0-9-]{0,61}[A-Za-z0-9])?)(?:\\.[A-Za-z0-9](?:[A-Za-z0-9-]{0,61}[A-Za-z0-9])?)*$|^\\[[0-9A-Fa-f:.]+]$");
+ private static final Pattern DATABASE = Pattern.compile("^[A-Za-z0-9_$.-]{1,63}$");
+ private static final Pattern IDENTIFIER = Pattern.compile("^[A-Za-z_][A-Za-z0-9_$]{0,62}$");
+ private static final Set SSL_MODES = Set.of("disable", "allow", "prefer", "require", "verify-ca", "verify-full");
+ static final Set READ_VERBS = Set.of("SELECT", "WITH", "VALUES", "TABLE");
+ static final Set WRITE_VERBS = Set.of("INSERT", "UPDATE", "MERGE");
+ static final Set DELETE_VERBS = Set.of("DELETE");
+
+ private final StorageSettings settings;
+ private final StorageEndpointGuard guard;
+ private final ObjectMapper objectMapper;
+ private final Map pools = new ConcurrentHashMap<>();
+
+ public PostgresStorageType(StorageSettings settings, StorageEndpointGuard guard, ObjectMapper objectMapper) {
+ this.settings = settings;
+ this.guard = guard;
+ this.objectMapper = objectMapper;
+ }
+
+ @Override
+ public String getName() {
+ return NAME;
+ }
+
+ @Override
+ public String getDescription() {
+ return "A PostgreSQL database. Reads are SELECT queries returning rows as JSON; writes are INSERT, UPDATE or "
+ + "MERGE statements; values are bound as parameters.";
+ }
+
+ @Override
+ public List credentialFields() {
+ return List.of(
+ new CredentialField("host", "Host", "e.g. db.example.org", false, true),
+ new CredentialField("port", "Port", "5432 when empty", false, false),
+ new CredentialField("database", "Database", null, false, true),
+ new CredentialField("schema", "Schema", "public when empty", false, false),
+ new CredentialField("sslMode", "SSL mode", "disable, allow, prefer, require, verify-ca or verify-full; prefer when empty", false, false),
+ new CredentialField("user", "User", null, false, true),
+ new CredentialField("password", "Password", null, true, true));
+ }
+
+ @Override
+ public List viewShapes() {
+ return List.of(ViewShape.JSON);
+ }
+
+ @Override
+ public List writableKinds() {
+ return List.of(new IOCapability(IOCapabilityType.JSON, false), new IOCapability(IOCapabilityType.TEXT, false));
+ }
+
+ @Override
+ public void validateConnection(StorageConnection connection) {
+ String host = connection.requireText("host");
+ if (!HOST.matcher(host).matches()) {
+ throw invalid(connection, "the host " + host + " is not a host name or address");
+ }
+ port(connection);
+ String database = connection.requireText("database");
+ if (!DATABASE.matcher(database).matches()) {
+ throw invalid(connection, "the database name " + database + " is not allowed");
+ }
+ String schema = connection.optionalText("schema");
+ if (schema != null && !IDENTIFIER.matcher(schema).matches()) {
+ throw invalid(connection, "the schema " + schema + " is not a plain identifier");
+ }
+ String sslMode = connection.optionalText("sslMode");
+ if (sslMode != null && !SSL_MODES.contains(sslMode)) {
+ throw invalid(connection, "the SSL mode must be one of " + SSL_MODES);
+ }
+ connection.requireText("user");
+ connection.requireText("password");
+ guard.check(connection, null, "host");
+ }
+
+ @Override
+ public List validateTemplate(String template) {
+ return validateStatement(template, WRITE_VERBS, "A write", true);
+ }
+
+ @Override
+ public List validateQuery(String query, ViewShape shape) {
+ List problems = validateStatement(query, READ_VERBS, "A read", false);
+ if (shape != null && shape != ViewShape.JSON) {
+ problems.add("A PostgreSQL read returns rows as JSON; it cannot return " + shape);
+ }
+ return problems;
+ }
+
+ @Override
+ public List validateDelete(String template) {
+ return validateStatement(template, DELETE_VERBS, "A delete", false);
+ }
+
+ @Override
+ public List validateInit(StorageInit init) {
+ List problems = new ArrayList<>();
+ if (init == null) {
+ return problems;
+ }
+ if (init.createIfMissing()) {
+ problems.add("A PostgreSQL storage has nothing to create by itself; put CREATE TABLE IF NOT EXISTS ... in the init script");
+ }
+ if (init.hasScript() && init.script().contains("${{")) {
+ problems.add("The init script cannot have placeholders: it runs as it is written, before any step");
+ }
+ return problems;
+ }
+
+ @Override
+ public StorageSession open(StorageConnection connection, StorageScope scope) {
+ validateConnection(connection);
+ try {
+ Connection jdbc = connection.operatorChosen() ? pool(connection).getConnection() : connect(connection);
+ return new PostgresSession(jdbc, connection, settings, objectMapper);
+ } catch (SQLException ex) {
+ throw new StorageException(StorageException.UNREACHABLE,
+ "Could not connect to " + connection.name() + ": " + ex.getMessage(), ex);
+ }
+ }
+
+ @Override
+ public void destroy() {
+ pools.values().forEach(HikariDataSource::close);
+ pools.clear();
+ }
+
+ private static List validateStatement(String sql, Set verbs, String what, boolean valueAllowed) {
+ List problems = new ArrayList<>();
+ if (sql == null || sql.isBlank()) {
+ problems.add(what + " needs a statement");
+ return problems;
+ }
+ SqlTemplate template = SqlTemplate.parse(sql);
+ problems.addAll(template.problems());
+ if (!verbs.contains(template.verb())) {
+ problems.add(what + " must start with " + String.join(" or ", verbs.stream().sorted().toList())
+ + (template.verb().isEmpty() ? "" : ", not " + template.verb()));
+ }
+ StoragePlaceholders.invalidNames(sql).forEach(name -> problems.add(
+ "${{" + name + "}} is not a usable name: start with a letter, then letters, digits, '-', '_' or '.'"));
+ if (!valueAllowed && StoragePlaceholders.names(sql).stream().anyMatch(StoragePlaceholders::isValueName)) {
+ problems.add(what + " has no value: ${{value}} only means something when writing");
+ }
+ return problems;
+ }
+
+ private static StorageException invalid(StorageConnection connection, String problem) {
+ return new StorageException(StorageException.CONNECTION_INVALID, connection.name() + ": " + problem);
+ }
+
+ private static int port(StorageConnection connection) {
+ String port = connection.optionalText("port");
+ if (port == null) {
+ return 5432;
+ }
+ try {
+ int value = Integer.parseInt(port);
+ if (value < 1 || value > 65535) {
+ throw invalid(connection, "the port must be between 1 and 65535");
+ }
+ return value;
+ } catch (NumberFormatException ex) {
+ throw invalid(connection, "the port " + port + " is not a number");
+ }
+ }
+
+ private String url(StorageConnection connection) {
+ return "jdbc:postgresql://" + connection.requireText("host") + ":" + port(connection) + "/"
+ + URLEncoder.encode(connection.requireText("database"), StandardCharsets.UTF_8);
+ }
+
+ /** Every driver property this connection has - set here, from checked settings, and nowhere else. */
+ private Properties properties(StorageConnection connection) {
+ Properties properties = new Properties();
+ properties.setProperty("user", connection.requireText("user"));
+ properties.setProperty("password", connection.requireText("password"));
+ properties.setProperty("sslmode", connection.optionalText("sslMode") == null ? "prefer" : connection.optionalText("sslMode"));
+ properties.setProperty("ApplicationName", "humainflow-storage");
+ properties.setProperty("connectTimeout", String.valueOf(settings.connectTimeoutSeconds()));
+ properties.setProperty("loginTimeout", String.valueOf(settings.connectTimeoutSeconds()));
+ // Strings go to the server untyped, so a placeholder bound as text fits an integer or a
+ // jsonb column the way a literal would, instead of failing as "character varying".
+ properties.setProperty("stringtype", "unspecified");
+ String schema = connection.optionalText("schema");
+ if (schema != null) {
+ properties.setProperty("currentSchema", schema);
+ }
+ return properties;
+ }
+
+ private Connection connect(StorageConnection connection) throws SQLException {
+ return DriverManager.getConnection(url(connection), properties(connection));
+ }
+
+ private HikariDataSource pool(StorageConnection connection) {
+ String key = connection.name() + "\u0000" + connection.settings();
+ return pools.computeIfAbsent(key, ignored -> {
+ HikariConfig config = new HikariConfig();
+ config.setPoolName("storage-" + connection.name().replaceAll("[^A-Za-z0-9_-]", "_"));
+ config.setJdbcUrl(url(connection));
+ config.setDataSourceProperties(properties(connection));
+ config.setMaximumPoolSize(4);
+ config.setMinimumIdle(0);
+ config.setIdleTimeout(60_000);
+ config.setConnectionTimeout(settings.connectTimeoutSeconds() * 1000L);
+ config.setInitializationFailTimeout(-1);
+ return new HikariDataSource(config);
+ });
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplate.java b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplate.java
new file mode 100644
index 0000000..b10e99a
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplate.java
@@ -0,0 +1,224 @@
+// 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.postgres;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Locale;
+
+/**
+ * A statement written with ${{name}} placeholders, turned into one the driver can bind.
+ *
+ * Every placeholder becomes a {@code ?} and its name joins {@link #parameters()} in order, so
+ * each value is sent to the server as a value, never spliced into the text: a value of
+ * {@code '; DROP TABLE x; --} is stored as those characters. That only holds for a placeholder
+ * standing where a value can - so one found inside a string literal, a quoted identifier, a
+ * dollar-quoted body or a comment is refused, with the fix spelled out, rather than left there as
+ * literal text. A {@code ?} the author wrote themselves, the jsonb operator say, is doubled so the
+ * driver does not take it for a parameter.
+ *
+ *
One statement only. The driver does not see to that - pgjdbc splits {@code INSERT ...; DROP
+ * TABLE x} into two and runs both, prepared or not - so a second statement is refused here.
+ */
+final class SqlTemplate {
+
+ private final String sql;
+ private final List parameters;
+ private final String verb;
+ private final List problems;
+
+ private SqlTemplate(String sql, List parameters, String verb, List problems) {
+ this.sql = sql;
+ this.parameters = parameters;
+ this.verb = verb;
+ this.problems = problems;
+ }
+
+ /** The statement to prepare. Meaningful only when {@link #problems()} is empty. */
+ String sql() {
+ return sql;
+ }
+
+ /** The placeholder bound to each {@code ?}, in order - a name used twice is bound twice. */
+ List parameters() {
+ return parameters;
+ }
+
+ /** The first keyword, upper case: SELECT, INSERT, WITH... Empty for a blank statement. */
+ String verb() {
+ return verb;
+ }
+
+ List problems() {
+ return problems;
+ }
+
+ static SqlTemplate parse(String template) {
+ List problems = new ArrayList<>();
+ List parameters = new ArrayList<>();
+ StringBuilder out = new StringBuilder();
+ String text = template == null ? "" : template;
+ int i = 0;
+ int length = text.length();
+ while (i < length) {
+ char c = text.charAt(i);
+ if (text.startsWith("${{", i)) {
+ int end = text.indexOf("}}", i + 3);
+ if (end < 0) {
+ problems.add("A placeholder starting at character " + (i + 1) + " is never closed with }}");
+ break;
+ }
+ parameters.add(text.substring(i + 3, end).trim());
+ out.append('?');
+ i = end + 2;
+ } else if (c == '\'') {
+ boolean escapes = i > 0 && (text.charAt(i - 1) == 'E' || text.charAt(i - 1) == 'e')
+ && (i < 2 || !Character.isLetterOrDigit(text.charAt(i - 2)));
+ i = copyQuoted(text, i, '\'', escapes, out, problems, "a string literal");
+ } else if (c == '"') {
+ i = copyQuoted(text, i, '"', false, out, problems, "a quoted identifier");
+ } else if (c == '$' && dollarTag(text, i) != null) {
+ String tag = dollarTag(text, i);
+ int close = text.indexOf(tag, i + tag.length());
+ int stop = close < 0 ? length : close + tag.length();
+ checkNoPlaceholder(text.substring(i, stop), problems, "a dollar-quoted body");
+ out.append(text, i, stop);
+ i = stop;
+ } else if (text.startsWith("--", i)) {
+ int stop = text.indexOf('\n', i);
+ stop = stop < 0 ? length : stop;
+ checkNoPlaceholder(text.substring(i, stop), problems, "a comment");
+ out.append(text, i, stop);
+ i = stop;
+ } else if (text.startsWith("/*", i)) {
+ int stop = blockCommentEnd(text, i);
+ checkNoPlaceholder(text.substring(i, stop), problems, "a comment");
+ out.append(text, i, stop);
+ i = stop;
+ } else if (c == '?') {
+ out.append("??");
+ i++;
+ } else if (c == ';') {
+ if (!onlyCommentsAndSpaceFrom(text, i + 1)) {
+ problems.add("Only one statement is allowed: remove what follows the ';' at character " + (i + 1));
+ break;
+ }
+ // A closing ';' is harmless, and dropped: the driver would otherwise split on it.
+ i = length;
+ } else {
+ out.append(c);
+ i++;
+ }
+ }
+ return new SqlTemplate(out.toString(), List.copyOf(parameters), firstKeyword(text), List.copyOf(problems));
+ }
+
+ private static int copyQuoted(String text, int start, char quote, boolean backslashEscapes, StringBuilder out,
+ List problems, String what) {
+ int i = start + 1;
+ while (i < text.length()) {
+ char c = text.charAt(i);
+ if (backslashEscapes && c == '\\') {
+ i += 2;
+ continue;
+ }
+ if (c == quote) {
+ if (i + 1 < text.length() && text.charAt(i + 1) == quote) {
+ i += 2;
+ continue;
+ }
+ i++;
+ break;
+ }
+ i++;
+ }
+ int stop = Math.min(i, text.length());
+ checkNoPlaceholder(text.substring(start, stop), problems, what);
+ out.append(text, start, stop);
+ return stop;
+ }
+
+ private static void checkNoPlaceholder(String part, List problems, String where) {
+ int at = part.indexOf("${{");
+ if (at >= 0) {
+ int end = part.indexOf("}}", at);
+ String placeholder = end < 0 ? part.substring(at) : part.substring(at, end + 2);
+ problems.add(placeholder + " is inside " + where + ", where it would stay literal text. Write it on its own, "
+ + "without quotes - the value is sent as a parameter: WHERE name = " + placeholder);
+ }
+ }
+
+ /** $tag$ or $$ opening a dollar-quoted string at i, or null. */
+ private static String dollarTag(String text, int i) {
+ int j = i + 1;
+ while (j < text.length() && (Character.isLetterOrDigit(text.charAt(j)) || text.charAt(j) == '_')) {
+ j++;
+ }
+ if (j < text.length() && text.charAt(j) == '$' && (j == i + 1 || !Character.isDigit(text.charAt(i + 1)))) {
+ return text.substring(i, j + 1);
+ }
+ return null;
+ }
+
+ private static boolean onlyCommentsAndSpaceFrom(String text, int start) {
+ int i = start;
+ while (i < text.length()) {
+ if (Character.isWhitespace(text.charAt(i)) || text.charAt(i) == ';') {
+ i++;
+ } else if (text.startsWith("--", i)) {
+ int stop = text.indexOf('\n', i);
+ i = stop < 0 ? text.length() : stop;
+ } else if (text.startsWith("/*", i)) {
+ i = blockCommentEnd(text, i);
+ } else {
+ return false;
+ }
+ }
+ return true;
+ }
+
+ /** Block comments nest in PostgreSQL. */
+ private static int blockCommentEnd(String text, int start) {
+ int depth = 0;
+ int i = start;
+ while (i < text.length()) {
+ if (text.startsWith("/*", i)) {
+ depth++;
+ i += 2;
+ } else if (text.startsWith("*/", i)) {
+ depth--;
+ i += 2;
+ if (depth == 0) {
+ return i;
+ }
+ } else {
+ i++;
+ }
+ }
+ return text.length();
+ }
+
+ private static String firstKeyword(String text) {
+ int i = 0;
+ while (i < text.length()) {
+ char c = text.charAt(i);
+ if (Character.isWhitespace(c) || c == '(') {
+ i++;
+ } else if (text.startsWith("--", i)) {
+ int stop = text.indexOf('\n', i);
+ i = stop < 0 ? text.length() : stop;
+ } else if (text.startsWith("/*", i)) {
+ i = blockCommentEnd(text, i);
+ } else {
+ break;
+ }
+ }
+ int start = i;
+ while (i < text.length() && Character.isLetter(text.charAt(i))) {
+ i++;
+ }
+ return text.substring(start, i).toUpperCase(Locale.ROOT);
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Keys.java b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Keys.java
new file mode 100644
index 0000000..4a56a65
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Keys.java
@@ -0,0 +1,136 @@
+// 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.s3;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.regex.Pattern;
+
+import it.cnr.isti.workflow.manager.storage.StorageException;
+import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
+
+/**
+ * Object keys and globs, checked and matched.
+ *
+ * A key is taken literally by S3, so {@code ..} is not a way up - but this service puts a prefix
+ * of its own in front of every key (the instance's root, an execution's scope), and a key that
+ * looks like it climbs out of that prefix is refused rather than trusted to be harmless everywhere.
+ * The checks run on the key as rendered, placeholders included, since a value coming from another
+ * step is exactly where a {@code ../} would come from.
+ */
+final class S3Keys {
+
+ private static final Pattern BUCKET = Pattern.compile("^[a-z0-9][a-z0-9.-]{1,61}[a-z0-9]$");
+ private static final int MAX_KEY_LENGTH = 1024;
+
+ private S3Keys() {
+ }
+
+ static boolean hasWildcard(String text) {
+ return text.indexOf('*') >= 0 || text.indexOf('?') >= 0;
+ }
+
+ /** Problems with a key or glob as written, placeholders left in; empty when there are none. */
+ static List problemsAsWritten(String text, String what) {
+ List problems = new ArrayList<>();
+ if (text == null || text.isBlank()) {
+ problems.add(what + " is empty");
+ return problems;
+ }
+ String literal = StoragePlaceholders.pattern().matcher(text).replaceAll("x");
+ if (literal.startsWith("/")) {
+ problems.add(what + " must not start with '/': it is relative to the storage's own root");
+ }
+ if (hasParentSegment(literal)) {
+ problems.add(what + " must not contain '..'");
+ }
+ if (literal.contains("//")) {
+ problems.add(what + " must not contain an empty segment ('//')");
+ }
+ StoragePlaceholders.invalidNames(text).forEach(name -> problems.add(
+ "${{" + name + "}} is not a usable name: start with a letter, then letters, digits, '-', '_' or '.'"));
+ return problems;
+ }
+
+ /** The key as it will be used, after placeholders: refused if it could leave its prefix. */
+ static String requireSafeKey(String key) {
+ if (key == null || key.isBlank()) {
+ throw new StorageException(StorageException.INVALID_LOCATION, "The object key is empty");
+ }
+ if (key.startsWith("/") || hasParentSegment(key) || key.contains("//") || key.indexOf('\0') >= 0) {
+ throw new StorageException(StorageException.INVALID_LOCATION,
+ "The object key " + key + " is not allowed: no leading '/', no '..', no empty segment");
+ }
+ if (key.length() > MAX_KEY_LENGTH) {
+ throw new StorageException(StorageException.INVALID_LOCATION,
+ "The object key is longer than " + MAX_KEY_LENGTH + " characters");
+ }
+ return key;
+ }
+
+ static boolean isValidBucket(String bucket) {
+ return bucket != null && BUCKET.matcher(bucket).matches() && !bucket.contains("..");
+ }
+
+ /** The part of a glob before its first wildcard: what to list, before matching. */
+ static String listingPrefix(String glob) {
+ int wildcard = firstWildcard(glob);
+ return wildcard < 0 ? glob : glob.substring(0, wildcard);
+ }
+
+ /**
+ * A glob as a regex: {@code *} stays inside one segment, {@code **} crosses them, {@code ?} is
+ * one character that is not '/'.
+ */
+ static Pattern globToRegex(String glob) {
+ StringBuilder regex = new StringBuilder("^");
+ for (int i = 0; i < glob.length(); i++) {
+ char c = glob.charAt(i);
+ if (c == '*') {
+ if (i + 1 < glob.length() && glob.charAt(i + 1) == '*') {
+ regex.append(".*");
+ i++;
+ } else {
+ regex.append("[^/]*");
+ }
+ } else if (c == '?') {
+ regex.append("[^/]");
+ } else {
+ regex.append(Pattern.quote(String.valueOf(c)));
+ }
+ }
+ return Pattern.compile(regex.append('$').toString());
+ }
+
+ /** The root prefix as a directory: empty, or ending in exactly one '/'. */
+ static String normalizePrefix(String prefix) {
+ if (prefix == null || prefix.isBlank()) {
+ return "";
+ }
+ String trimmed = prefix.trim();
+ while (trimmed.startsWith("/")) {
+ trimmed = trimmed.substring(1);
+ }
+ return trimmed.isEmpty() ? "" : (trimmed.endsWith("/") ? trimmed : trimmed + "/");
+ }
+
+ private static int firstWildcard(String glob) {
+ int star = glob.indexOf('*');
+ int question = glob.indexOf('?');
+ if (star < 0) {
+ return question;
+ }
+ return question < 0 ? star : Math.min(star, question);
+ }
+
+ private static boolean hasParentSegment(String key) {
+ for (String segment : key.split("/", -1)) {
+ if (segment.equals("..")) {
+ return true;
+ }
+ }
+ return false;
+ }
+}
diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Session.java b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Session.java
new file mode 100644
index 0000000..4b92ba5
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Session.java
@@ -0,0 +1,315 @@
+// 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.s3;
+
+import java.io.ByteArrayInputStream;
+import java.io.File;
+import java.io.IOException;
+import java.io.InputStream;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.List;
+import java.util.Map;
+import java.util.function.Function;
+import java.util.regex.Pattern;
+
+import io.minio.BucketExistsArgs;
+import io.minio.GetObjectArgs;
+import io.minio.ListObjectsArgs;
+import io.minio.MakeBucketArgs;
+import io.minio.MinioClient;
+import io.minio.PutObjectArgs;
+import io.minio.RemoveObjectArgs;
+import io.minio.Result;
+import io.minio.StatObjectArgs;
+import io.minio.StatObjectResponse;
+import io.minio.errors.ErrorResponseException;
+import io.minio.messages.Item;
+import it.cnr.isti.workflow.manager.storage.StorageConnection;
+import it.cnr.isti.workflow.manager.storage.StorageException;
+import it.cnr.isti.workflow.manager.storage.StorageInit;
+import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
+import it.cnr.isti.workflow.manager.storage.StorageReadResult;
+import it.cnr.isti.workflow.manager.storage.StorageSession;
+import it.cnr.isti.workflow.manager.storage.StorageSettings;
+import it.cnr.isti.workflow.manager.storage.StorageWriteResult;
+import it.cnr.isti.workflow.manager.storage.ViewShape;
+import tools.jackson.databind.JsonNode;
+import tools.jackson.databind.ObjectMapper;
+import tools.jackson.databind.node.ArrayNode;
+
+/**
+ * One piece of work against one bucket. Keys the flow sees are relative to {@code prefix} - the
+ * instance's root plus the execution's scope - which is added on the way in and taken off on the
+ * way out, so a flow never learns, or depends on, where its data sits.
+ */
+final class S3Session implements StorageSession {
+
+ private final MinioClient client;
+ private final StorageConnection connection;
+ private final String bucket;
+ private final String prefix;
+ private final StorageSettings settings;
+ private final ObjectMapper objectMapper;
+
+ S3Session(MinioClient client, StorageConnection connection, String bucket, String prefix, StorageSettings settings,
+ ObjectMapper objectMapper) {
+ this.client = client;
+ this.connection = connection;
+ this.bucket = bucket;
+ this.prefix = prefix;
+ this.settings = settings;
+ this.objectMapper = objectMapper;
+ }
+
+ @Override
+ public void initialize(StorageInit init) {
+ boolean exists = call("check the bucket " + bucket, () -> client.bucketExists(BucketExistsArgs.builder().bucket(bucket).build()));
+ if (exists) {
+ return;
+ }
+ if (init == null || !init.createIfMissing()) {
+ throw new StorageException(StorageException.INIT_FAILED,
+ "The bucket " + bucket + " of " + connection.name() + " does not exist; create it, or let the storage node create it if missing");
+ }
+ if (!connection.allowInit()) {
+ throw new StorageException(StorageException.INIT_FAILED,
+ "The bucket " + bucket + " of " + connection.name() + " does not exist, and this storage may not be prepared by flows");
+ }
+ call("create the bucket " + bucket, () -> {
+ client.makeBucket(MakeBucketArgs.builder().bucket(bucket).build());
+ return null;
+ });
+ }
+
+ @Override
+ public StorageWriteResult write(String template, Object value, String contentType, Function values) {
+ requireWritable("write to");
+ String key = S3Keys.requireSafeKey(StoragePlaceholders.render(template, values));
+ Content content = contentOf(value, contentType);
+ call("write " + key, () -> {
+ try (InputStream stream = content.open()) {
+ client.putObject(PutObjectArgs.builder().bucket(bucket).object(prefix + key)
+ .stream(stream, content.size(), -1).contentType(content.contentType()).build());
+ }
+ return null;
+ });
+ return new StorageWriteResult(key, 1, null);
+ }
+
+ @Override
+ public StorageReadResult read(String query, ViewShape shape, Function values) {
+ String rendered = S3Keys.requireSafeKey(StoragePlaceholders.render(query, values));
+ Matches matches = match(rendered);
+ return switch (shape == null ? ViewShape.FILES : shape) {
+ case KEYS -> new StorageReadResult(matches.keys(), matches.keys().size(), matches.truncated());
+ case FILES -> new StorageReadResult(downloadAll(matches.keys()), matches.keys().size(), matches.truncated());
+ case TEXTS -> new StorageReadResult(textsOf(matches.keys()), matches.keys().size(), matches.truncated());
+ case JSON -> new StorageReadResult(jsonOf(matches.keys()), matches.keys().size(), matches.truncated());
+ };
+ }
+
+ @Override
+ public StorageReadResult list(String query, Function values) {
+ String rendered = query == null || query.isBlank() ? "**" : StoragePlaceholders.render(query, values);
+ Matches matches = match(S3Keys.requireSafeKey(rendered));
+ return new StorageReadResult(matches.keys(), matches.keys().size(), matches.truncated());
+ }
+
+ @Override
+ public StorageWriteResult delete(String template, Function values) {
+ requireWritable("delete from");
+ String key = S3Keys.requireSafeKey(StoragePlaceholders.render(template, values));
+ if (S3Keys.hasWildcard(key)) {
+ throw new StorageException(StorageException.INVALID_LOCATION, "A delete names one object, not a glob: " + key);
+ }
+ boolean existed = stat(key) != null;
+ if (existed) {
+ call("delete " + key, () -> {
+ client.removeObject(RemoveObjectArgs.builder().bucket(bucket).object(prefix + key).build());
+ return null;
+ });
+ }
+ return new StorageWriteResult(key, existed ? 1 : 0, null);
+ }
+
+ @Override
+ public void close() {
+ // The HTTP client is shared by every session and outlives this one.
+ }
+
+ private void requireWritable(String verb) {
+ if (connection.readOnly()) {
+ throw new StorageException(StorageException.OPERATION_NOT_ALLOWED,
+ connection.name() + " is read-only: nothing can " + verb + " it");
+ }
+ }
+
+ private record Matches(List keys, boolean truncated) {
+ }
+
+ /** A key names itself; a glob names what it matches, at most max-keys of it. */
+ private Matches match(String keyOrGlob) {
+ if (!S3Keys.hasWildcard(keyOrGlob)) {
+ return new Matches(stat(keyOrGlob) == null ? List.of() : List.of(keyOrGlob), false);
+ }
+ Pattern pattern = S3Keys.globToRegex(keyOrGlob);
+ String listing = prefix + S3Keys.listingPrefix(keyOrGlob);
+ List keys = new ArrayList<>();
+ boolean truncated = call("list " + keyOrGlob, () -> {
+ for (Result- result : client.listObjects(ListObjectsArgs.builder().bucket(bucket).prefix(listing)
+ .recursive(true).build())) {
+ Item item = result.get();
+ if (item.isDir()) {
+ continue;
+ }
+ String key = item.objectName().substring(prefix.length());
+ if (!pattern.matcher(key).matches()) {
+ continue;
+ }
+ if (keys.size() == settings.maxKeys()) {
+ return true;
+ }
+ keys.add(key);
+ }
+ return false;
+ });
+ return new Matches(List.copyOf(keys), truncated);
+ }
+
+ private StatObjectResponse stat(String key) {
+ try {
+ return client.statObject(StatObjectArgs.builder().bucket(bucket).object(prefix + key).build());
+ } catch (ErrorResponseException ex) {
+ if ("NoSuchKey".equals(ex.errorResponse().code()) || "NoSuchObject".equals(ex.errorResponse().code())) {
+ return null;
+ }
+ throw failure("read " + key, ex);
+ } catch (Exception ex) {
+ throw failure("read " + key, ex);
+ }
+ }
+
+ private List downloadAll(List keys) {
+ List files = new ArrayList<>();
+ for (String key : keys) {
+ byte[] bytes = bytesOf(key);
+ try {
+ String name = key.substring(key.lastIndexOf('/') + 1);
+ Path file = Files.createTempDirectory("storage-").resolve(name.isBlank() ? "object" : name);
+ Files.write(file, bytes);
+ files.add(file.toFile());
+ } catch (IOException ex) {
+ throw new StorageException(StorageException.UNREACHABLE, "Could not keep " + key + " as a file: " + ex.getMessage(), ex);
+ }
+ }
+ return files;
+ }
+
+ private List textsOf(List keys) {
+ return keys.stream().map(key -> new String(bytesOf(key), StandardCharsets.UTF_8)).toList();
+ }
+
+ private ArrayNode jsonOf(List keys) {
+ ArrayNode array = objectMapper.createArrayNode();
+ for (String key : keys) {
+ try {
+ array.add(objectMapper.readTree(bytesOf(key)));
+ } catch (StorageException ex) {
+ throw ex;
+ } catch (Exception ex) {
+ throw new StorageException(StorageException.CONTENT_INVALID,
+ key + " is not JSON, so it cannot be read as JSON; read it as TEXTS or FILES instead");
+ }
+ }
+ return array;
+ }
+
+ private byte[] bytesOf(String key) {
+ StatObjectResponse stat = stat(key);
+ if (stat == null) {
+ throw new StorageException(StorageException.NOT_FOUND, key + " is not in " + connection.name());
+ }
+ if (stat.size() > settings.maxObjectBytes()) {
+ throw new StorageException(StorageException.TOO_LARGE, key + " is " + stat.size()
+ + " bytes, more than the " + settings.maxObjectBytes() + " a read may take");
+ }
+ return call("read " + key, () -> {
+ try (InputStream stream = client.getObject(GetObjectArgs.builder().bucket(bucket).object(prefix + key).build())) {
+ return stream.readAllBytes();
+ }
+ });
+ }
+
+ private record Content(byte[] bytes, File file, String contentType) {
+ long size() throws IOException {
+ return file != null ? Files.size(file.toPath()) : bytes.length;
+ }
+
+ InputStream open() throws IOException {
+ return file != null ? Files.newInputStream(file.toPath()) : new ByteArrayInputStream(bytes);
+ }
+ }
+
+ private Content contentOf(Object value, String contentType) {
+ String declared = contentType == null || contentType.isBlank() ? null : contentType.trim();
+ if (value instanceof File file) {
+ if (!file.isFile()) {
+ throw new StorageException(StorageException.NOT_FOUND, "The file to write, " + file.getName() + ", is not there");
+ }
+ String type = declared;
+ if (type == null) {
+ try {
+ type = Files.probeContentType(file.toPath());
+ } catch (IOException ignored) {
+ type = null;
+ }
+ }
+ return new Content(null, file, type == null ? "application/octet-stream" : type);
+ }
+ if (value instanceof JsonNode || value instanceof Map, ?> || value instanceof Collection>) {
+ byte[] json = objectMapper.writeValueAsBytes(value);
+ return new Content(json, null, declared == null ? "application/json" : declared);
+ }
+ String text = value == null ? "" : String.valueOf(value);
+ return new Content(text.getBytes(StandardCharsets.UTF_8), null, declared == null ? "text/plain; charset=utf-8" : declared);
+ }
+
+ @FunctionalInterface
+ private interface S3Call {
+ T run() throws Exception;
+ }
+
+ private T call(String what, S3Call call) {
+ try {
+ return call.run();
+ } catch (StorageException ex) {
+ throw ex;
+ } catch (Exception ex) {
+ throw failure(what, ex);
+ }
+ }
+
+ private StorageException failure(String what, Exception ex) {
+ if (ex instanceof ErrorResponseException response) {
+ String code = response.errorResponse().code();
+ String message = response.errorResponse().message();
+ if ("NoSuchBucket".equals(code)) {
+ return new StorageException(StorageException.NOT_FOUND, "The bucket " + bucket + " of " + connection.name() + " does not exist", ex);
+ }
+ if ("AccessDenied".equals(code) || "InvalidAccessKeyId".equals(code) || "SignatureDoesNotMatch".equals(code)) {
+ return new StorageException(StorageException.OPERATION_NOT_ALLOWED,
+ connection.name() + " refused to " + what + ": " + code + (message == null ? "" : " - " + message), ex);
+ }
+ return new StorageException(StorageException.UNREACHABLE,
+ connection.name() + " could not " + what + ": " + code + (message == null ? "" : " - " + message), ex);
+ }
+ return new StorageException(StorageException.UNREACHABLE,
+ "Could not " + what + " on " + connection.name() + ": " + ex.getMessage(), ex);
+ }
+}
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
new file mode 100644
index 0000000..c30de23
--- /dev/null
+++ b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageType.java
@@ -0,0 +1,185 @@
+// 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.s3;
+
+import java.net.URI;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.TimeUnit;
+
+import org.springframework.stereotype.Component;
+
+import io.minio.MinioClient;
+import it.cnr.isti.workflow.manager.blocks.IOCapability;
+import it.cnr.isti.workflow.manager.blocks.IOCapabilityType;
+import it.cnr.isti.workflow.manager.storage.CredentialField;
+import it.cnr.isti.workflow.manager.storage.StorageConnection;
+import it.cnr.isti.workflow.manager.storage.StorageEndpointGuard;
+import it.cnr.isti.workflow.manager.storage.StorageException;
+import it.cnr.isti.workflow.manager.storage.StorageInit;
+import it.cnr.isti.workflow.manager.storage.StoragePlaceholders;
+import it.cnr.isti.workflow.manager.storage.StorageScope;
+import it.cnr.isti.workflow.manager.storage.StorageSession;
+import it.cnr.isti.workflow.manager.storage.StorageSettings;
+import it.cnr.isti.workflow.manager.storage.StorageType;
+import it.cnr.isti.workflow.manager.storage.ViewShape;
+import okhttp3.OkHttpClient;
+import tools.jackson.databind.ObjectMapper;
+
+/**
+ * An S3-compatible object store: MinIO, AWS S3, Ceph and the rest.
+ *
+ *
+ * - A template is an object key, placeholders rendered in: {@code reports/${{context.iteration}}.json}.
+ *
- A query is a key or a glob - {@code pippo*}, {@code reports/**}{@code /*.json} - and a read
+ * returns every object it matches, in the shape asked.
+ *
- Init can create the bucket when it is missing; there is no script.
+ *
+ */
+@Component
+public class S3StorageType implements StorageType {
+
+ public static final String NAME = "S3";
+
+ private final StorageSettings settings;
+ private final StorageEndpointGuard guard;
+ private final ObjectMapper objectMapper;
+ private final OkHttpClient httpClient;
+
+ public S3StorageType(StorageSettings settings, StorageEndpointGuard guard, ObjectMapper objectMapper) {
+ this.settings = settings;
+ this.guard = guard;
+ this.objectMapper = objectMapper;
+ // One client, one connection pool, for every session: a MinioClient is cheap to build on it.
+ this.httpClient = new OkHttpClient.Builder()
+ .connectTimeout(settings.connectTimeoutSeconds(), TimeUnit.SECONDS)
+ .readTimeout(Math.max(60, settings.connectTimeoutSeconds()), TimeUnit.SECONDS)
+ .writeTimeout(Math.max(60, settings.connectTimeoutSeconds()), TimeUnit.SECONDS)
+ .build();
+ }
+
+ @Override
+ public String getName() {
+ return NAME;
+ }
+
+ @Override
+ public String getDescription() {
+ return "An S3-compatible object store (MinIO, AWS S3, Ceph). Writes put an object under a key; "
+ + "reads take a key or a glob such as reports/*.json.";
+ }
+
+ @Override
+ public List credentialFields() {
+ return List.of(
+ new CredentialField("endpoint", "Endpoint URL", "e.g. https://minio.example.org", false, true),
+ new CredentialField("region", "Region", "Leave empty unless the store asks for one", false, false),
+ new CredentialField("bucket", "Bucket", "The bucket this connection works in", false, true),
+ new CredentialField("rootPrefix", "Root prefix", "Optional folder every key goes under", false, false),
+ new CredentialField("accessKey", "Access key", null, false, true),
+ new CredentialField("secretKey", "Secret key", null, true, true));
+ }
+
+ @Override
+ public List viewShapes() {
+ return List.of(ViewShape.FILES, ViewShape.TEXTS, ViewShape.JSON, ViewShape.KEYS);
+ }
+
+ @Override
+ public List writableKinds() {
+ return List.of(new IOCapability(IOCapabilityType.FILE, false), new IOCapability(IOCapabilityType.TEXT, false),
+ new IOCapability(IOCapabilityType.JSON, false));
+ }
+
+ @Override
+ public void validateConnection(StorageConnection connection) {
+ String endpoint = connection.requireText("endpoint");
+ try {
+ URI uri = new URI(endpoint);
+ if (!"http".equalsIgnoreCase(uri.getScheme()) && !"https".equalsIgnoreCase(uri.getScheme())
+ || uri.getHost() == null) {
+ throw new StorageException(StorageException.CONNECTION_INVALID,
+ connection.name() + ": the endpoint must be an http(s) URL, e.g. https://minio.example.org");
+ }
+ } catch (java.net.URISyntaxException ex) {
+ throw new StorageException(StorageException.CONNECTION_INVALID,
+ connection.name() + ": the endpoint is not a URL: " + endpoint);
+ }
+ String bucket = connection.requireText("bucket");
+ if (!S3Keys.isValidBucket(bucket)) {
+ throw new StorageException(StorageException.CONNECTION_INVALID,
+ connection.name() + ": " + bucket + " is not a valid bucket name (3-63 lowercase letters, digits, '.' or '-')");
+ }
+ connection.requireText("accessKey");
+ connection.requireText("secretKey");
+ String rootPrefix = connection.optionalText("rootPrefix");
+ if (rootPrefix != null && !S3Keys.problemsAsWritten(rootPrefix, "The root prefix").isEmpty()) {
+ throw new StorageException(StorageException.CONNECTION_INVALID,
+ connection.name() + ": " + String.join("; ", S3Keys.problemsAsWritten(rootPrefix, "the root prefix")));
+ }
+ guard.check(connection, "endpoint", null);
+ }
+
+ @Override
+ public List validateTemplate(String template) {
+ List problems = S3Keys.problemsAsWritten(template, "The object key");
+ if (template != null && S3Keys.hasWildcard(template)) {
+ problems.add("The object key cannot contain '*' or '?': it names one object");
+ }
+ if (template != null && StoragePlaceholders.names(template).contains(StoragePlaceholders.VALUE)) {
+ problems.add("${{value}} is the whole content being written and cannot be part of its key; use ${{value.}} for a field of a JSON value");
+ }
+ return problems;
+ }
+
+ @Override
+ public List validateQuery(String query, ViewShape shape) {
+ List problems = S3Keys.problemsAsWritten(query, "The key or glob");
+ if (shape != null && !viewShapes().contains(shape)) {
+ problems.add("An S3 read cannot return " + shape);
+ }
+ if (query != null && StoragePlaceholders.names(query).stream().anyMatch(StoragePlaceholders::isValueName)) {
+ problems.add("A read has no value: ${{value}} only means something when writing");
+ }
+ return problems;
+ }
+
+ @Override
+ public List validateDelete(String template) {
+ List problems = S3Keys.problemsAsWritten(template, "The object key");
+ if (template != null && S3Keys.hasWildcard(template)) {
+ problems.add("A delete names one object: '*' and '?' are not allowed");
+ }
+ return problems;
+ }
+
+ @Override
+ public List validateInit(StorageInit init) {
+ List problems = new ArrayList<>();
+ if (init != null && init.hasScript()) {
+ problems.add("An S3 storage has no init script; leave it empty");
+ }
+ return problems;
+ }
+
+ @Override
+ public StorageSession open(StorageConnection connection, StorageScope scope) {
+ validateConnection(connection);
+ MinioClient.Builder builder = MinioClient.builder()
+ .endpoint(connection.requireText("endpoint"))
+ .credentials(connection.requireText("accessKey"), connection.requireText("secretKey"))
+ .httpClient(httpClient);
+ String region = connection.optionalText("region");
+ if (region != null) {
+ builder.region(region);
+ }
+ MinioClient client = builder.build();
+ if (connection.flag("pathStyle", false)) {
+ client.disableVirtualStyleEndpoint();
+ }
+ String prefix = S3Keys.normalizePrefix(connection.optionalText("rootPrefix")) + S3Keys.normalizePrefix(scope.prefix());
+ return new S3Session(client, connection, connection.requireText("bucket"), prefix, settings, objectMapper);
+ }
+}
diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties
index a74e9cd..9de0a20 100644
--- a/src/main/resources/application.properties
+++ b/src/main/resources/application.properties
@@ -160,3 +160,16 @@ logging.level.root=ERROR
# Jackson 3 changed default field visibility from ANY (Jackson 2) to PUBLIC_ONLY.
# Restore to ANY to match Jackson 2 behavior for package-private fields in domain model.
spring.jackson.visibility.field=any
+
+# Storage node and block. The catalog of storages offered to every flow, with their credentials -
+# keep the file outside the image. Empty uses classpath:storages.json, which lists none.
+app.storage.instances.file=${STORAGE_INSTANCES_FILE:}
+# A person's own storage connection (from the vault) is checked like an LLM endpoint, behind its own
+# switch: private and local addresses are refused unless allowed here.
+app.storage.endpoint.allow-private-network=${STORAGE_ENDPOINT_ALLOW_PRIVATE_NETWORK:false}
+app.storage.endpoint.allowed-host-patterns=${STORAGE_ENDPOINT_ALLOWED_HOST_PATTERNS:}
+app.storage.max-object-bytes=${STORAGE_MAX_OBJECT_BYTES:33554432}
+app.storage.max-keys=${STORAGE_MAX_KEYS:1000}
+app.storage.sql.statement-timeout-seconds=${STORAGE_SQL_STATEMENT_TIMEOUT_SECONDS:30}
+app.storage.sql.max-rows=${STORAGE_SQL_MAX_ROWS:1000}
+app.storage.connect-timeout-seconds=${STORAGE_CONNECT_TIMEOUT_SECONDS:10}
diff --git a/src/main/resources/storages.json b/src/main/resources/storages.json
new file mode 100644
index 0000000..93044c0
--- /dev/null
+++ b/src/main/resources/storages.json
@@ -0,0 +1 @@
+{ "storages": [] }
diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProviderTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProviderTest.java
new file mode 100644
index 0000000..688688d
--- /dev/null
+++ b/src/test/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProviderTest.java
@@ -0,0 +1,82 @@
+// 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.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.nio.file.Files;
+import java.nio.file.Path;
+import java.util.List;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+import org.springframework.core.io.DefaultResourceLoader;
+
+import it.cnr.isti.workflow.manager.storage.postgres.PostgresStorageType;
+import it.cnr.isti.workflow.manager.storage.s3.S3StorageType;
+import tools.jackson.databind.ObjectMapper;
+
+/** The operator's catalog: loaded whole, or the service does not start. */
+class StorageInstancesProviderTest {
+
+ @TempDir
+ Path folder;
+
+ private final ObjectMapper mapper = new ObjectMapper();
+ private final StorageEndpointGuard guard = new StorageEndpointGuard(false, "");
+ private final StorageTypes types = new StorageTypes(List.of(
+ new S3StorageType(StorageSettings.defaults(), guard, mapper),
+ new PostgresStorageType(StorageSettings.defaults(), guard, mapper)));
+
+ @Test
+ void withoutAFileTheCatalogIsEmpty() {
+ assertTrue(new StorageInstancesProvider(mapper, new DefaultResourceLoader(), types, "").all().isEmpty());
+ }
+
+ @Test
+ void aValidCatalogIsLoadedAndAnOperatorsPrivateAddressIsFine() throws Exception {
+ StorageInstancesProvider provider = load("""
+ { "storages": [
+ { "id": "flow-files", "name": "Flow files", "type": "S3", "allowInit": true,
+ "connection": { "endpoint": "http://127.0.0.1:9000", "bucket": "flow-files",
+ "accessKey": "a", "secretKey": "s" } },
+ { "id": "reports-db", "type": "postgresql", "readOnly": true,
+ "connection": { "host": "10.0.0.5", "database": "reports", "user": "u", "password": "p" } }
+ ] }""");
+
+ assertEquals(List.of("flow-files", "reports-db"), provider.all().stream().map(StorageInstancesProvider.StorageInstance::id).toList());
+ assertEquals(1, provider.ofType("PostgreSQL").size());
+ assertTrue(provider.find("reports-db").orElseThrow().readOnly());
+ assertFalse(provider.find("missing").isPresent());
+ }
+
+ @Test
+ void anEntryThatDoesNotHoldTogetherStopsTheServiceSayingWhich() {
+ IllegalStateException unknownType = assertThrows(IllegalStateException.class, () -> load("""
+ { "storages": [ { "id": "x", "type": "Mongo", "connection": {} } ] }"""));
+ assertTrue(unknownType.getMessage().contains("(x) has unknown type Mongo"), unknownType.getMessage());
+
+ IllegalStateException badBucket = assertThrows(IllegalStateException.class, () -> load("""
+ { "storages": [ { "id": "files", "type": "S3", "connection":
+ { "endpoint": "http://minio:9000", "bucket": "Bad_Bucket", "accessKey": "a", "secretKey": "s" } } ] }"""));
+ assertTrue(badBucket.getMessage().contains("(files)") && badBucket.getMessage().contains("not a valid bucket name"),
+ badBucket.getMessage());
+
+ IllegalStateException twice = assertThrows(IllegalStateException.class, () -> load("""
+ { "storages": [
+ { "id": "db", "type": "PostgreSQL", "connection": { "host": "db", "database": "d", "user": "u", "password": "p" } },
+ { "id": "db", "type": "PostgreSQL", "connection": { "host": "db", "database": "d", "user": "u", "password": "p" } } ] }"""));
+ assertTrue(twice.getMessage().contains("id db is used twice"), twice.getMessage());
+ }
+
+ private StorageInstancesProvider load(String json) throws Exception {
+ Path file = folder.resolve("storages.json");
+ Files.writeString(file, json);
+ return new StorageInstancesProvider(mapper, new DefaultResourceLoader(), types, file.toString());
+ }
+}
diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageTypeTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageTypeTest.java
new file mode 100644
index 0000000..c3f8c7c
--- /dev/null
+++ b/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageTypeTest.java
@@ -0,0 +1,198 @@
+// 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.postgres;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.util.List;
+import java.util.Map;
+import java.util.function.Function;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.testcontainers.junit.jupiter.Container;
+import org.testcontainers.junit.jupiter.Testcontainers;
+import org.testcontainers.postgresql.PostgreSQLContainer;
+
+import it.cnr.isti.workflow.manager.storage.StorageConnection;
+import it.cnr.isti.workflow.manager.storage.StorageEndpointGuard;
+import it.cnr.isti.workflow.manager.storage.StorageException;
+import it.cnr.isti.workflow.manager.storage.StorageInit;
+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.StorageSettings;
+import it.cnr.isti.workflow.manager.storage.StorageWriteResult;
+import it.cnr.isti.workflow.manager.storage.ViewShape;
+import tools.jackson.databind.JsonNode;
+import tools.jackson.databind.ObjectMapper;
+
+/** Against a real PostgreSQL: values bound, reads kept read-only, and refusals said plainly. */
+@Testcontainers(disabledWithoutDocker = true)
+class PostgresStorageTypeTest {
+
+ @Container
+ static final PostgreSQLContainer POSTGRES = new PostgreSQLContainer("postgres:17-alpine");
+
+ private static final ObjectMapper MAPPER = new ObjectMapper();
+ private static PostgresStorageType type;
+
+ @BeforeAll
+ static void createType() {
+ type = new PostgresStorageType(StorageSettings.defaults(), new StorageEndpointGuard(false, ""), MAPPER);
+ try (StorageSession session = type.open(connection(false, true), StorageScope.shared())) {
+ session.initialize(new StorageInit(false, """
+ CREATE TABLE IF NOT EXISTS rounds (id serial PRIMARY KEY, brief text NOT NULL, verdict text, score int, data jsonb);
+ CREATE TABLE IF NOT EXISTS notes (id serial PRIMARY KEY, body text);"""));
+ }
+ }
+
+ @AfterAll
+ static void closePools() {
+ type.destroy();
+ }
+
+ @Test
+ void aWriteBindsTheFieldsOfItsValueAndAReadReturnsRowsAsJson() {
+ JsonNode round = MAPPER.readTree("""
+ {"brief": "Build a list", "verdict": "rejected", "score": 3, "extra": {"k": [1, 2]}}""");
+ try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) {
+ StorageWriteResult written = session.write(
+ "INSERT INTO rounds(brief, verdict, score, data) VALUES (${{value.brief}}, ${{value.verdict}}, ${{value.score}}, ${{value.extra}}) RETURNING id",
+ round, null, name -> null);
+ assertEquals(1, written.affected());
+ assertTrue(written.rows().get(0).get("id").isNumber());
+
+ StorageReadResult read = session.read("SELECT brief, score, data FROM rounds WHERE verdict = ${{verdict}} ORDER BY id",
+ ViewShape.JSON, values(Map.of("verdict", "rejected")));
+ JsonNode rows = (JsonNode) read.value();
+ assertEquals("Build a list", rows.get(0).get("brief").asString());
+ assertEquals(3, rows.get(0).get("score").asInt());
+ assertEquals(2, rows.get(0).get("data").get("k").get(1).asInt(), "jsonb comes back as JSON, not a string");
+ }
+ }
+
+ @Test
+ void aValueThatLooksLikeSqlIsStoredAsTheCharactersItIs() {
+ String hostile = "x'); DROP TABLE notes; --";
+ try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) {
+ session.write("INSERT INTO notes(body) VALUES (${{value}})", hostile, null, name -> null);
+
+ JsonNode rows = (JsonNode) session.read("SELECT body FROM notes WHERE body = ${{body}}", ViewShape.JSON,
+ values(Map.of("body", hostile))).value();
+ assertEquals(hostile, rows.get(0).get("body").asString());
+ assertTrue(((List>) session.list(null, name -> null).value()).contains("notes"), "the table is still there");
+ }
+ }
+
+ @Test
+ void aReadThatTriesToChangeDataIsStoppedByTheReadOnlyTransaction() {
+ try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) {
+ StorageException refusal = assertThrows(StorageException.class, () -> session.read(
+ "WITH gone AS (DELETE FROM notes RETURNING *) SELECT count(*) FROM gone", ViewShape.JSON, name -> null));
+ assertEquals(StorageException.OPERATION_NOT_ALLOWED, refusal.getErrorCode());
+ assertTrue(refusal.getMessage().contains("A read cannot change data"), refusal.getMessage());
+ }
+ }
+
+ @Test
+ void theVerbIsCheckedAtRunTimeTooAndTheServersReasonIsKept() {
+ try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) {
+ StorageException notASelect = assertThrows(StorageException.class,
+ () -> session.read("DELETE FROM notes", ViewShape.JSON, name -> null));
+ assertEquals(StorageException.SQL_REJECTED, notASelect.getErrorCode());
+
+ StorageException missingTable = assertThrows(StorageException.class,
+ () -> session.read("SELECT * FROM nowhere", ViewShape.JSON, name -> null));
+ assertEquals(StorageException.SQL_REJECTED, missingTable.getErrorCode());
+ assertTrue(missingTable.getMessage().contains("relation \"nowhere\" does not exist"), missingTable.getMessage());
+
+ StorageException twoStatements = assertThrows(StorageException.class,
+ () -> session.read("SELECT 1; DROP TABLE notes", ViewShape.JSON, name -> null));
+ assertEquals(StorageException.SQL_REJECTED, twoStatements.getErrorCode());
+ assertTrue(twoStatements.getMessage().startsWith("Only one statement is allowed"), twoStatements.getMessage());
+
+ // The driver would have run both: it is the check before it that stops the second.
+ StorageException smuggled = assertThrows(StorageException.class, () -> session.write(
+ "INSERT INTO notes(body) VALUES (${{value}}); DROP TABLE notes", "x", null, name -> null));
+ assertEquals(StorageException.SQL_REJECTED, smuggled.getErrorCode());
+ assertTrue(((List>) session.list(null, name -> null).value()).contains("notes"));
+ }
+ }
+
+ @Test
+ void aDeleteReturnsHowManyRowsWentAndAReadOnlyStorageRefusesToChange() {
+ try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) {
+ session.write("INSERT INTO notes(body) VALUES (${{value}})", "to delete", null, name -> null);
+ StorageWriteResult deleted = session.delete("DELETE FROM notes WHERE body = ${{body}}",
+ values(Map.of("body", "to delete")));
+ assertEquals(1, deleted.affected());
+ assertNull(deleted.rows());
+ }
+ try (StorageSession readOnly = type.open(connection(true, false), StorageScope.shared())) {
+ StorageException refusal = assertThrows(StorageException.class,
+ () -> readOnly.write("INSERT INTO notes(body) VALUES (${{value}})", "no", null, name -> null));
+ assertEquals(StorageException.OPERATION_NOT_ALLOWED, refusal.getErrorCode());
+ }
+ }
+
+ @Test
+ void aMissingParameterIsNamed() {
+ try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) {
+ StorageException refusal = assertThrows(StorageException.class,
+ () -> session.read("SELECT * FROM notes WHERE body = ${{body}}", ViewShape.JSON, name -> null));
+ assertEquals(StorageException.PLACEHOLDER_MISSING, refusal.getErrorCode());
+ assertTrue(refusal.getMessage().contains("${{body}}"), refusal.getMessage());
+ }
+ }
+
+ @Test
+ void aPersonsOwnConnectionToAPrivateAddressIsRefusedAnOperatorsIsNot() {
+ StorageConnection personal = new StorageConnection("your PostgreSQL connection", PostgresStorageType.NAME,
+ connection(false, false).settings(), false, true, false);
+ StorageException refusal = assertThrows(StorageException.class, () -> type.open(personal, StorageScope.shared()));
+ assertEquals(StorageException.CONNECTION_INVALID, refusal.getErrorCode());
+ assertTrue(refusal.getMessage().contains("private, local or otherwise disallowed"), refusal.getMessage());
+ }
+
+ @Test
+ void anInitScriptNeedsAStorageThatMayBePrepared() {
+ try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) {
+ StorageException refusal = assertThrows(StorageException.class,
+ () -> session.initialize(new StorageInit(false, "CREATE TABLE x (id int)")));
+ assertEquals(StorageException.INIT_FAILED, refusal.getErrorCode());
+ }
+ }
+
+ @Test
+ void editorChecksSayWhatIsWrong() {
+ assertTrue(type.validateQuery("INSERT INTO t VALUES (1)", ViewShape.JSON).getFirst().contains("must start with"));
+ assertTrue(type.validateQuery("SELECT 1", ViewShape.FILES).getFirst().contains("cannot return FILES"));
+ assertTrue(type.validateTemplate("UPDATE t SET a = '${{value}}'").getFirst().contains("inside a string literal"));
+ assertTrue(type.validateTemplate("INSERT INTO t(a) VALUES (${{value.a}})").isEmpty());
+ assertTrue(type.validateDelete("DELETE FROM t WHERE id = ${{value}}").getFirst().contains("has no value"));
+ assertTrue(type.validateInit(new StorageInit(false, "CREATE TABLE ${{t}} (id int)")).getFirst().contains("cannot have placeholders"));
+ }
+
+ private static Function values(Map values) {
+ return values::get;
+ }
+
+ private static StorageConnection connection(boolean readOnly, boolean allowInit) {
+ JsonNode settings = MAPPER.valueToTree(Map.of(
+ "host", POSTGRES.getHost(),
+ "port", String.valueOf(POSTGRES.getFirstMappedPort()),
+ "database", POSTGRES.getDatabaseName(),
+ "user", POSTGRES.getUsername(),
+ "password", POSTGRES.getPassword(),
+ "sslMode", "disable"));
+ return new StorageConnection("test database" + (readOnly ? " (read-only)" : "") + (allowInit ? " (init)" : ""),
+ PostgresStorageType.NAME, settings, readOnly, allowInit, true);
+ }
+}
diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplateTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplateTest.java
new file mode 100644
index 0000000..b8b5f9e
--- /dev/null
+++ b/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplateTest.java
@@ -0,0 +1,81 @@
+// 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.postgres;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.util.List;
+
+import org.junit.jupiter.api.Test;
+
+/** Placeholders become bound parameters, and only where a value can stand. */
+class SqlTemplateTest {
+
+ @Test
+ void placeholdersBecomeParametersInOrderAndAUsedTwiceNameIsBoundTwice() {
+ SqlTemplate template = SqlTemplate.parse(
+ "INSERT INTO rounds(brief, verdict, again) VALUES (${{value.brief}}, ${{ value.verdict }}, ${{value.brief}})");
+
+ assertEquals("INSERT INTO rounds(brief, verdict, again) VALUES (?, ?, ?)", template.sql());
+ assertEquals(List.of("value.brief", "value.verdict", "value.brief"), template.parameters());
+ assertEquals("INSERT", template.verb());
+ assertTrue(template.problems().isEmpty());
+ }
+
+ @Test
+ void aQuestionMarkTheAuthorWroteIsNotTakenForAParameter() {
+ SqlTemplate template = SqlTemplate.parse("SELECT * FROM docs WHERE data ? 'tag' AND id = ${{id}}");
+
+ assertEquals("SELECT * FROM docs WHERE data ?? 'tag' AND id = ?", template.sql());
+ assertEquals(List.of("id"), template.parameters());
+ }
+
+ @Test
+ void aPlaceholderInsideAStringLiteralIsRefusedWithTheFix() {
+ SqlTemplate template = SqlTemplate.parse("SELECT * FROM t WHERE name = '${{name}}'");
+
+ assertEquals(1, template.problems().size());
+ String problem = template.problems().getFirst();
+ assertTrue(problem.contains("${{name}} is inside a string literal"), problem);
+ assertTrue(problem.contains("WHERE name = ${{name}}"), problem);
+ }
+
+ @Test
+ void placeholdersInCommentsQuotedIdentifiersAndDollarBodiesAreRefusedToo() {
+ assertTrue(SqlTemplate.parse("SELECT 1 -- ${{x}}").problems().getFirst().contains("inside a comment"));
+ assertTrue(SqlTemplate.parse("SELECT 1 /* a /* nested ${{x}} */ */").problems().getFirst().contains("inside a comment"));
+ assertTrue(SqlTemplate.parse("SELECT \"${{x}}\" FROM t").problems().getFirst().contains("quoted identifier"));
+ assertTrue(SqlTemplate.parse("SELECT $body$ ${{x}} $body$").problems().getFirst().contains("dollar-quoted"));
+ }
+
+ @Test
+ void quotesAreSkippedWithTheirEscapesSoWhatFollowsIsStillSeen() {
+ SqlTemplate template = SqlTemplate.parse("SELECT 'it''s', E'a\\'b' FROM t WHERE id = ${{id}}");
+
+ assertTrue(template.problems().isEmpty(), template.problems().toString());
+ assertEquals(List.of("id"), template.parameters());
+ assertEquals("SELECT 'it''s', E'a\\'b' FROM t WHERE id = ?", template.sql());
+ }
+
+ @Test
+ void theVerbIsTheFirstKeywordPastCommentsAndParentheses() {
+ assertEquals("SELECT", SqlTemplate.parse(" -- rows\n /* all */ (SELECT 1)").verb());
+ assertEquals("WITH", SqlTemplate.parse("with x as (select 1) select * from x").verb());
+ assertEquals("", SqlTemplate.parse(" ").verb());
+ }
+
+ @Test
+ void oneStatementOnlyThoughAClosingSemicolonIsFine() {
+ SqlTemplate closed = SqlTemplate.parse("SELECT * FROM t WHERE id = ${{id}}; -- done\n");
+ assertTrue(closed.problems().isEmpty(), closed.problems().toString());
+ assertEquals("SELECT * FROM t WHERE id = ?", closed.sql().trim());
+
+ SqlTemplate two = SqlTemplate.parse("INSERT INTO t VALUES (${{value}}); DROP TABLE t");
+ assertTrue(two.problems().getFirst().startsWith("Only one statement is allowed"), two.problems().toString());
+
+ assertTrue(SqlTemplate.parse("SELECT ';' AS semicolon").problems().isEmpty(), "a ';' in a literal is text");
+ }
+}
diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3KeysTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3KeysTest.java
new file mode 100644
index 0000000..5bd741a
--- /dev/null
+++ b/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3KeysTest.java
@@ -0,0 +1,67 @@
+// 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.s3;
+
+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.util.regex.Pattern;
+
+import org.junit.jupiter.api.Test;
+
+import it.cnr.isti.workflow.manager.storage.StorageException;
+
+class S3KeysTest {
+
+ @Test
+ void aSingleStarStaysInsideOneFolderAndADoubleStarCrossesThem() {
+ Pattern pippo = S3Keys.globToRegex("pippo*");
+ assertTrue(pippo.matcher("pippo.txt").matches());
+ assertTrue(pippo.matcher("pippo").matches());
+ assertFalse(pippo.matcher("pippo/inside.txt").matches());
+ assertFalse(pippo.matcher("a-pippo.txt").matches());
+
+ Pattern deep = S3Keys.globToRegex("reports/**.json");
+ assertTrue(deep.matcher("reports/2026/09/r.json").matches());
+ assertFalse(deep.matcher("reports/r.txt").matches());
+
+ assertTrue(S3Keys.globToRegex("r?.md").matcher("r1.md").matches());
+ assertFalse(S3Keys.globToRegex("r?.md").matcher("r/.md").matches());
+ assertTrue(S3Keys.globToRegex("a+b(1).txt").matcher("a+b(1).txt").matches(), "regex characters are literal");
+ }
+
+ @Test
+ void aGlobIsListedFromItsFixedPart() {
+ assertEquals("reports/2026-", S3Keys.listingPrefix("reports/2026-*.json"));
+ assertEquals("", S3Keys.listingPrefix("*.json"));
+ assertEquals("exact.txt", S3Keys.listingPrefix("exact.txt"));
+ }
+
+ @Test
+ void aKeyThatCouldLeaveItsPrefixIsRefusedAfterPlaceholdersAreIn() {
+ assertEquals("notes/1.txt", S3Keys.requireSafeKey("notes/1.txt"));
+ for (String key : new String[] { "../other/secret", "notes/../../x", "/absolute", "a//b", "" }) {
+ StorageException refusal = assertThrows(StorageException.class, () -> S3Keys.requireSafeKey(key), key);
+ assertEquals(StorageException.INVALID_LOCATION, refusal.getErrorCode());
+ }
+ }
+
+ @Test
+ void theKeyAsWrittenIsCheckedWithItsPlaceholdersStandingInForValues() {
+ assertTrue(S3Keys.problemsAsWritten("reports/${{context.iteration}}.json", "The key").isEmpty());
+ assertFalse(S3Keys.problemsAsWritten("../${{x}}", "The key").isEmpty());
+ assertFalse(S3Keys.problemsAsWritten("r/${{1bad}}", "The key").isEmpty());
+ }
+
+ @Test
+ void bucketNamesFollowTheS3Rules() {
+ assertTrue(S3Keys.isValidBucket("flow-data"));
+ assertFalse(S3Keys.isValidBucket("Flow_Data"));
+ assertFalse(S3Keys.isValidBucket("ab"));
+ assertFalse(S3Keys.isValidBucket("a..b"));
+ }
+}
diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageTypeTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageTypeTest.java
new file mode 100644
index 0000000..1400740
--- /dev/null
+++ b/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageTypeTest.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.storage.s3;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.io.File;
+import java.nio.file.Files;
+import java.util.List;
+import java.util.Map;
+
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.testcontainers.containers.MinIOContainer;
+import org.testcontainers.junit.jupiter.Container;
+import org.testcontainers.junit.jupiter.Testcontainers;
+
+import it.cnr.isti.workflow.manager.storage.StorageConnection;
+import it.cnr.isti.workflow.manager.storage.StorageEndpointGuard;
+import it.cnr.isti.workflow.manager.storage.StorageException;
+import it.cnr.isti.workflow.manager.storage.StorageInit;
+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.StorageSettings;
+import it.cnr.isti.workflow.manager.storage.ViewShape;
+import tools.jackson.databind.JsonNode;
+import tools.jackson.databind.ObjectMapper;
+
+/** Against a real MinIO: keys, globs, scopes and every shape a read can take. */
+@Testcontainers(disabledWithoutDocker = true)
+class S3StorageTypeTest {
+
+ @Container
+ static final MinIOContainer MINIO = new MinIOContainer("minio/minio:latest");
+
+ private static final ObjectMapper MAPPER = new ObjectMapper();
+ private static S3StorageType type;
+
+ @BeforeAll
+ static void createBucket() {
+ type = new S3StorageType(StorageSettings.defaults(), new StorageEndpointGuard(false, ""), MAPPER);
+ try (StorageSession session = type.open(connection("flow-data", false, true), StorageScope.shared())) {
+ session.initialize(new StorageInit(true, null));
+ }
+ }
+
+ @Test
+ void whatIsWrittenCanBeReadBackAsFilesTextsJsonOrKeys() throws Exception {
+ File report = Files.createTempFile("report", ".md").toFile();
+ Files.writeString(report.toPath(), "# Round 1");
+ try (StorageSession session = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-shapes"))) {
+ assertEquals("pippo-1.md", session.write("pippo-${{context.iteration}}.md", report, null,
+ Map.of("context.iteration", 1)::get).location());
+ session.write("pippo-2.json", MAPPER.readTree("{\"round\": 2}"), null, name -> null);
+ session.write("notes/pippo-3.txt", "not at the top level", null, name -> null);
+ session.write("other.txt", "not pippo", null, name -> null);
+
+ StorageReadResult keys = session.read("pippo*", ViewShape.KEYS, name -> null);
+ assertEquals(List.of("pippo-1.md", "pippo-2.json"), keys.value(), "one * stays in its folder");
+
+ StorageReadResult deep = session.read("**pippo*", ViewShape.KEYS, name -> null);
+ assertEquals(3, deep.count());
+
+ @SuppressWarnings("unchecked")
+ List files = (List) session.read("pippo-1.md", ViewShape.FILES, name -> null).value();
+ assertEquals("pippo-1.md", files.getFirst().getName());
+ assertEquals("# Round 1", Files.readString(files.getFirst().toPath()));
+
+ assertEquals(List.of("not at the top level"), session.read("notes/*", ViewShape.TEXTS, name -> null).value());
+
+ JsonNode json = (JsonNode) session.read("pippo-2.json", ViewShape.JSON, name -> null).value();
+ assertEquals(2, json.get(0).get("round").asInt());
+ }
+ }
+
+ @Test
+ void aScopePerExecutionSeesOnlyWhatThatExecutionWrote() {
+ try (StorageSession first = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-a"))) {
+ first.write("brief.txt", "from a", null, name -> null);
+ }
+ try (StorageSession second = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-b"))) {
+ assertEquals(0, second.read("brief.txt", ViewShape.TEXTS, name -> null).count());
+ second.write("brief.txt", "from b", null, name -> null);
+ assertEquals(List.of("from b"), second.read("*", ViewShape.TEXTS, name -> null).value());
+ }
+ try (StorageSession shared = type.open(connection("flow-data", false, false), StorageScope.shared())) {
+ assertTrue(((List>) shared.list("run-*/brief.txt", name -> null).value()).containsAll(List.of("run-a/brief.txt", "run-b/brief.txt")));
+ }
+ }
+
+ @Test
+ void aReadThatMatchesNothingIsEmptyRatherThanAFailure() {
+ try (StorageSession session = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-empty"))) {
+ assertEquals(0, session.read("notes/*.json", ViewShape.JSON, name -> null).count());
+ assertEquals(0, session.read("first-round.txt", ViewShape.FILES, name -> null).count());
+ }
+ }
+
+ @Test
+ void aKeyFromAValueCannotClimbOutOfItsScope() {
+ try (StorageSession session = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-safe"))) {
+ StorageException refusal = assertThrows(StorageException.class, () -> session.write("${{name}}", "x", null,
+ Map.of("name", "../run-a/brief.txt")::get));
+ assertEquals(StorageException.INVALID_LOCATION, refusal.getErrorCode());
+ }
+ }
+
+ @Test
+ void deleteRemovesOneObjectAndSaysWhetherItWasThere() {
+ try (StorageSession session = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-delete"))) {
+ session.write("gone.txt", "bye", null, name -> null);
+ assertEquals(1, session.delete("gone.txt", name -> null).affected());
+ assertEquals(0, session.delete("gone.txt", name -> null).affected());
+ }
+ }
+
+ @Test
+ void aReadOnlyStorageRefusesWritesAndAMissingBucketIsSaidAtInit() {
+ try (StorageSession readOnly = type.open(connection("flow-data", true, false), StorageScope.shared())) {
+ assertEquals(StorageException.OPERATION_NOT_ALLOWED,
+ assertThrows(StorageException.class, () -> readOnly.write("x.txt", "x", null, name -> null)).getErrorCode());
+ }
+ try (StorageSession missing = type.open(connection("not-created", false, false), StorageScope.shared())) {
+ StorageException refusal = assertThrows(StorageException.class, () -> missing.initialize(new StorageInit(true, null)));
+ assertEquals(StorageException.INIT_FAILED, refusal.getErrorCode());
+ assertTrue(refusal.getMessage().contains("may not be prepared by flows"), refusal.getMessage());
+ }
+ }
+
+ @Test
+ void wrongCredentialsAreSaidAsARefusal() {
+ StorageConnection wrong = new StorageConnection("test store", S3StorageType.NAME, MAPPER.valueToTree(Map.of(
+ "endpoint", MINIO.getS3URL(), "bucket", "flow-data", "accessKey", "nobody", "secretKey", "wrong-secret")),
+ false, false, true);
+ try (StorageSession session = type.open(wrong, StorageScope.shared())) {
+ StorageException refusal = assertThrows(StorageException.class, () -> session.read("*", ViewShape.KEYS, name -> null));
+ assertEquals(StorageException.OPERATION_NOT_ALLOWED, refusal.getErrorCode());
+ }
+ }
+
+ @Test
+ void editorChecksSayWhatIsWrong() {
+ assertTrue(type.validateTemplate("reports/*.json").stream().anyMatch(p -> p.contains("cannot contain '*'")));
+ assertTrue(type.validateTemplate("${{value}}.txt").stream().anyMatch(p -> p.contains("cannot be part of its key")));
+ assertTrue(type.validateTemplate("reports/${{value.id}}.json").isEmpty());
+ assertTrue(type.validateQuery("../x", ViewShape.TEXTS).stream().anyMatch(p -> p.contains("'..'")));
+ assertTrue(type.validateDelete("old/*").stream().anyMatch(p -> p.contains("names one object")));
+ assertTrue(type.validateInit(new StorageInit(false, "CREATE TABLE t (id int)")).getFirst().contains("no init script"));
+ }
+
+ private static StorageConnection connection(String bucket, boolean readOnly, boolean allowInit) {
+ return new StorageConnection("test store", S3StorageType.NAME, MAPPER.valueToTree(Map.of(
+ "endpoint", MINIO.getS3URL(), "bucket", bucket, "rootPrefix", "flows",
+ "accessKey", MINIO.getUserName(), "secretKey", MINIO.getPassword())),
+ readOnly, allowInit, true);
+ }
+}