Add pluggable storage types: S3-compatible and PostgreSQL

The groundwork for a Storage node: a StorageType SPI collected like the
LLM providers, with two types - an S3-compatible object store (MinIO,
AWS, Ceph) on the MinIO client, and PostgreSQL spoken to in SQL.

Every type answers the same four things - a template to write to, a
query to read with, a script to prepare, the shapes a read comes back
in - and checks them itself, so the node stays generic.

- S3: keys and globs (pippo*, reports/**.json) under an instance root
  and an optional per-execution prefix that a key cannot climb out of.
- PostgreSQL: placeholders are bound, never spliced in, and refused
  inside literals or comments; reads run READ ONLY; one statement only,
  checked here because pgjdbc splits and runs 'INSERT ...; DROP ...';
  the JDBC URL and driver properties are built from checked fields.
- storages.json holds the operator's instances, checked at startup;
  /storage/types and /storage/instances expose them without settings.
- A person's own connection goes through a storage endpoint guard
  (app.storage.endpoint.*), like an LLM endpoint.

Tested against real MinIO and PostgreSQL with Testcontainers.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Lucio Lelii 2026-09-28 11:35:47 +02:00
parent 7b2e811ab4
commit c472c80674
36 changed files with 3013 additions and 0 deletions

16
pom.xml
View File

@ -114,6 +114,12 @@
<version>0.13.0</version>
<scope>runtime</scope>
</dependency>
<!-- S3-compatible object storage (MinIO, AWS S3, Ceph) for the Storage node -->
<dependency>
<groupId>io.minio</groupId>
<artifactId>minio</artifactId>
<version>8.5.17</version>
</dependency>
<!-- PostgreSQL database -->
<dependency>
<groupId>org.postgresql</groupId>
@ -177,6 +183,16 @@
<artifactId>testcontainers-ollama</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>testcontainers-minio</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>testcontainers-postgresql</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-webmvc-test</artifactId>

View File

@ -0,0 +1,61 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.configurations.retrievers;
import java.util.List;
import java.util.Map;
import org.springframework.http.HttpStatus;
import org.springframework.stereotype.Component;
import org.springframework.web.server.ResponseStatusException;
import it.cnr.isti.workflow.manager.storage.StorageInstancesProvider;
import it.cnr.isti.workflow.manager.storage.StorageType;
import it.cnr.isti.workflow.manager.storage.StorageTypes;
/**
* The storage node's dropdowns: the types, the catalog instances of a type, and the shapes a read
* of that type can come back in. Instances are listed by id alone - their settings, credentials
* included, never leave the service.
*/
@Component
public class StorageFieldRetriever implements DynamicFieldRetriever {
public static final String CATEGORY = "Storage";
private final StorageTypes types;
private final StorageInstancesProvider instances;
public StorageFieldRetriever(StorageTypes types, StorageInstancesProvider instances) {
this.types = types;
this.instances = instances;
}
@Override
public String getCategory() {
return CATEGORY;
}
@Override
public List<String> retrieve(String parameter, Map<String, String> params) {
return switch (parameter) {
case "types" -> types.all().stream().map(StorageType::getName).toList();
case "instances" -> {
String type = params == null ? null : params.get("storageType");
yield (type == null || type.isBlank() ? instances.all() : instances.ofType(type)).stream()
.map(StorageInstancesProvider.StorageInstance::id)
.sorted()
.toList();
}
case "shapes" -> {
String type = params == null ? null : params.get("storageType");
yield types.find(type).map(found -> found.viewShapes().stream().map(Enum::name).toList())
.orElse(List.of());
}
default -> throw new ResponseStatusException(HttpStatus.NOT_FOUND,
"Unknown Storage retriever parameter: " + parameter);
};
}
}

View File

@ -0,0 +1,46 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.controllers;
import java.util.List;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.security.SecurityRequirement;
import it.cnr.isti.workflow.manager.storage.StorageInstanceView;
import it.cnr.isti.workflow.manager.storage.StorageInstancesProvider;
import it.cnr.isti.workflow.manager.storage.StorageTypeMetadata;
import it.cnr.isti.workflow.manager.storage.StorageTypes;
@RestController
@RequestMapping("/storage")
@SecurityRequirement(name = "bearerAuth")
public class StorageController {
private final StorageTypes types;
private final StorageInstancesProvider instances;
public StorageController(StorageTypes types, StorageInstancesProvider instances) {
this.types = types;
this.instances = instances;
}
@GetMapping("/types")
@Operation(summary = "List storage types",
description = "The kinds of storage a flow can use, what a read of each returns, and the fields of a personal connection to one.")
public List<StorageTypeMetadata> types() {
return types.all().stream().map(StorageTypeMetadata::of).toList();
}
@GetMapping("/instances")
@Operation(summary = "List catalog storages",
description = "The storages this service offers to every flow. Connection settings are never returned.")
public List<StorageInstanceView> instances() {
return instances.all().stream().map(StorageInstanceView::of).toList();
}
}

View File

@ -15,4 +15,9 @@ public class NodeExecutionException extends RuntimeException {
super(message);
this.errorCode = errorCode;
}
public NodeExecutionException(String errorCode, String message, Throwable cause) {
super(message, cause);
this.errorCode = errorCode;
}
}

View File

@ -74,6 +74,19 @@ public class OutboundEndpointGuard {
if (!StringUtils.hasText(host)) {
throw new IllegalArgumentException("Endpoint URL must include a host: " + url);
}
validateHost(host);
}
/**
* The address rules alone, for a connection named by a host rather than a URL - a database's.
*
* @throws IllegalArgumentException if the host is blank, or - unless allow-listed or the check
* is disabled - resolves to a disallowed address
*/
public void validateHost(String host) {
if (!StringUtils.hasText(host)) {
throw new IllegalArgumentException("Endpoint host cannot be empty");
}
if (isAllowedHost(host) || allowPrivateNetwork) {
return;
}

View File

@ -0,0 +1,13 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
/**
* One field of a connection to a storage of some type, as a person fills it in when they save their
* own connection in the vault. The whole connection is encrypted together; {@code secret} only says
* which fields the form masks.
*/
public record CredentialField(String key, String label, String description, boolean secret, boolean required) {
}

View File

@ -0,0 +1,46 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
import tools.jackson.databind.JsonNode;
/**
* Everything needed to reach one storage: its type, the settings that type reads (an endpoint and
* a bucket, a host and a database) including the secret ones, and what the operator allows on it.
*
* @param name how errors and events name it - the catalog instance, or "your PostgreSQL connection"
* @param operatorChosen true for a catalog instance: its address was chosen by whoever runs this
* service, and is not checked against the private-network guard that a person's own
* connection is
*/
public record StorageConnection(String name, String type, JsonNode settings, boolean readOnly, boolean allowInit,
boolean operatorChosen) {
/** Only for the error text of a missing setting: settings themselves are never shown. */
public String requireText(String key) {
JsonNode value = settings == null ? null : settings.get(key);
if (value == null || value.isNull() || value.asString().isBlank()) {
throw new StorageException(StorageException.CONNECTION_INVALID,
name + " has no " + key + " set");
}
return value.asString().trim();
}
public String optionalText(String key) {
JsonNode value = settings == null ? null : settings.get(key);
if (value == null || value.isNull() || value.asString().isBlank()) {
return null;
}
return value.asString().trim();
}
public boolean flag(String key, boolean fallback) {
JsonNode value = settings == null ? null : settings.get(key);
if (value == null || value.isNull()) {
return fallback;
}
return value.isBoolean() ? value.asBoolean() : Boolean.parseBoolean(value.asString());
}
}

View File

@ -0,0 +1,62 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
import 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.
*
* <p>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));
}
}
}

View File

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

View File

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

View File

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

View File

@ -0,0 +1,146 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
import 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}.
*
* <p>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.
*
* <p>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<StorageInstance> 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<StorageInstance> all() {
return instances;
}
public List<StorageInstance> ofType(String type) {
return instances.stream().filter(instance -> instance.type().equalsIgnoreCase(type)).toList();
}
public Optional<StorageInstance> find(String id) {
return instances.stream().filter(instance -> instance.id().equals(id)).findFirst();
}
private static List<StorageInstance> 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<StorageInstance> loaded = new ArrayList<>();
Set<String> ids = new HashSet<>();
for (StorageInstance instance : catalog.storages() == null ? List.<StorageInstance>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<StorageInstance> 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);
}
}
}

View File

@ -0,0 +1,13 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
/** What a step can ask of a storage, one operation at a time. */
public enum StorageOperation {
READ,
WRITE,
LIST,
DELETE
}

View File

@ -0,0 +1,157 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
import 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.
*
* <p>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.<field>} 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<String> names(String template) {
Set<String> 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<String> parameters(String template) {
Set<String> 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<String> invalidNames(String template) {
Set<String> 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<String, Object> 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<String, Object> 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<String, Object> withValue(Object value, Function<String, Object> 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;
}
}

View File

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

View File

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

View File

@ -0,0 +1,32 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
import java.util.function.Function;
/**
* One conversation with a storage, opened for an operation and closed after it.
*
* <p>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<String, Object> values);
StorageReadResult read(String query, ViewShape shape, Function<String, Object> values);
/** What there is: the keys under a prefix or glob, or the tables of a database. */
StorageReadResult list(String query, Function<String, Object> values);
StorageWriteResult delete(String template, Function<String, Object> values);
@Override
void close();
}

View File

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

View File

@ -0,0 +1,57 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
import 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.
*
* <p>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<CredentialField> credentialFields();
/** The shapes a read can come back in. The first is the default. */
List<ViewShape> viewShapes();
/** What a write accepts, for the kinds of value the target's port takes. */
List<IOCapability> 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<String> validateTemplate(String template);
/** Problems with a read query for the shape asked. */
List<String> validateQuery(String query, ViewShape shape);
/** Problems with a delete template. */
List<String> validateDelete(String template);
List<String> validateInit(StorageInit init);
/** The placeholders of a template or query that bind as values rather than render as text. */
default Set<String> parameters(String templateOrQuery) {
return StoragePlaceholders.parameters(templateOrQuery);
}
StorageSession open(StorageConnection connection, StorageScope scope);
}

View File

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

View File

@ -0,0 +1,40 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
import 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<StorageType> types;
public StorageTypes(List<StorageType> types) {
this.types = types.stream()
.sorted(Comparator.comparing(StorageType::getName, String.CASE_INSENSITIVE_ORDER))
.toList();
}
public List<StorageType> all() {
return types;
}
public Optional<StorageType> 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));
}
}

View File

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

View File

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

View File

@ -0,0 +1,320 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage.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<String, Object> values) {
requireWritable("write to");
return change(template, PostgresStorageType.WRITE_VERBS, "write", StoragePlaceholders.withValue(value, values));
}
@Override
public StorageReadResult read(String query, ViewShape shape, Function<String, Object> 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<String, Object> 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<String> tables = new ArrayList<>();
try (ResultSet rows = statement.executeQuery()) {
while (rows.next()) {
tables.add(rows.getString(1));
}
}
boolean truncated = tables.size() > settings.maxKeys();
List<String> kept = truncated ? tables.subList(0, settings.maxKeys()) : tables;
return new StorageReadResult(List.copyOf(kept), kept.size(), truncated);
}
});
}
@Override
public StorageWriteResult delete(String template, Function<String, Object> 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<String> verbs, String what, Function<String, Object> 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<String> 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<String, Object> values) throws SQLException {
PreparedStatement statement = jdbc.prepareStatement(template.sql());
List<String> 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> {
T run() throws SQLException;
}
/** One statement's transaction: read-only when asked, with the statement timeout set inside it. */
private <T> T inTransaction(boolean readOnly, String what, Work<T> 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);
}
}

View File

@ -0,0 +1,274 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage.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.
*
* <ul>
* <li>A query is a {@code SELECT} (or {@code WITH}), run in a read-only transaction, returning its
* rows as JSON.
* <li>A template is an {@code INSERT}, {@code UPDATE} or {@code MERGE}; a delete is a
* {@code DELETE}. Both may end in {@code RETURNING}.
* <li>Init is a script - {@code CREATE TABLE IF NOT EXISTS ...} - run in one transaction.
* </ul>
*
* <p>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<String> SSL_MODES = Set.of("disable", "allow", "prefer", "require", "verify-ca", "verify-full");
static final Set<String> READ_VERBS = Set.of("SELECT", "WITH", "VALUES", "TABLE");
static final Set<String> WRITE_VERBS = Set.of("INSERT", "UPDATE", "MERGE");
static final Set<String> DELETE_VERBS = Set.of("DELETE");
private final StorageSettings settings;
private final StorageEndpointGuard guard;
private final ObjectMapper objectMapper;
private final Map<String, HikariDataSource> 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<CredentialField> 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<ViewShape> viewShapes() {
return List.of(ViewShape.JSON);
}
@Override
public List<IOCapability> 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<String> validateTemplate(String template) {
return validateStatement(template, WRITE_VERBS, "A write", true);
}
@Override
public List<String> validateQuery(String query, ViewShape shape) {
List<String> 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<String> validateDelete(String template) {
return validateStatement(template, DELETE_VERBS, "A delete", false);
}
@Override
public List<String> validateInit(StorageInit init) {
List<String> 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<String> validateStatement(String sql, Set<String> verbs, String what, boolean valueAllowed) {
List<String> 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);
});
}
}

View File

@ -0,0 +1,224 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage.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.
*
* <p>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.
*
* <p>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<String> parameters;
private final String verb;
private final List<String> problems;
private SqlTemplate(String sql, List<String> parameters, String verb, List<String> 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<String> parameters() {
return parameters;
}
/** The first keyword, upper case: SELECT, INSERT, WITH... Empty for a blank statement. */
String verb() {
return verb;
}
List<String> problems() {
return problems;
}
static SqlTemplate parse(String template) {
List<String> problems = new ArrayList<>();
List<String> 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<String> 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<String> 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);
}
}

View File

@ -0,0 +1,136 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage.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.
*
* <p>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<String> problemsAsWritten(String text, String what) {
List<String> 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;
}
}

View File

@ -0,0 +1,315 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage.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<String, Object> 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<String, Object> 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<String, Object> 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<String, Object> 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<String> 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<String> keys = new ArrayList<>();
boolean truncated = call("list " + keyOrGlob, () -> {
for (Result<Item> 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<File> downloadAll(List<String> keys) {
List<File> 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<String> textsOf(List<String> keys) {
return keys.stream().map(key -> new String(bytesOf(key), StandardCharsets.UTF_8)).toList();
}
private ArrayNode jsonOf(List<String> 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> {
T run() throws Exception;
}
private <T> T call(String what, S3Call<T> 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);
}
}

View File

@ -0,0 +1,185 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage.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.
*
* <ul>
* <li>A template is an object key, placeholders rendered in: {@code reports/${{context.iteration}}.json}.
* <li>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.
* <li>Init can create the bucket when it is missing; there is no script.
* </ul>
*/
@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<CredentialField> 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<ViewShape> viewShapes() {
return List.of(ViewShape.FILES, ViewShape.TEXTS, ViewShape.JSON, ViewShape.KEYS);
}
@Override
public List<IOCapability> 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<String> validateTemplate(String template) {
List<String> 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.<field>}} for a field of a JSON value");
}
return problems;
}
@Override
public List<String> validateQuery(String query, ViewShape shape) {
List<String> 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<String> validateDelete(String template) {
List<String> 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<String> validateInit(StorageInit init) {
List<String> 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);
}
}

View File

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

View File

@ -0,0 +1 @@
{ "storages": [] }

View File

@ -0,0 +1,82 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.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());
}
}

View File

@ -0,0 +1,198 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.storage.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<String, Object> values(Map<String, Object> 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);
}
}

View File

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

View File

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

View File

@ -0,0 +1,162 @@
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
// SPDX-License-Identifier: AGPL-3.0-or-later
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
package it.cnr.isti.workflow.manager.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.<String, Object>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<File> files = (List<File>) 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.<String, Object>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);
}
}