From c472c806745288d3b9d699ea2df687302fc36239 Mon Sep 17 00:00:00 2001 From: Lucio Lelii Date: Mon, 28 Sep 2026 11:35:47 +0200 Subject: [PATCH] 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) --- pom.xml | 16 + .../retrievers/StorageFieldRetriever.java | 61 ++++ .../controllers/StorageController.java | 46 +++ .../executions/NodeExecutionException.java | 5 + .../llms/providers/OutboundEndpointGuard.java | 13 + .../manager/storage/CredentialField.java | 13 + .../manager/storage/StorageConnection.java | 46 +++ .../manager/storage/StorageEndpointGuard.java | 62 ++++ .../manager/storage/StorageException.java | 30 ++ .../workflow/manager/storage/StorageInit.java | 25 ++ .../manager/storage/StorageInstanceView.java | 15 + .../storage/StorageInstancesProvider.java | 146 ++++++++ .../manager/storage/StorageOperation.java | 13 + .../manager/storage/StoragePlaceholders.java | 157 +++++++++ .../manager/storage/StorageReadResult.java | 15 + .../manager/storage/StorageScope.java | 25 ++ .../manager/storage/StorageSession.java | 32 ++ .../manager/storage/StorageSettings.java | 56 +++ .../workflow/manager/storage/StorageType.java | 57 ++++ .../manager/storage/StorageTypeMetadata.java | 26 ++ .../manager/storage/StorageTypes.java | 40 +++ .../manager/storage/StorageWriteResult.java | 17 + .../workflow/manager/storage/ViewShape.java | 39 +++ .../storage/postgres/PostgresSession.java | 320 ++++++++++++++++++ .../storage/postgres/PostgresStorageType.java | 274 +++++++++++++++ .../manager/storage/postgres/SqlTemplate.java | 224 ++++++++++++ .../workflow/manager/storage/s3/S3Keys.java | 136 ++++++++ .../manager/storage/s3/S3Session.java | 315 +++++++++++++++++ .../manager/storage/s3/S3StorageType.java | 185 ++++++++++ src/main/resources/application.properties | 13 + src/main/resources/storages.json | 1 + .../storage/StorageInstancesProviderTest.java | 82 +++++ .../postgres/PostgresStorageTypeTest.java | 198 +++++++++++ .../storage/postgres/SqlTemplateTest.java | 81 +++++ .../manager/storage/s3/S3KeysTest.java | 67 ++++ .../manager/storage/s3/S3StorageTypeTest.java | 162 +++++++++ 36 files changed, 3013 insertions(+) create mode 100644 src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/StorageFieldRetriever.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/controllers/StorageController.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/CredentialField.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageConnection.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageEndpointGuard.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageException.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageInit.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstanceView.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProvider.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageOperation.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StoragePlaceholders.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageReadResult.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageScope.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageSession.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageSettings.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageType.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypeMetadata.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypes.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/StorageWriteResult.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/ViewShape.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresSession.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageType.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplate.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Keys.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Session.java create mode 100644 src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageType.java create mode 100644 src/main/resources/storages.json create mode 100644 src/test/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProviderTest.java create mode 100644 src/test/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageTypeTest.java create mode 100644 src/test/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplateTest.java create mode 100644 src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3KeysTest.java create mode 100644 src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageTypeTest.java diff --git a/pom.xml b/pom.xml index 4591c94..5394e95 100644 --- a/pom.xml +++ b/pom.xml @@ -114,6 +114,12 @@ 0.13.0 runtime + + + io.minio + minio + 8.5.17 + org.postgresql @@ -177,6 +183,16 @@ testcontainers-ollama test + + org.testcontainers + testcontainers-minio + test + + + org.testcontainers + testcontainers-postgresql + test + org.springframework.boot spring-boot-webmvc-test diff --git a/src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/StorageFieldRetriever.java b/src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/StorageFieldRetriever.java new file mode 100644 index 0000000..175f7db --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/configurations/retrievers/StorageFieldRetriever.java @@ -0,0 +1,61 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.configurations.retrievers; + +import java.util.List; +import java.util.Map; + +import org.springframework.http.HttpStatus; +import org.springframework.stereotype.Component; +import org.springframework.web.server.ResponseStatusException; + +import it.cnr.isti.workflow.manager.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 retrieve(String parameter, Map 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); + }; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/controllers/StorageController.java b/src/main/java/it/cnr/isti/workflow/manager/controllers/StorageController.java new file mode 100644 index 0000000..7c80ffd --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/controllers/StorageController.java @@ -0,0 +1,46 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.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 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 instances() { + return instances.all().stream().map(StorageInstanceView::of).toList(); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/executions/NodeExecutionException.java b/src/main/java/it/cnr/isti/workflow/manager/executions/NodeExecutionException.java index f1abdcb..2b19337 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/executions/NodeExecutionException.java +++ b/src/main/java/it/cnr/isti/workflow/manager/executions/NodeExecutionException.java @@ -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; + } } diff --git a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/OutboundEndpointGuard.java b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/OutboundEndpointGuard.java index 7a77e18..4987ade 100644 --- a/src/main/java/it/cnr/isti/workflow/manager/llms/providers/OutboundEndpointGuard.java +++ b/src/main/java/it/cnr/isti/workflow/manager/llms/providers/OutboundEndpointGuard.java @@ -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; } diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/CredentialField.java b/src/main/java/it/cnr/isti/workflow/manager/storage/CredentialField.java new file mode 100644 index 0000000..a8edd26 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/CredentialField.java @@ -0,0 +1,13 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +/** + * 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) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageConnection.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageConnection.java new file mode 100644 index 0000000..ac05a7d --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageConnection.java @@ -0,0 +1,46 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import tools.jackson.databind.JsonNode; + +/** + * 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()); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageEndpointGuard.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageEndpointGuard.java new file mode 100644 index 0000000..d803671 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageEndpointGuard.java @@ -0,0 +1,62 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; + +import it.cnr.isti.workflow.manager.llms.providers.OutboundEndpointGuard; + +/** + * Where a person's own storage connection may point. + * + *

The same rules as an LLM endpoint - see {@link OutboundEndpointGuard} - behind a separate + * setting, {@code app.storage.endpoint.*}: an operator may well open one to a private network and + * not the other. Applied when the connection is saved and again right before every connection, + * which is the check that protects; a catalog instance was chosen by the operator and does not + * pass through here. + */ +@Component +public class StorageEndpointGuard { + + private final OutboundEndpointGuard guard; + + public StorageEndpointGuard( + @Value("${app.storage.endpoint.allow-private-network:false}") boolean allowPrivateNetwork, + @Value("${app.storage.endpoint.allowed-host-patterns:}") String allowedHostPatterns) { + this.guard = new OutboundEndpointGuard(allowPrivateNetwork, allowedHostPatterns); + } + + /** An http(s) endpoint, such as an object store's. */ + public void validateUrl(String url) { + try { + guard.validate(url); + } catch (IllegalArgumentException ex) { + throw new StorageException(StorageException.CONNECTION_INVALID, ex.getMessage()); + } + } + + /** A bare host, such as a database's. */ + public void validateHost(String host) { + try { + guard.validateHost(host); + } catch (IllegalArgumentException ex) { + throw new StorageException(StorageException.CONNECTION_INVALID, ex.getMessage()); + } + } + + /** Checks a person's own connection; lets an operator's through. */ + public void check(StorageConnection connection, String urlKey, String hostKey) { + if (connection.operatorChosen()) { + return; + } + if (urlKey != null) { + validateUrl(connection.requireText(urlKey)); + } + if (hostKey != null) { + validateHost(connection.requireText(hostKey)); + } + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageException.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageException.java new file mode 100644 index 0000000..e4cb0b6 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageException.java @@ -0,0 +1,30 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import it.cnr.isti.workflow.manager.executions.NodeExecutionException; + +/** A storage refusing, or failing, what a step asked of it - always said in the step's terms. */ +public class StorageException extends NodeExecutionException { + + public static final String CONNECTION_INVALID = "STORAGE_CONNECTION_INVALID"; + public static final String UNREACHABLE = "STORAGE_UNREACHABLE"; + public static final String NOT_FOUND = "STORAGE_OBJECT_NOT_FOUND"; + public static final String TOO_LARGE = "STORAGE_OBJECT_TOO_LARGE"; + public static final String OPERATION_NOT_ALLOWED = "STORAGE_OPERATION_NOT_ALLOWED"; + public static final String PLACEHOLDER_MISSING = "STORAGE_PLACEHOLDER_MISSING"; + public static final String INVALID_LOCATION = "STORAGE_INVALID_LOCATION"; + public static final String SQL_REJECTED = "STORAGE_SQL_REJECTED"; + public static final String INIT_FAILED = "STORAGE_INIT_FAILED"; + public static final String CONTENT_INVALID = "STORAGE_CONTENT_INVALID"; + + public StorageException(String errorCode, String message) { + super(errorCode, message); + } + + public StorageException(String errorCode, String message, Throwable cause) { + super(errorCode, message, cause); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInit.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInit.java new file mode 100644 index 0000000..3036999 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInit.java @@ -0,0 +1,25 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +/** + * What a storage node prepares before the first step of an execution runs. Both are idempotent by + * contract, because every execution runs them again. + * + * @param createIfMissing create the container the connection names (an object store's bucket) when + * it does not exist yet + * @param script statements that prepare a database - {@code CREATE TABLE IF NOT EXISTS ...} - run + * in one transaction; blank for none + */ +public record StorageInit(boolean createIfMissing, String script) { + + public static StorageInit none() { + return new StorageInit(false, null); + } + + public boolean hasScript() { + return script != null && !script.isBlank(); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstanceView.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstanceView.java new file mode 100644 index 0000000..59e979a --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstanceView.java @@ -0,0 +1,15 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +/** A catalog storage as a flow author sees it: what it is and what it allows, not how to reach it. */ +public record StorageInstanceView(String id, String name, String description, String type, boolean readOnly, + boolean allowInit) { + + public static StorageInstanceView of(StorageInstancesProvider.StorageInstance instance) { + return new StorageInstanceView(instance.id(), instance.displayName(), instance.description(), instance.type(), + instance.readOnly(), instance.allowInit()); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProvider.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProvider.java new file mode 100644 index 0000000..9c5b8ce --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProvider.java @@ -0,0 +1,146 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import java.io.InputStream; +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Objects; +import java.util.Optional; +import java.util.Set; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.core.io.ClassPathResource; +import org.springframework.core.io.Resource; +import org.springframework.core.io.ResourceLoader; +import org.springframework.stereotype.Component; +import org.springframework.util.StringUtils; + +import com.fasterxml.jackson.annotation.JsonCreator; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import com.fasterxml.jackson.annotation.JsonProperty; +import tools.jackson.databind.JsonNode; +import tools.jackson.databind.ObjectMapper; + +/** + * The storages whoever runs this service offers to every flow, from {@code storages.json}. + * + *

Loaded the way the MCP server catalog is: {@code app.storage.instances.file} names the file, + * and without it the (empty) one on the classpath is used. Every entry is checked by its own type + * as it loads, and a file that does not hold together stops the service from starting: an instance + * that looked fine in the editor and failed at every run would be found much later and by someone + * else. + * + *

The file holds the credentials of the instances it lists, like {@code mcp-servers.json} does, + * so it belongs outside the image, readable by the service alone. + */ +@Component +public class StorageInstancesProvider { + + private final List instances; + + public StorageInstancesProvider(ObjectMapper objectMapper, ResourceLoader resourceLoader, StorageTypes types, + @Value("${app.storage.instances.file:}") String location) { + this.instances = load(Objects.requireNonNull(objectMapper), Objects.requireNonNull(resourceLoader), + Objects.requireNonNull(types), location); + } + + public List all() { + return instances; + } + + public List ofType(String type) { + return instances.stream().filter(instance -> instance.type().equalsIgnoreCase(type)).toList(); + } + + public Optional find(String id) { + return instances.stream().filter(instance -> instance.id().equals(id)).findFirst(); + } + + private static List load(ObjectMapper objectMapper, ResourceLoader resourceLoader, + StorageTypes types, String location) { + Catalog catalog; + try (InputStream input = resolve(resourceLoader, location).getInputStream()) { + catalog = objectMapper.readValue(input, Catalog.class); + } catch (Exception ex) { + throw new IllegalStateException("Unable to load the storage catalog from " + describe(location), ex); + } + List loaded = new ArrayList<>(); + Set ids = new HashSet<>(); + for (StorageInstance instance : catalog.storages() == null ? List.of() : catalog.storages()) { + String where = "Storage catalog " + describe(location) + ", entry " + (loaded.size() + 1); + if (instance == null || !StringUtils.hasText(instance.id())) { + throw new IllegalStateException(where + " has no id"); + } + if (!ids.add(instance.id())) { + throw new IllegalStateException(where + ": id " + instance.id() + " is used twice"); + } + StorageType type = types.find(instance.type()).orElseThrow(() -> new IllegalStateException( + where + " (" + instance.id() + ") has unknown type " + instance.type() + + "; known types: " + types.all().stream().map(StorageType::getName).toList())); + try { + type.validateConnection(instance.connection(type)); + } catch (StorageException ex) { + throw new IllegalStateException(where + " (" + instance.id() + "): " + ex.getMessage(), ex); + } + loaded.add(instance); + } + return List.copyOf(loaded); + } + + private static Resource resolve(ResourceLoader resourceLoader, String location) { + if (!StringUtils.hasText(location)) { + return new ClassPathResource("storages.json"); + } + String resolved = location.startsWith("classpath:") || location.startsWith("file:") ? location : "file:" + location; + Resource resource = resourceLoader.getResource(resolved); + if (!resource.exists()) { + throw new IllegalStateException("Storage catalog not found at " + resolved); + } + return resource; + } + + private static String describe(String location) { + return StringUtils.hasText(location) ? location : "classpath:storages.json"; + } + + @JsonIgnoreProperties(ignoreUnknown = true) + private record Catalog(List storages) { + } + + /** + * One storage of the catalog. + * + * @param readOnly only reads: a node writing to it, or a block deleting from it, is refused + * while the flow is edited + * @param allowInit whether a flow may prepare it - create its bucket, run its DDL script + */ + @JsonIgnoreProperties(ignoreUnknown = true) + public record StorageInstance(String id, String name, String description, String type, boolean readOnly, + boolean allowInit, JsonNode connection) { + + @JsonCreator + static StorageInstance of(@JsonProperty("id") String id, + @JsonProperty("name") String name, + @JsonProperty("description") String description, + @JsonProperty("type") String type, + @JsonProperty("readOnly") Boolean readOnly, + @JsonProperty("allowInit") Boolean allowInit, + @JsonProperty("connection") JsonNode connection) { + // Left out means no: an instance is writable-but-not-preparable only when the file says so. + return new StorageInstance(id, name, description, type, Boolean.TRUE.equals(readOnly), + Boolean.TRUE.equals(allowInit), connection); + } + + public String displayName() { + return StringUtils.hasText(name) ? name : id; + } + + public StorageConnection connection(StorageType type) { + return new StorageConnection(displayName(), type.getName(), connection, readOnly, allowInit, true); + } + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageOperation.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageOperation.java new file mode 100644 index 0000000..c3973a1 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageOperation.java @@ -0,0 +1,13 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +/** What a step can ask of a storage, one operation at a time. */ +public enum StorageOperation { + READ, + WRITE, + LIST, + DELETE +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StoragePlaceholders.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StoragePlaceholders.java new file mode 100644 index 0000000..9feed59 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StoragePlaceholders.java @@ -0,0 +1,157 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import java.util.LinkedHashSet; +import java.util.Set; +import java.util.function.Function; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +import it.cnr.isti.workflow.manager.executions.ExecutionTemplateResolver; + +/** + * The ${{name}} placeholders of a storage template - an object key, a glob, a statement. + * + *

Names come in three kinds, and the storage node treats them differently: {@code context.*} + * and {@code global.*} are the execution's own and always known; {@code value} and + * {@code value.} are the value being written; anything else is a parameter the node takes + * from another step through a port. + */ +public final class StoragePlaceholders { + + private static final Pattern PLACEHOLDER = Pattern.compile("\\$\\{\\{\\s*([^}]*?)\\s*}}"); + private static final Pattern VALID_NAME = Pattern.compile("^[A-Za-z][A-Za-z0-9_.-]*$"); + + public static final String VALUE = "value"; + + private StoragePlaceholders() { + } + + /** Every name, in the order it first appears. */ + public static Set names(String template) { + Set names = new LinkedHashSet<>(); + if (template == null) { + return names; + } + Matcher matcher = PLACEHOLDER.matcher(template); + while (matcher.find()) { + names.add(matcher.group(1)); + } + return names; + } + + /** The names another step has to supply: neither the execution's own nor the written value. */ + public static Set parameters(String template) { + Set parameters = new LinkedHashSet<>(); + for (String name : names(template)) { + if (!isExecutionName(name) && !isValueName(name)) { + parameters.add(name); + } + } + return parameters; + } + + public static boolean isExecutionName(String name) { + return name.startsWith("context.") || name.startsWith("global.") || name.startsWith("project."); + } + + public static boolean isValueName(String name) { + return name.equals(VALUE) || name.startsWith(VALUE + "."); + } + + /** Names that would not survive as a port or a lookup, for the node's validation. */ + public static Set invalidNames(String template) { + Set invalid = new LinkedHashSet<>(); + for (String name : names(template)) { + if (!VALID_NAME.matcher(name).matches()) { + invalid.add(name); + } + } + return invalid; + } + + /** + * The template with every placeholder replaced by its value as text - for an object key or a + * glob, never for a statement, whose values are bound instead. + */ + public static String render(String template, Function values) { + if (template == null) { + return null; + } + Matcher matcher = PLACEHOLDER.matcher(template); + StringBuilder rendered = new StringBuilder(); + while (matcher.find()) { + Object value = require(values, matcher.group(1)); + matcher.appendReplacement(rendered, Matcher.quoteReplacement(ExecutionTemplateResolver.formatValue(value))); + } + matcher.appendTail(rendered); + return rendered.toString(); + } + + public static Object require(Function values, String name) { + Object value = values == null ? null : values.apply(name); + if (value == null) { + throw new StorageException(StorageException.PLACEHOLDER_MISSING, + "No value for ${{" + name + "}}"); + } + return value; + } + + /** + * Lookups for a write: {@code value} is the value itself, {@code value.a.b} walks into a JSON + * value (or a map) field by field, and every other name goes to {@code values}. + */ + public static Function withValue(Object value, Function values) { + return name -> { + if (name.equals(VALUE)) { + return value; + } + if (name.startsWith(VALUE + ".")) { + return field(value, name.substring(VALUE.length() + 1)); + } + return values == null ? null : values.apply(name); + }; + } + + private static Object field(Object value, String path) { + Object current = value; + for (String part : path.split("\\.")) { + if (current instanceof tools.jackson.databind.JsonNode node) { + tools.jackson.databind.JsonNode next = node.get(part); + if (next == null || next.isNull()) { + return null; + } + current = next.isValueNode() ? scalar(next) : next; + } else if (current instanceof java.util.Map map) { + current = map.get(part); + } else { + return null; + } + if (current == null) { + return null; + } + } + return current; + } + + private static Object scalar(tools.jackson.databind.JsonNode node) { + if (node.isString()) { + return node.asString(); + } + if (node.isBoolean()) { + return node.asBoolean(); + } + if (node.isNumber()) { + return node.numberValue(); + } + return node.toString(); + } + + /** The pattern, for a type that has to find placeholders itself - inside a statement, say. */ + public static Pattern pattern() { + return PLACEHOLDER; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageReadResult.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageReadResult.java new file mode 100644 index 0000000..307dc8a --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageReadResult.java @@ -0,0 +1,15 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +/** + * What a read found. + * + * @param value in the shape asked for: a list of files, a list of texts or names, or one JSON value + * @param count how many objects or rows it holds + * @param truncated true when more matched than a read may return, and the rest were left out + */ +public record StorageReadResult(Object value, int count, boolean truncated) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageScope.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageScope.java new file mode 100644 index 0000000..5da6313 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageScope.java @@ -0,0 +1,25 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +/** + * Where inside the connection an execution works. Shared means where the connection already + * points; per execution adds a prefix of that execution's own, so every run starts empty and none + * sees another's data. + */ +public record StorageScope(String prefix) { + + public static StorageScope shared() { + return new StorageScope(""); + } + + public static StorageScope perExecution(String rootExecutionId) { + return new StorageScope(rootExecutionId + "/"); + } + + public StorageScope { + prefix = prefix == null ? "" : prefix; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSession.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSession.java new file mode 100644 index 0000000..335069d --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSession.java @@ -0,0 +1,32 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import java.util.function.Function; + +/** + * One conversation with a storage, opened for an operation and closed after it. + * + *

The same four things for every type, each read the way the type reads it: a {@code template} + * is an object key or a statement, a {@code query} is a glob or a {@code SELECT}. Placeholders are + * looked up through {@code values}, and a type decides whether it renders them into the text (a + * key) or binds them (a statement). + */ +public interface StorageSession extends AutoCloseable { + + void initialize(StorageInit init); + + StorageWriteResult write(String template, Object value, String contentType, Function values); + + StorageReadResult read(String query, ViewShape shape, Function values); + + /** What there is: the keys under a prefix or glob, or the tables of a database. */ + StorageReadResult list(String query, Function values); + + StorageWriteResult delete(String template, Function values); + + @Override + void close(); +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSettings.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSettings.java new file mode 100644 index 0000000..12db891 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageSettings.java @@ -0,0 +1,56 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; + +/** The limits every storage works within, whoever's connection it is. */ +@Component +public class StorageSettings { + + private final long maxObjectBytes; + private final int maxKeys; + private final int sqlStatementTimeoutSeconds; + private final int sqlMaxRows; + private final int connectTimeoutSeconds; + + public StorageSettings( + @Value("${app.storage.max-object-bytes:33554432}") long maxObjectBytes, + @Value("${app.storage.max-keys:1000}") int maxKeys, + @Value("${app.storage.sql.statement-timeout-seconds:30}") int sqlStatementTimeoutSeconds, + @Value("${app.storage.sql.max-rows:1000}") int sqlMaxRows, + @Value("${app.storage.connect-timeout-seconds:10}") int connectTimeoutSeconds) { + this.maxObjectBytes = maxObjectBytes; + this.maxKeys = maxKeys; + this.sqlStatementTimeoutSeconds = sqlStatementTimeoutSeconds; + this.sqlMaxRows = sqlMaxRows; + this.connectTimeoutSeconds = connectTimeoutSeconds; + } + + public static StorageSettings defaults() { + return new StorageSettings(32L * 1024 * 1024, 1000, 30, 1000, 10); + } + + public long maxObjectBytes() { + return maxObjectBytes; + } + + public int maxKeys() { + return maxKeys; + } + + public int sqlStatementTimeoutSeconds() { + return sqlStatementTimeoutSeconds; + } + + public int sqlMaxRows() { + return sqlMaxRows; + } + + public int connectTimeoutSeconds() { + return connectTimeoutSeconds; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageType.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageType.java new file mode 100644 index 0000000..85c0a40 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageType.java @@ -0,0 +1,57 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import java.util.List; +import java.util.Set; + +import it.cnr.isti.workflow.manager.blocks.IOCapability; + +/** + * A kind of storage a flow can keep data in - an S3-compatible object store, a PostgreSQL database. + * + *

Pluggable the way LLM providers are: each is a Spring bean, collected by {@link StorageTypes} + * and listed on {@code /storage/types}. The storage node and the storage block speak only to this + * interface, in four generic terms - a template to write to, a query to read with, a script to + * prepare, a list of shapes to read in - and each type says what they mean for it and checks them + * while the flow is edited, so a mistake shows up in the editor rather than halfway through a run. + */ +public interface StorageType { + + /** How flows, catalog entries and vault credentials name this type, e.g. "S3". */ + String getName(); + + String getDescription(); + + /** The fields of a person's own connection, as the vault form asks for them. */ + List credentialFields(); + + /** The shapes a read can come back in. The first is the default. */ + List viewShapes(); + + /** What a write accepts, for the kinds of value the target's port takes. */ + List writableKinds(); + + /** A connection's settings, checked when the catalog loads or a person saves their own. */ + void validateConnection(StorageConnection connection); + + /** Problems with a write template, in words for the flow author; empty when there are none. */ + List validateTemplate(String template); + + /** Problems with a read query for the shape asked. */ + List validateQuery(String query, ViewShape shape); + + /** Problems with a delete template. */ + List validateDelete(String template); + + List validateInit(StorageInit init); + + /** The placeholders of a template or query that bind as values rather than render as text. */ + default Set parameters(String templateOrQuery) { + return StoragePlaceholders.parameters(templateOrQuery); + } + + StorageSession open(StorageConnection connection, StorageScope scope); +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypeMetadata.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypeMetadata.java new file mode 100644 index 0000000..5ce69a9 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypeMetadata.java @@ -0,0 +1,26 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import java.util.List; + +/** + * What the editor and the vault form need to know about a storage type - never anything about a + * connection to one. + * + * @param viewShapes what a read can return, the default first + * @param writableKinds the kinds of value a write accepts (FILE, TEXT, JSON) + * @param credentialFields the fields of a person's own connection, for the vault form + */ +public record StorageTypeMetadata(String name, String description, List viewShapes, + List writableKinds, List credentialFields) { + + public static StorageTypeMetadata of(StorageType type) { + return new StorageTypeMetadata(type.getName(), type.getDescription(), + type.viewShapes().stream().map(Enum::name).toList(), + type.writableKinds().stream().map(kind -> kind.type().name()).distinct().toList(), + type.credentialFields()); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypes.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypes.java new file mode 100644 index 0000000..027818e --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageTypes.java @@ -0,0 +1,40 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import java.util.Comparator; +import java.util.List; +import java.util.Optional; + +import org.springframework.stereotype.Component; + +/** Every {@link StorageType} this service has, found by the name flows use for it. */ +@Component +public class StorageTypes { + + private final List types; + + public StorageTypes(List types) { + this.types = types.stream() + .sorted(Comparator.comparing(StorageType::getName, String.CASE_INSENSITIVE_ORDER)) + .toList(); + } + + public List all() { + return types; + } + + public Optional find(String name) { + if (name == null || name.isBlank()) { + return Optional.empty(); + } + return types.stream().filter(type -> type.getName().equalsIgnoreCase(name.trim())).findFirst(); + } + + public StorageType require(String name) { + return find(name).orElseThrow(() -> new StorageException(StorageException.CONNECTION_INVALID, + "Unknown storage type: " + name)); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/StorageWriteResult.java b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageWriteResult.java new file mode 100644 index 0000000..0815ca8 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/StorageWriteResult.java @@ -0,0 +1,17 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import tools.jackson.databind.JsonNode; + +/** + * What a write or a delete did. + * + * @param location where it went: the object key, or null for a statement + * @param affected objects or rows touched + * @param rows what a statement's {@code RETURNING} gave back, null when it had none + */ +public record StorageWriteResult(String location, long affected, JsonNode rows) { +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/ViewShape.java b/src/main/java/it/cnr/isti/workflow/manager/storage/ViewShape.java new file mode 100644 index 0000000..3fdaa7b --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/ViewShape.java @@ -0,0 +1,39 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import it.cnr.isti.workflow.manager.ios.IOType; + +/** + * What a read hands to the step that asked for it. A type offers the shapes that mean something for + * it - an object store has files, a database has rows - and a node asking for another is refused + * while the flow is edited. + */ +public enum ViewShape { + /** Each matching object as a file, for a node that takes files. */ + FILES(IOType.FILE, true), + /** Each matching object's content as text. */ + TEXTS(IOType.TEXT, true), + /** Only the names of what matched. */ + KEYS(IOType.TEXT, true), + /** One JSON value: the rows of a query, or the parsed content of the matching objects. */ + JSON(IOType.JSON, false); + + private final IOType portType; + private final boolean multiple; + + ViewShape(IOType portType, boolean multiple) { + this.portType = portType; + this.multiple = multiple; + } + + public IOType portType() { + return portType; + } + + public boolean multiple() { + return multiple; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresSession.java b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresSession.java new file mode 100644 index 0000000..bdfb200 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresSession.java @@ -0,0 +1,320 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage.postgres; + +import java.io.File; +import java.math.BigDecimal; +import java.math.BigInteger; +import java.sql.Array; +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.ResultSetMetaData; +import java.sql.SQLException; +import java.sql.Statement; +import java.sql.Types; +import java.util.ArrayList; +import java.util.Base64; +import java.util.Collection; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.function.Function; + +import it.cnr.isti.workflow.manager.storage.StorageConnection; +import it.cnr.isti.workflow.manager.storage.StorageException; +import it.cnr.isti.workflow.manager.storage.StorageInit; +import it.cnr.isti.workflow.manager.storage.StoragePlaceholders; +import it.cnr.isti.workflow.manager.storage.StorageReadResult; +import it.cnr.isti.workflow.manager.storage.StorageSession; +import it.cnr.isti.workflow.manager.storage.StorageSettings; +import it.cnr.isti.workflow.manager.storage.StorageWriteResult; +import it.cnr.isti.workflow.manager.storage.ViewShape; +import tools.jackson.databind.JsonNode; +import tools.jackson.databind.ObjectMapper; +import tools.jackson.databind.node.ArrayNode; +import tools.jackson.databind.node.ObjectNode; + +/** One connection, for one piece of work: every statement in a transaction of its own. */ +final class PostgresSession implements StorageSession { + + private final Connection jdbc; + private final StorageConnection connection; + private final StorageSettings settings; + private final ObjectMapper objectMapper; + + PostgresSession(Connection jdbc, StorageConnection connection, StorageSettings settings, ObjectMapper objectMapper) { + this.jdbc = jdbc; + this.connection = connection; + this.settings = settings; + this.objectMapper = objectMapper; + } + + @Override + public void initialize(StorageInit init) { + if (init == null || !init.hasScript()) { + return; + } + if (!connection.allowInit()) { + throw new StorageException(StorageException.INIT_FAILED, + connection.name() + " may not be prepared by flows: its init script cannot run"); + } + inTransaction(false, "run the init script", () -> { + try (Statement statement = jdbc.createStatement()) { + statement.execute(init.script()); + } + return null; + }); + } + + @Override + public StorageWriteResult write(String template, Object value, String contentType, Function values) { + requireWritable("write to"); + return change(template, PostgresStorageType.WRITE_VERBS, "write", StoragePlaceholders.withValue(value, values)); + } + + @Override + public StorageReadResult read(String query, ViewShape shape, Function values) { + SqlTemplate template = checked(query, PostgresStorageType.READ_VERBS, "read"); + return inTransaction(true, "read", () -> { + try (PreparedStatement statement = prepare(template, values)) { + statement.setMaxRows(settings.sqlMaxRows() + 1); + try (ResultSet rows = statement.executeQuery()) { + ArrayNode array = rowsOf(rows); + boolean truncated = array.size() > settings.sqlMaxRows(); + if (truncated) { + array.remove(array.size() - 1); + } + return new StorageReadResult(array, array.size(), truncated); + } + } + }); + } + + @Override + public StorageReadResult list(String query, Function values) { + String schema = query == null || query.isBlank() ? null : StoragePlaceholders.render(query, values).trim(); + return inTransaction(true, "list the tables", () -> { + try (PreparedStatement statement = jdbc.prepareStatement( + "SELECT table_name FROM information_schema.tables WHERE table_schema = coalesce(?, current_schema())" + + " AND table_type IN ('BASE TABLE', 'VIEW') ORDER BY table_name")) { + statement.setString(1, schema); + statement.setMaxRows(settings.maxKeys() + 1); + List tables = new ArrayList<>(); + try (ResultSet rows = statement.executeQuery()) { + while (rows.next()) { + tables.add(rows.getString(1)); + } + } + boolean truncated = tables.size() > settings.maxKeys(); + List kept = truncated ? tables.subList(0, settings.maxKeys()) : tables; + return new StorageReadResult(List.copyOf(kept), kept.size(), truncated); + } + }); + } + + @Override + public StorageWriteResult delete(String template, Function values) { + requireWritable("delete from"); + return change(template, PostgresStorageType.DELETE_VERBS, "delete", values); + } + + @Override + public void close() { + try { + jdbc.close(); + } catch (SQLException ignored) { + // Returned to the pool, or dropped: either way nothing more can be done with it. + } + } + + private StorageWriteResult change(String sql, Set verbs, String what, Function values) { + SqlTemplate template = checked(sql, verbs, what); + return inTransaction(false, what, () -> { + try (PreparedStatement statement = prepare(template, values)) { + if (statement.execute()) { + try (ResultSet rows = statement.getResultSet()) { + ArrayNode returned = rowsOf(rows); + return new StorageWriteResult(null, returned.size(), returned); + } + } + return new StorageWriteResult(null, Math.max(0, statement.getUpdateCount()), null); + } + }); + } + + private void requireWritable(String verb) { + if (connection.readOnly()) { + throw new StorageException(StorageException.OPERATION_NOT_ALLOWED, + connection.name() + " is read-only: nothing can " + verb + " it"); + } + } + + /** Checked again here, not only in the editor: a flow saved before a rule existed still meets it. */ + private SqlTemplate checked(String sql, Set verbs, String what) { + SqlTemplate template = SqlTemplate.parse(sql); + if (!template.problems().isEmpty()) { + throw new StorageException(StorageException.SQL_REJECTED, String.join("; ", template.problems())); + } + if (!verbs.contains(template.verb())) { + throw new StorageException(StorageException.SQL_REJECTED, + "A " + what + " must start with " + String.join(" or ", verbs.stream().sorted().toList()) + + ", not " + (template.verb().isEmpty() ? "nothing" : template.verb())); + } + return template; + } + + private PreparedStatement prepare(SqlTemplate template, Function values) throws SQLException { + PreparedStatement statement = jdbc.prepareStatement(template.sql()); + List parameters = template.parameters(); + for (int i = 0; i < parameters.size(); i++) { + bind(statement, i + 1, parameters.get(i), values == null ? null : values.apply(parameters.get(i))); + } + return statement; + } + + private void bind(PreparedStatement statement, int index, String name, Object value) throws SQLException { + if (value instanceof JsonNode node) { + if (node.isNull() || node.isMissingNode()) { + value = null; + } else if (node.isString()) { + value = node.asString(); + } else if (node.isBoolean()) { + value = node.asBoolean(); + } else if (node.isNumber()) { + value = node.numberValue(); + } else { + statement.setObject(index, node.toString(), Types.OTHER); + return; + } + } + if (value == null) { + if (!StoragePlaceholders.isValueName(name)) { + throw new StorageException(StorageException.PLACEHOLDER_MISSING, "No value for ${{" + name + "}}"); + } + statement.setNull(index, Types.OTHER); + } else if (value instanceof File) { + throw new StorageException(StorageException.SQL_REJECTED, + "${{" + name + "}} is a file, and a file cannot be bound into a statement"); + } else if (value instanceof Boolean flag) { + statement.setBoolean(index, flag); + } else if (value instanceof Integer || value instanceof Long || value instanceof Short) { + statement.setLong(index, ((Number) value).longValue()); + } else if (value instanceof BigInteger || value instanceof BigDecimal || value instanceof Double + || value instanceof Float) { + statement.setBigDecimal(index, new BigDecimal(value.toString())); + } else if (value instanceof Map || value instanceof Collection) { + statement.setObject(index, objectMapper.writeValueAsString(value), Types.OTHER); + } else { + statement.setString(index, value.toString()); + } + } + + private ArrayNode rowsOf(ResultSet rows) throws SQLException { + ArrayNode array = objectMapper.createArrayNode(); + ResultSetMetaData meta = rows.getMetaData(); + while (rows.next()) { + ObjectNode row = array.addObject(); + for (int column = 1; column <= meta.getColumnCount(); column++) { + row.set(meta.getColumnLabel(column), cell(rows, column, meta.getColumnTypeName(column))); + } + } + return array; + } + + private JsonNode cell(ResultSet rows, int column, String typeName) throws SQLException { + Object value = rows.getObject(column); + if (value == null) { + return objectMapper.nullNode(); + } + if ("json".equalsIgnoreCase(typeName) || "jsonb".equalsIgnoreCase(typeName)) { + try { + return objectMapper.readTree(rows.getString(column)); + } catch (Exception ex) { + return objectMapper.getNodeFactory().stringNode(rows.getString(column)); + } + } + if (value instanceof Array array) { + ArrayNode items = objectMapper.createArrayNode(); + for (Object item : (Object[]) array.getArray()) { + items.add(item == null ? objectMapper.nullNode() : objectMapper.valueToTree(item instanceof Number || item instanceof Boolean ? item : item.toString())); + } + return items; + } + if (value instanceof byte[] bytes) { + return objectMapper.getNodeFactory().stringNode(Base64.getEncoder().encodeToString(bytes)); + } + if (value instanceof Number || value instanceof Boolean) { + return objectMapper.valueToTree(value); + } + return objectMapper.getNodeFactory().stringNode(value.toString()); + } + + @FunctionalInterface + private interface Work { + T run() throws SQLException; + } + + /** One statement's transaction: read-only when asked, with the statement timeout set inside it. */ + private T inTransaction(boolean readOnly, String what, Work work) { + try { + jdbc.setAutoCommit(false); + jdbc.setReadOnly(readOnly); + try (Statement timeout = jdbc.createStatement()) { + timeout.execute("SET LOCAL statement_timeout = " + (settings.sqlStatementTimeoutSeconds() * 1000L)); + } + T result = work.run(); + if (readOnly) { + jdbc.rollback(); + } else { + jdbc.commit(); + } + return result; + } catch (SQLException ex) { + rollbackQuietly(); + throw failure(what, ex); + } catch (RuntimeException ex) { + rollbackQuietly(); + throw ex; + } finally { + try { + jdbc.setReadOnly(false); + jdbc.setAutoCommit(true); + } catch (SQLException ignored) { + // The connection is closed right after; a pooled one is reset by the pool. + } + } + } + + private void rollbackQuietly() { + try { + jdbc.rollback(); + } catch (SQLException ignored) { + // Already failed; the original error is the one worth reporting. + } + } + + private StorageException failure(String what, SQLException ex) { + String state = ex.getSQLState() == null ? "" : ex.getSQLState(); + if (state.equals("25006")) { + return new StorageException(StorageException.OPERATION_NOT_ALLOWED, + "A read cannot change data in " + connection.name() + ": " + ex.getMessage(), ex); + } + if (state.startsWith("08") || state.equals("57P01") || state.equals("53300")) { + return new StorageException(StorageException.UNREACHABLE, + "Could not " + what + " on " + connection.name() + ": " + ex.getMessage(), ex); + } + if (state.equals("57014")) { + return new StorageException(StorageException.SQL_REJECTED, "The " + what + " on " + connection.name() + + " took longer than " + settings.sqlStatementTimeoutSeconds() + " seconds and was stopped", ex); + } + if (state.startsWith("28") || state.equals("42501")) { + return new StorageException(StorageException.OPERATION_NOT_ALLOWED, + connection.name() + " refused to " + what + ": " + ex.getMessage(), ex); + } + return new StorageException(what.equals("run the init script") ? StorageException.INIT_FAILED : StorageException.SQL_REJECTED, + connection.name() + " refused the " + what + ": " + ex.getMessage(), ex); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageType.java b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageType.java new file mode 100644 index 0000000..98634f1 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageType.java @@ -0,0 +1,274 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage.postgres; + +import java.net.URLEncoder; +import java.nio.charset.StandardCharsets; +import java.sql.Connection; +import java.sql.DriverManager; +import java.sql.SQLException; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.regex.Pattern; + +import org.springframework.beans.factory.DisposableBean; +import org.springframework.stereotype.Component; + +import com.zaxxer.hikari.HikariConfig; +import com.zaxxer.hikari.HikariDataSource; + +import it.cnr.isti.workflow.manager.blocks.IOCapability; +import it.cnr.isti.workflow.manager.blocks.IOCapabilityType; +import it.cnr.isti.workflow.manager.storage.CredentialField; +import it.cnr.isti.workflow.manager.storage.StorageConnection; +import it.cnr.isti.workflow.manager.storage.StorageEndpointGuard; +import it.cnr.isti.workflow.manager.storage.StorageException; +import it.cnr.isti.workflow.manager.storage.StorageInit; +import it.cnr.isti.workflow.manager.storage.StoragePlaceholders; +import it.cnr.isti.workflow.manager.storage.StorageScope; +import it.cnr.isti.workflow.manager.storage.StorageSession; +import it.cnr.isti.workflow.manager.storage.StorageSettings; +import it.cnr.isti.workflow.manager.storage.StorageType; +import it.cnr.isti.workflow.manager.storage.ViewShape; +import tools.jackson.databind.ObjectMapper; + +/** + * A PostgreSQL database, spoken to in SQL. + * + *

    + *
  • A query is a {@code SELECT} (or {@code WITH}), run in a read-only transaction, returning its + * rows as JSON. + *
  • A template is an {@code INSERT}, {@code UPDATE} or {@code MERGE}; a delete is a + * {@code DELETE}. Both may end in {@code RETURNING}. + *
  • Init is a script - {@code CREATE TABLE IF NOT EXISTS ...} - run in one transaction. + *
+ * + *

Placeholders are bound, never spliced in: see {@link SqlTemplate}. The JDBC URL is built here + * from a host, a port and a database, and the driver's properties are set here too. Nothing a + * person types reaches either, because pgjdbc's URL parameters are themselves a way in: a + * {@code socketFactory} or {@code loggerFile} of the caller's choosing ran code or wrote files on + * the server (CVE-2022-21724). + */ +@Component +public class PostgresStorageType implements StorageType, DisposableBean { + + public static final String NAME = "PostgreSQL"; + + private static final Pattern HOST = Pattern.compile("^(?:[A-Za-z0-9](?:[A-Za-z0-9-]{0,61}[A-Za-z0-9])?)(?:\\.[A-Za-z0-9](?:[A-Za-z0-9-]{0,61}[A-Za-z0-9])?)*$|^\\[[0-9A-Fa-f:.]+]$"); + private static final Pattern DATABASE = Pattern.compile("^[A-Za-z0-9_$.-]{1,63}$"); + private static final Pattern IDENTIFIER = Pattern.compile("^[A-Za-z_][A-Za-z0-9_$]{0,62}$"); + private static final Set SSL_MODES = Set.of("disable", "allow", "prefer", "require", "verify-ca", "verify-full"); + static final Set READ_VERBS = Set.of("SELECT", "WITH", "VALUES", "TABLE"); + static final Set WRITE_VERBS = Set.of("INSERT", "UPDATE", "MERGE"); + static final Set DELETE_VERBS = Set.of("DELETE"); + + private final StorageSettings settings; + private final StorageEndpointGuard guard; + private final ObjectMapper objectMapper; + private final Map pools = new ConcurrentHashMap<>(); + + public PostgresStorageType(StorageSettings settings, StorageEndpointGuard guard, ObjectMapper objectMapper) { + this.settings = settings; + this.guard = guard; + this.objectMapper = objectMapper; + } + + @Override + public String getName() { + return NAME; + } + + @Override + public String getDescription() { + return "A PostgreSQL database. Reads are SELECT queries returning rows as JSON; writes are INSERT, UPDATE or " + + "MERGE statements; values are bound as parameters."; + } + + @Override + public List credentialFields() { + return List.of( + new CredentialField("host", "Host", "e.g. db.example.org", false, true), + new CredentialField("port", "Port", "5432 when empty", false, false), + new CredentialField("database", "Database", null, false, true), + new CredentialField("schema", "Schema", "public when empty", false, false), + new CredentialField("sslMode", "SSL mode", "disable, allow, prefer, require, verify-ca or verify-full; prefer when empty", false, false), + new CredentialField("user", "User", null, false, true), + new CredentialField("password", "Password", null, true, true)); + } + + @Override + public List viewShapes() { + return List.of(ViewShape.JSON); + } + + @Override + public List writableKinds() { + return List.of(new IOCapability(IOCapabilityType.JSON, false), new IOCapability(IOCapabilityType.TEXT, false)); + } + + @Override + public void validateConnection(StorageConnection connection) { + String host = connection.requireText("host"); + if (!HOST.matcher(host).matches()) { + throw invalid(connection, "the host " + host + " is not a host name or address"); + } + port(connection); + String database = connection.requireText("database"); + if (!DATABASE.matcher(database).matches()) { + throw invalid(connection, "the database name " + database + " is not allowed"); + } + String schema = connection.optionalText("schema"); + if (schema != null && !IDENTIFIER.matcher(schema).matches()) { + throw invalid(connection, "the schema " + schema + " is not a plain identifier"); + } + String sslMode = connection.optionalText("sslMode"); + if (sslMode != null && !SSL_MODES.contains(sslMode)) { + throw invalid(connection, "the SSL mode must be one of " + SSL_MODES); + } + connection.requireText("user"); + connection.requireText("password"); + guard.check(connection, null, "host"); + } + + @Override + public List validateTemplate(String template) { + return validateStatement(template, WRITE_VERBS, "A write", true); + } + + @Override + public List validateQuery(String query, ViewShape shape) { + List problems = validateStatement(query, READ_VERBS, "A read", false); + if (shape != null && shape != ViewShape.JSON) { + problems.add("A PostgreSQL read returns rows as JSON; it cannot return " + shape); + } + return problems; + } + + @Override + public List validateDelete(String template) { + return validateStatement(template, DELETE_VERBS, "A delete", false); + } + + @Override + public List validateInit(StorageInit init) { + List problems = new ArrayList<>(); + if (init == null) { + return problems; + } + if (init.createIfMissing()) { + problems.add("A PostgreSQL storage has nothing to create by itself; put CREATE TABLE IF NOT EXISTS ... in the init script"); + } + if (init.hasScript() && init.script().contains("${{")) { + problems.add("The init script cannot have placeholders: it runs as it is written, before any step"); + } + return problems; + } + + @Override + public StorageSession open(StorageConnection connection, StorageScope scope) { + validateConnection(connection); + try { + Connection jdbc = connection.operatorChosen() ? pool(connection).getConnection() : connect(connection); + return new PostgresSession(jdbc, connection, settings, objectMapper); + } catch (SQLException ex) { + throw new StorageException(StorageException.UNREACHABLE, + "Could not connect to " + connection.name() + ": " + ex.getMessage(), ex); + } + } + + @Override + public void destroy() { + pools.values().forEach(HikariDataSource::close); + pools.clear(); + } + + private static List validateStatement(String sql, Set verbs, String what, boolean valueAllowed) { + List problems = new ArrayList<>(); + if (sql == null || sql.isBlank()) { + problems.add(what + " needs a statement"); + return problems; + } + SqlTemplate template = SqlTemplate.parse(sql); + problems.addAll(template.problems()); + if (!verbs.contains(template.verb())) { + problems.add(what + " must start with " + String.join(" or ", verbs.stream().sorted().toList()) + + (template.verb().isEmpty() ? "" : ", not " + template.verb())); + } + StoragePlaceholders.invalidNames(sql).forEach(name -> problems.add( + "${{" + name + "}} is not a usable name: start with a letter, then letters, digits, '-', '_' or '.'")); + if (!valueAllowed && StoragePlaceholders.names(sql).stream().anyMatch(StoragePlaceholders::isValueName)) { + problems.add(what + " has no value: ${{value}} only means something when writing"); + } + return problems; + } + + private static StorageException invalid(StorageConnection connection, String problem) { + return new StorageException(StorageException.CONNECTION_INVALID, connection.name() + ": " + problem); + } + + private static int port(StorageConnection connection) { + String port = connection.optionalText("port"); + if (port == null) { + return 5432; + } + try { + int value = Integer.parseInt(port); + if (value < 1 || value > 65535) { + throw invalid(connection, "the port must be between 1 and 65535"); + } + return value; + } catch (NumberFormatException ex) { + throw invalid(connection, "the port " + port + " is not a number"); + } + } + + private String url(StorageConnection connection) { + return "jdbc:postgresql://" + connection.requireText("host") + ":" + port(connection) + "/" + + URLEncoder.encode(connection.requireText("database"), StandardCharsets.UTF_8); + } + + /** Every driver property this connection has - set here, from checked settings, and nowhere else. */ + private Properties properties(StorageConnection connection) { + Properties properties = new Properties(); + properties.setProperty("user", connection.requireText("user")); + properties.setProperty("password", connection.requireText("password")); + properties.setProperty("sslmode", connection.optionalText("sslMode") == null ? "prefer" : connection.optionalText("sslMode")); + properties.setProperty("ApplicationName", "humainflow-storage"); + properties.setProperty("connectTimeout", String.valueOf(settings.connectTimeoutSeconds())); + properties.setProperty("loginTimeout", String.valueOf(settings.connectTimeoutSeconds())); + // Strings go to the server untyped, so a placeholder bound as text fits an integer or a + // jsonb column the way a literal would, instead of failing as "character varying". + properties.setProperty("stringtype", "unspecified"); + String schema = connection.optionalText("schema"); + if (schema != null) { + properties.setProperty("currentSchema", schema); + } + return properties; + } + + private Connection connect(StorageConnection connection) throws SQLException { + return DriverManager.getConnection(url(connection), properties(connection)); + } + + private HikariDataSource pool(StorageConnection connection) { + String key = connection.name() + "\u0000" + connection.settings(); + return pools.computeIfAbsent(key, ignored -> { + HikariConfig config = new HikariConfig(); + config.setPoolName("storage-" + connection.name().replaceAll("[^A-Za-z0-9_-]", "_")); + config.setJdbcUrl(url(connection)); + config.setDataSourceProperties(properties(connection)); + config.setMaximumPoolSize(4); + config.setMinimumIdle(0); + config.setIdleTimeout(60_000); + config.setConnectionTimeout(settings.connectTimeoutSeconds() * 1000L); + config.setInitializationFailTimeout(-1); + return new HikariDataSource(config); + }); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplate.java b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplate.java new file mode 100644 index 0000000..b10e99a --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplate.java @@ -0,0 +1,224 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage.postgres; + +import java.util.ArrayList; +import java.util.List; +import java.util.Locale; + +/** + * A statement written with ${{name}} placeholders, turned into one the driver can bind. + * + *

Every placeholder becomes a {@code ?} and its name joins {@link #parameters()} in order, so + * each value is sent to the server as a value, never spliced into the text: a value of + * {@code '; DROP TABLE x; --} is stored as those characters. That only holds for a placeholder + * standing where a value can - so one found inside a string literal, a quoted identifier, a + * dollar-quoted body or a comment is refused, with the fix spelled out, rather than left there as + * literal text. A {@code ?} the author wrote themselves, the jsonb operator say, is doubled so the + * driver does not take it for a parameter. + * + *

One statement only. The driver does not see to that - pgjdbc splits {@code INSERT ...; DROP + * TABLE x} into two and runs both, prepared or not - so a second statement is refused here. + */ +final class SqlTemplate { + + private final String sql; + private final List parameters; + private final String verb; + private final List problems; + + private SqlTemplate(String sql, List parameters, String verb, List problems) { + this.sql = sql; + this.parameters = parameters; + this.verb = verb; + this.problems = problems; + } + + /** The statement to prepare. Meaningful only when {@link #problems()} is empty. */ + String sql() { + return sql; + } + + /** The placeholder bound to each {@code ?}, in order - a name used twice is bound twice. */ + List parameters() { + return parameters; + } + + /** The first keyword, upper case: SELECT, INSERT, WITH... Empty for a blank statement. */ + String verb() { + return verb; + } + + List problems() { + return problems; + } + + static SqlTemplate parse(String template) { + List problems = new ArrayList<>(); + List parameters = new ArrayList<>(); + StringBuilder out = new StringBuilder(); + String text = template == null ? "" : template; + int i = 0; + int length = text.length(); + while (i < length) { + char c = text.charAt(i); + if (text.startsWith("${{", i)) { + int end = text.indexOf("}}", i + 3); + if (end < 0) { + problems.add("A placeholder starting at character " + (i + 1) + " is never closed with }}"); + break; + } + parameters.add(text.substring(i + 3, end).trim()); + out.append('?'); + i = end + 2; + } else if (c == '\'') { + boolean escapes = i > 0 && (text.charAt(i - 1) == 'E' || text.charAt(i - 1) == 'e') + && (i < 2 || !Character.isLetterOrDigit(text.charAt(i - 2))); + i = copyQuoted(text, i, '\'', escapes, out, problems, "a string literal"); + } else if (c == '"') { + i = copyQuoted(text, i, '"', false, out, problems, "a quoted identifier"); + } else if (c == '$' && dollarTag(text, i) != null) { + String tag = dollarTag(text, i); + int close = text.indexOf(tag, i + tag.length()); + int stop = close < 0 ? length : close + tag.length(); + checkNoPlaceholder(text.substring(i, stop), problems, "a dollar-quoted body"); + out.append(text, i, stop); + i = stop; + } else if (text.startsWith("--", i)) { + int stop = text.indexOf('\n', i); + stop = stop < 0 ? length : stop; + checkNoPlaceholder(text.substring(i, stop), problems, "a comment"); + out.append(text, i, stop); + i = stop; + } else if (text.startsWith("/*", i)) { + int stop = blockCommentEnd(text, i); + checkNoPlaceholder(text.substring(i, stop), problems, "a comment"); + out.append(text, i, stop); + i = stop; + } else if (c == '?') { + out.append("??"); + i++; + } else if (c == ';') { + if (!onlyCommentsAndSpaceFrom(text, i + 1)) { + problems.add("Only one statement is allowed: remove what follows the ';' at character " + (i + 1)); + break; + } + // A closing ';' is harmless, and dropped: the driver would otherwise split on it. + i = length; + } else { + out.append(c); + i++; + } + } + return new SqlTemplate(out.toString(), List.copyOf(parameters), firstKeyword(text), List.copyOf(problems)); + } + + private static int copyQuoted(String text, int start, char quote, boolean backslashEscapes, StringBuilder out, + List problems, String what) { + int i = start + 1; + while (i < text.length()) { + char c = text.charAt(i); + if (backslashEscapes && c == '\\') { + i += 2; + continue; + } + if (c == quote) { + if (i + 1 < text.length() && text.charAt(i + 1) == quote) { + i += 2; + continue; + } + i++; + break; + } + i++; + } + int stop = Math.min(i, text.length()); + checkNoPlaceholder(text.substring(start, stop), problems, what); + out.append(text, start, stop); + return stop; + } + + private static void checkNoPlaceholder(String part, List problems, String where) { + int at = part.indexOf("${{"); + if (at >= 0) { + int end = part.indexOf("}}", at); + String placeholder = end < 0 ? part.substring(at) : part.substring(at, end + 2); + problems.add(placeholder + " is inside " + where + ", where it would stay literal text. Write it on its own, " + + "without quotes - the value is sent as a parameter: WHERE name = " + placeholder); + } + } + + /** $tag$ or $$ opening a dollar-quoted string at i, or null. */ + private static String dollarTag(String text, int i) { + int j = i + 1; + while (j < text.length() && (Character.isLetterOrDigit(text.charAt(j)) || text.charAt(j) == '_')) { + j++; + } + if (j < text.length() && text.charAt(j) == '$' && (j == i + 1 || !Character.isDigit(text.charAt(i + 1)))) { + return text.substring(i, j + 1); + } + return null; + } + + private static boolean onlyCommentsAndSpaceFrom(String text, int start) { + int i = start; + while (i < text.length()) { + if (Character.isWhitespace(text.charAt(i)) || text.charAt(i) == ';') { + i++; + } else if (text.startsWith("--", i)) { + int stop = text.indexOf('\n', i); + i = stop < 0 ? text.length() : stop; + } else if (text.startsWith("/*", i)) { + i = blockCommentEnd(text, i); + } else { + return false; + } + } + return true; + } + + /** Block comments nest in PostgreSQL. */ + private static int blockCommentEnd(String text, int start) { + int depth = 0; + int i = start; + while (i < text.length()) { + if (text.startsWith("/*", i)) { + depth++; + i += 2; + } else if (text.startsWith("*/", i)) { + depth--; + i += 2; + if (depth == 0) { + return i; + } + } else { + i++; + } + } + return text.length(); + } + + private static String firstKeyword(String text) { + int i = 0; + while (i < text.length()) { + char c = text.charAt(i); + if (Character.isWhitespace(c) || c == '(') { + i++; + } else if (text.startsWith("--", i)) { + int stop = text.indexOf('\n', i); + i = stop < 0 ? text.length() : stop; + } else if (text.startsWith("/*", i)) { + i = blockCommentEnd(text, i); + } else { + break; + } + } + int start = i; + while (i < text.length() && Character.isLetter(text.charAt(i))) { + i++; + } + return text.substring(start, i).toUpperCase(Locale.ROOT); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Keys.java b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Keys.java new file mode 100644 index 0000000..4a56a65 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Keys.java @@ -0,0 +1,136 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage.s3; + +import java.util.ArrayList; +import java.util.List; +import java.util.regex.Pattern; + +import it.cnr.isti.workflow.manager.storage.StorageException; +import it.cnr.isti.workflow.manager.storage.StoragePlaceholders; + +/** + * Object keys and globs, checked and matched. + * + *

A key is taken literally by S3, so {@code ..} is not a way up - but this service puts a prefix + * of its own in front of every key (the instance's root, an execution's scope), and a key that + * looks like it climbs out of that prefix is refused rather than trusted to be harmless everywhere. + * The checks run on the key as rendered, placeholders included, since a value coming from another + * step is exactly where a {@code ../} would come from. + */ +final class S3Keys { + + private static final Pattern BUCKET = Pattern.compile("^[a-z0-9][a-z0-9.-]{1,61}[a-z0-9]$"); + private static final int MAX_KEY_LENGTH = 1024; + + private S3Keys() { + } + + static boolean hasWildcard(String text) { + return text.indexOf('*') >= 0 || text.indexOf('?') >= 0; + } + + /** Problems with a key or glob as written, placeholders left in; empty when there are none. */ + static List problemsAsWritten(String text, String what) { + List problems = new ArrayList<>(); + if (text == null || text.isBlank()) { + problems.add(what + " is empty"); + return problems; + } + String literal = StoragePlaceholders.pattern().matcher(text).replaceAll("x"); + if (literal.startsWith("/")) { + problems.add(what + " must not start with '/': it is relative to the storage's own root"); + } + if (hasParentSegment(literal)) { + problems.add(what + " must not contain '..'"); + } + if (literal.contains("//")) { + problems.add(what + " must not contain an empty segment ('//')"); + } + StoragePlaceholders.invalidNames(text).forEach(name -> problems.add( + "${{" + name + "}} is not a usable name: start with a letter, then letters, digits, '-', '_' or '.'")); + return problems; + } + + /** The key as it will be used, after placeholders: refused if it could leave its prefix. */ + static String requireSafeKey(String key) { + if (key == null || key.isBlank()) { + throw new StorageException(StorageException.INVALID_LOCATION, "The object key is empty"); + } + if (key.startsWith("/") || hasParentSegment(key) || key.contains("//") || key.indexOf('\0') >= 0) { + throw new StorageException(StorageException.INVALID_LOCATION, + "The object key " + key + " is not allowed: no leading '/', no '..', no empty segment"); + } + if (key.length() > MAX_KEY_LENGTH) { + throw new StorageException(StorageException.INVALID_LOCATION, + "The object key is longer than " + MAX_KEY_LENGTH + " characters"); + } + return key; + } + + static boolean isValidBucket(String bucket) { + return bucket != null && BUCKET.matcher(bucket).matches() && !bucket.contains(".."); + } + + /** The part of a glob before its first wildcard: what to list, before matching. */ + static String listingPrefix(String glob) { + int wildcard = firstWildcard(glob); + return wildcard < 0 ? glob : glob.substring(0, wildcard); + } + + /** + * A glob as a regex: {@code *} stays inside one segment, {@code **} crosses them, {@code ?} is + * one character that is not '/'. + */ + static Pattern globToRegex(String glob) { + StringBuilder regex = new StringBuilder("^"); + for (int i = 0; i < glob.length(); i++) { + char c = glob.charAt(i); + if (c == '*') { + if (i + 1 < glob.length() && glob.charAt(i + 1) == '*') { + regex.append(".*"); + i++; + } else { + regex.append("[^/]*"); + } + } else if (c == '?') { + regex.append("[^/]"); + } else { + regex.append(Pattern.quote(String.valueOf(c))); + } + } + return Pattern.compile(regex.append('$').toString()); + } + + /** The root prefix as a directory: empty, or ending in exactly one '/'. */ + static String normalizePrefix(String prefix) { + if (prefix == null || prefix.isBlank()) { + return ""; + } + String trimmed = prefix.trim(); + while (trimmed.startsWith("/")) { + trimmed = trimmed.substring(1); + } + return trimmed.isEmpty() ? "" : (trimmed.endsWith("/") ? trimmed : trimmed + "/"); + } + + private static int firstWildcard(String glob) { + int star = glob.indexOf('*'); + int question = glob.indexOf('?'); + if (star < 0) { + return question; + } + return question < 0 ? star : Math.min(star, question); + } + + private static boolean hasParentSegment(String key) { + for (String segment : key.split("/", -1)) { + if (segment.equals("..")) { + return true; + } + } + return false; + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Session.java b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Session.java new file mode 100644 index 0000000..4b92ba5 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3Session.java @@ -0,0 +1,315 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage.s3; + +import java.io.ByteArrayInputStream; +import java.io.File; +import java.io.IOException; +import java.io.InputStream; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Map; +import java.util.function.Function; +import java.util.regex.Pattern; + +import io.minio.BucketExistsArgs; +import io.minio.GetObjectArgs; +import io.minio.ListObjectsArgs; +import io.minio.MakeBucketArgs; +import io.minio.MinioClient; +import io.minio.PutObjectArgs; +import io.minio.RemoveObjectArgs; +import io.minio.Result; +import io.minio.StatObjectArgs; +import io.minio.StatObjectResponse; +import io.minio.errors.ErrorResponseException; +import io.minio.messages.Item; +import it.cnr.isti.workflow.manager.storage.StorageConnection; +import it.cnr.isti.workflow.manager.storage.StorageException; +import it.cnr.isti.workflow.manager.storage.StorageInit; +import it.cnr.isti.workflow.manager.storage.StoragePlaceholders; +import it.cnr.isti.workflow.manager.storage.StorageReadResult; +import it.cnr.isti.workflow.manager.storage.StorageSession; +import it.cnr.isti.workflow.manager.storage.StorageSettings; +import it.cnr.isti.workflow.manager.storage.StorageWriteResult; +import it.cnr.isti.workflow.manager.storage.ViewShape; +import tools.jackson.databind.JsonNode; +import tools.jackson.databind.ObjectMapper; +import tools.jackson.databind.node.ArrayNode; + +/** + * One piece of work against one bucket. Keys the flow sees are relative to {@code prefix} - the + * instance's root plus the execution's scope - which is added on the way in and taken off on the + * way out, so a flow never learns, or depends on, where its data sits. + */ +final class S3Session implements StorageSession { + + private final MinioClient client; + private final StorageConnection connection; + private final String bucket; + private final String prefix; + private final StorageSettings settings; + private final ObjectMapper objectMapper; + + S3Session(MinioClient client, StorageConnection connection, String bucket, String prefix, StorageSettings settings, + ObjectMapper objectMapper) { + this.client = client; + this.connection = connection; + this.bucket = bucket; + this.prefix = prefix; + this.settings = settings; + this.objectMapper = objectMapper; + } + + @Override + public void initialize(StorageInit init) { + boolean exists = call("check the bucket " + bucket, () -> client.bucketExists(BucketExistsArgs.builder().bucket(bucket).build())); + if (exists) { + return; + } + if (init == null || !init.createIfMissing()) { + throw new StorageException(StorageException.INIT_FAILED, + "The bucket " + bucket + " of " + connection.name() + " does not exist; create it, or let the storage node create it if missing"); + } + if (!connection.allowInit()) { + throw new StorageException(StorageException.INIT_FAILED, + "The bucket " + bucket + " of " + connection.name() + " does not exist, and this storage may not be prepared by flows"); + } + call("create the bucket " + bucket, () -> { + client.makeBucket(MakeBucketArgs.builder().bucket(bucket).build()); + return null; + }); + } + + @Override + public StorageWriteResult write(String template, Object value, String contentType, Function values) { + requireWritable("write to"); + String key = S3Keys.requireSafeKey(StoragePlaceholders.render(template, values)); + Content content = contentOf(value, contentType); + call("write " + key, () -> { + try (InputStream stream = content.open()) { + client.putObject(PutObjectArgs.builder().bucket(bucket).object(prefix + key) + .stream(stream, content.size(), -1).contentType(content.contentType()).build()); + } + return null; + }); + return new StorageWriteResult(key, 1, null); + } + + @Override + public StorageReadResult read(String query, ViewShape shape, Function values) { + String rendered = S3Keys.requireSafeKey(StoragePlaceholders.render(query, values)); + Matches matches = match(rendered); + return switch (shape == null ? ViewShape.FILES : shape) { + case KEYS -> new StorageReadResult(matches.keys(), matches.keys().size(), matches.truncated()); + case FILES -> new StorageReadResult(downloadAll(matches.keys()), matches.keys().size(), matches.truncated()); + case TEXTS -> new StorageReadResult(textsOf(matches.keys()), matches.keys().size(), matches.truncated()); + case JSON -> new StorageReadResult(jsonOf(matches.keys()), matches.keys().size(), matches.truncated()); + }; + } + + @Override + public StorageReadResult list(String query, Function values) { + String rendered = query == null || query.isBlank() ? "**" : StoragePlaceholders.render(query, values); + Matches matches = match(S3Keys.requireSafeKey(rendered)); + return new StorageReadResult(matches.keys(), matches.keys().size(), matches.truncated()); + } + + @Override + public StorageWriteResult delete(String template, Function values) { + requireWritable("delete from"); + String key = S3Keys.requireSafeKey(StoragePlaceholders.render(template, values)); + if (S3Keys.hasWildcard(key)) { + throw new StorageException(StorageException.INVALID_LOCATION, "A delete names one object, not a glob: " + key); + } + boolean existed = stat(key) != null; + if (existed) { + call("delete " + key, () -> { + client.removeObject(RemoveObjectArgs.builder().bucket(bucket).object(prefix + key).build()); + return null; + }); + } + return new StorageWriteResult(key, existed ? 1 : 0, null); + } + + @Override + public void close() { + // The HTTP client is shared by every session and outlives this one. + } + + private void requireWritable(String verb) { + if (connection.readOnly()) { + throw new StorageException(StorageException.OPERATION_NOT_ALLOWED, + connection.name() + " is read-only: nothing can " + verb + " it"); + } + } + + private record Matches(List keys, boolean truncated) { + } + + /** A key names itself; a glob names what it matches, at most max-keys of it. */ + private Matches match(String keyOrGlob) { + if (!S3Keys.hasWildcard(keyOrGlob)) { + return new Matches(stat(keyOrGlob) == null ? List.of() : List.of(keyOrGlob), false); + } + Pattern pattern = S3Keys.globToRegex(keyOrGlob); + String listing = prefix + S3Keys.listingPrefix(keyOrGlob); + List keys = new ArrayList<>(); + boolean truncated = call("list " + keyOrGlob, () -> { + for (Result result : client.listObjects(ListObjectsArgs.builder().bucket(bucket).prefix(listing) + .recursive(true).build())) { + Item item = result.get(); + if (item.isDir()) { + continue; + } + String key = item.objectName().substring(prefix.length()); + if (!pattern.matcher(key).matches()) { + continue; + } + if (keys.size() == settings.maxKeys()) { + return true; + } + keys.add(key); + } + return false; + }); + return new Matches(List.copyOf(keys), truncated); + } + + private StatObjectResponse stat(String key) { + try { + return client.statObject(StatObjectArgs.builder().bucket(bucket).object(prefix + key).build()); + } catch (ErrorResponseException ex) { + if ("NoSuchKey".equals(ex.errorResponse().code()) || "NoSuchObject".equals(ex.errorResponse().code())) { + return null; + } + throw failure("read " + key, ex); + } catch (Exception ex) { + throw failure("read " + key, ex); + } + } + + private List downloadAll(List keys) { + List files = new ArrayList<>(); + for (String key : keys) { + byte[] bytes = bytesOf(key); + try { + String name = key.substring(key.lastIndexOf('/') + 1); + Path file = Files.createTempDirectory("storage-").resolve(name.isBlank() ? "object" : name); + Files.write(file, bytes); + files.add(file.toFile()); + } catch (IOException ex) { + throw new StorageException(StorageException.UNREACHABLE, "Could not keep " + key + " as a file: " + ex.getMessage(), ex); + } + } + return files; + } + + private List textsOf(List keys) { + return keys.stream().map(key -> new String(bytesOf(key), StandardCharsets.UTF_8)).toList(); + } + + private ArrayNode jsonOf(List keys) { + ArrayNode array = objectMapper.createArrayNode(); + for (String key : keys) { + try { + array.add(objectMapper.readTree(bytesOf(key))); + } catch (StorageException ex) { + throw ex; + } catch (Exception ex) { + throw new StorageException(StorageException.CONTENT_INVALID, + key + " is not JSON, so it cannot be read as JSON; read it as TEXTS or FILES instead"); + } + } + return array; + } + + private byte[] bytesOf(String key) { + StatObjectResponse stat = stat(key); + if (stat == null) { + throw new StorageException(StorageException.NOT_FOUND, key + " is not in " + connection.name()); + } + if (stat.size() > settings.maxObjectBytes()) { + throw new StorageException(StorageException.TOO_LARGE, key + " is " + stat.size() + + " bytes, more than the " + settings.maxObjectBytes() + " a read may take"); + } + return call("read " + key, () -> { + try (InputStream stream = client.getObject(GetObjectArgs.builder().bucket(bucket).object(prefix + key).build())) { + return stream.readAllBytes(); + } + }); + } + + private record Content(byte[] bytes, File file, String contentType) { + long size() throws IOException { + return file != null ? Files.size(file.toPath()) : bytes.length; + } + + InputStream open() throws IOException { + return file != null ? Files.newInputStream(file.toPath()) : new ByteArrayInputStream(bytes); + } + } + + private Content contentOf(Object value, String contentType) { + String declared = contentType == null || contentType.isBlank() ? null : contentType.trim(); + if (value instanceof File file) { + if (!file.isFile()) { + throw new StorageException(StorageException.NOT_FOUND, "The file to write, " + file.getName() + ", is not there"); + } + String type = declared; + if (type == null) { + try { + type = Files.probeContentType(file.toPath()); + } catch (IOException ignored) { + type = null; + } + } + return new Content(null, file, type == null ? "application/octet-stream" : type); + } + if (value instanceof JsonNode || value instanceof Map || value instanceof Collection) { + byte[] json = objectMapper.writeValueAsBytes(value); + return new Content(json, null, declared == null ? "application/json" : declared); + } + String text = value == null ? "" : String.valueOf(value); + return new Content(text.getBytes(StandardCharsets.UTF_8), null, declared == null ? "text/plain; charset=utf-8" : declared); + } + + @FunctionalInterface + private interface S3Call { + T run() throws Exception; + } + + private T call(String what, S3Call call) { + try { + return call.run(); + } catch (StorageException ex) { + throw ex; + } catch (Exception ex) { + throw failure(what, ex); + } + } + + private StorageException failure(String what, Exception ex) { + if (ex instanceof ErrorResponseException response) { + String code = response.errorResponse().code(); + String message = response.errorResponse().message(); + if ("NoSuchBucket".equals(code)) { + return new StorageException(StorageException.NOT_FOUND, "The bucket " + bucket + " of " + connection.name() + " does not exist", ex); + } + if ("AccessDenied".equals(code) || "InvalidAccessKeyId".equals(code) || "SignatureDoesNotMatch".equals(code)) { + return new StorageException(StorageException.OPERATION_NOT_ALLOWED, + connection.name() + " refused to " + what + ": " + code + (message == null ? "" : " - " + message), ex); + } + return new StorageException(StorageException.UNREACHABLE, + connection.name() + " could not " + what + ": " + code + (message == null ? "" : " - " + message), ex); + } + return new StorageException(StorageException.UNREACHABLE, + "Could not " + what + " on " + connection.name() + ": " + ex.getMessage(), ex); + } +} diff --git a/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageType.java b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageType.java new file mode 100644 index 0000000..c30de23 --- /dev/null +++ b/src/main/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageType.java @@ -0,0 +1,185 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage.s3; + +import java.net.URI; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.TimeUnit; + +import org.springframework.stereotype.Component; + +import io.minio.MinioClient; +import it.cnr.isti.workflow.manager.blocks.IOCapability; +import it.cnr.isti.workflow.manager.blocks.IOCapabilityType; +import it.cnr.isti.workflow.manager.storage.CredentialField; +import it.cnr.isti.workflow.manager.storage.StorageConnection; +import it.cnr.isti.workflow.manager.storage.StorageEndpointGuard; +import it.cnr.isti.workflow.manager.storage.StorageException; +import it.cnr.isti.workflow.manager.storage.StorageInit; +import it.cnr.isti.workflow.manager.storage.StoragePlaceholders; +import it.cnr.isti.workflow.manager.storage.StorageScope; +import it.cnr.isti.workflow.manager.storage.StorageSession; +import it.cnr.isti.workflow.manager.storage.StorageSettings; +import it.cnr.isti.workflow.manager.storage.StorageType; +import it.cnr.isti.workflow.manager.storage.ViewShape; +import okhttp3.OkHttpClient; +import tools.jackson.databind.ObjectMapper; + +/** + * An S3-compatible object store: MinIO, AWS S3, Ceph and the rest. + * + *

    + *
  • A template is an object key, placeholders rendered in: {@code reports/${{context.iteration}}.json}. + *
  • A query is a key or a glob - {@code pippo*}, {@code reports/**}{@code /*.json} - and a read + * returns every object it matches, in the shape asked. + *
  • Init can create the bucket when it is missing; there is no script. + *
+ */ +@Component +public class S3StorageType implements StorageType { + + public static final String NAME = "S3"; + + private final StorageSettings settings; + private final StorageEndpointGuard guard; + private final ObjectMapper objectMapper; + private final OkHttpClient httpClient; + + public S3StorageType(StorageSettings settings, StorageEndpointGuard guard, ObjectMapper objectMapper) { + this.settings = settings; + this.guard = guard; + this.objectMapper = objectMapper; + // One client, one connection pool, for every session: a MinioClient is cheap to build on it. + this.httpClient = new OkHttpClient.Builder() + .connectTimeout(settings.connectTimeoutSeconds(), TimeUnit.SECONDS) + .readTimeout(Math.max(60, settings.connectTimeoutSeconds()), TimeUnit.SECONDS) + .writeTimeout(Math.max(60, settings.connectTimeoutSeconds()), TimeUnit.SECONDS) + .build(); + } + + @Override + public String getName() { + return NAME; + } + + @Override + public String getDescription() { + return "An S3-compatible object store (MinIO, AWS S3, Ceph). Writes put an object under a key; " + + "reads take a key or a glob such as reports/*.json."; + } + + @Override + public List credentialFields() { + return List.of( + new CredentialField("endpoint", "Endpoint URL", "e.g. https://minio.example.org", false, true), + new CredentialField("region", "Region", "Leave empty unless the store asks for one", false, false), + new CredentialField("bucket", "Bucket", "The bucket this connection works in", false, true), + new CredentialField("rootPrefix", "Root prefix", "Optional folder every key goes under", false, false), + new CredentialField("accessKey", "Access key", null, false, true), + new CredentialField("secretKey", "Secret key", null, true, true)); + } + + @Override + public List viewShapes() { + return List.of(ViewShape.FILES, ViewShape.TEXTS, ViewShape.JSON, ViewShape.KEYS); + } + + @Override + public List writableKinds() { + return List.of(new IOCapability(IOCapabilityType.FILE, false), new IOCapability(IOCapabilityType.TEXT, false), + new IOCapability(IOCapabilityType.JSON, false)); + } + + @Override + public void validateConnection(StorageConnection connection) { + String endpoint = connection.requireText("endpoint"); + try { + URI uri = new URI(endpoint); + if (!"http".equalsIgnoreCase(uri.getScheme()) && !"https".equalsIgnoreCase(uri.getScheme()) + || uri.getHost() == null) { + throw new StorageException(StorageException.CONNECTION_INVALID, + connection.name() + ": the endpoint must be an http(s) URL, e.g. https://minio.example.org"); + } + } catch (java.net.URISyntaxException ex) { + throw new StorageException(StorageException.CONNECTION_INVALID, + connection.name() + ": the endpoint is not a URL: " + endpoint); + } + String bucket = connection.requireText("bucket"); + if (!S3Keys.isValidBucket(bucket)) { + throw new StorageException(StorageException.CONNECTION_INVALID, + connection.name() + ": " + bucket + " is not a valid bucket name (3-63 lowercase letters, digits, '.' or '-')"); + } + connection.requireText("accessKey"); + connection.requireText("secretKey"); + String rootPrefix = connection.optionalText("rootPrefix"); + if (rootPrefix != null && !S3Keys.problemsAsWritten(rootPrefix, "The root prefix").isEmpty()) { + throw new StorageException(StorageException.CONNECTION_INVALID, + connection.name() + ": " + String.join("; ", S3Keys.problemsAsWritten(rootPrefix, "the root prefix"))); + } + guard.check(connection, "endpoint", null); + } + + @Override + public List validateTemplate(String template) { + List problems = S3Keys.problemsAsWritten(template, "The object key"); + if (template != null && S3Keys.hasWildcard(template)) { + problems.add("The object key cannot contain '*' or '?': it names one object"); + } + if (template != null && StoragePlaceholders.names(template).contains(StoragePlaceholders.VALUE)) { + problems.add("${{value}} is the whole content being written and cannot be part of its key; use ${{value.}} for a field of a JSON value"); + } + return problems; + } + + @Override + public List validateQuery(String query, ViewShape shape) { + List problems = S3Keys.problemsAsWritten(query, "The key or glob"); + if (shape != null && !viewShapes().contains(shape)) { + problems.add("An S3 read cannot return " + shape); + } + if (query != null && StoragePlaceholders.names(query).stream().anyMatch(StoragePlaceholders::isValueName)) { + problems.add("A read has no value: ${{value}} only means something when writing"); + } + return problems; + } + + @Override + public List validateDelete(String template) { + List problems = S3Keys.problemsAsWritten(template, "The object key"); + if (template != null && S3Keys.hasWildcard(template)) { + problems.add("A delete names one object: '*' and '?' are not allowed"); + } + return problems; + } + + @Override + public List validateInit(StorageInit init) { + List problems = new ArrayList<>(); + if (init != null && init.hasScript()) { + problems.add("An S3 storage has no init script; leave it empty"); + } + return problems; + } + + @Override + public StorageSession open(StorageConnection connection, StorageScope scope) { + validateConnection(connection); + MinioClient.Builder builder = MinioClient.builder() + .endpoint(connection.requireText("endpoint")) + .credentials(connection.requireText("accessKey"), connection.requireText("secretKey")) + .httpClient(httpClient); + String region = connection.optionalText("region"); + if (region != null) { + builder.region(region); + } + MinioClient client = builder.build(); + if (connection.flag("pathStyle", false)) { + client.disableVirtualStyleEndpoint(); + } + String prefix = S3Keys.normalizePrefix(connection.optionalText("rootPrefix")) + S3Keys.normalizePrefix(scope.prefix()); + return new S3Session(client, connection, connection.requireText("bucket"), prefix, settings, objectMapper); + } +} diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties index a74e9cd..9de0a20 100644 --- a/src/main/resources/application.properties +++ b/src/main/resources/application.properties @@ -160,3 +160,16 @@ logging.level.root=ERROR # Jackson 3 changed default field visibility from ANY (Jackson 2) to PUBLIC_ONLY. # Restore to ANY to match Jackson 2 behavior for package-private fields in domain model. spring.jackson.visibility.field=any + +# Storage node and block. The catalog of storages offered to every flow, with their credentials - +# keep the file outside the image. Empty uses classpath:storages.json, which lists none. +app.storage.instances.file=${STORAGE_INSTANCES_FILE:} +# A person's own storage connection (from the vault) is checked like an LLM endpoint, behind its own +# switch: private and local addresses are refused unless allowed here. +app.storage.endpoint.allow-private-network=${STORAGE_ENDPOINT_ALLOW_PRIVATE_NETWORK:false} +app.storage.endpoint.allowed-host-patterns=${STORAGE_ENDPOINT_ALLOWED_HOST_PATTERNS:} +app.storage.max-object-bytes=${STORAGE_MAX_OBJECT_BYTES:33554432} +app.storage.max-keys=${STORAGE_MAX_KEYS:1000} +app.storage.sql.statement-timeout-seconds=${STORAGE_SQL_STATEMENT_TIMEOUT_SECONDS:30} +app.storage.sql.max-rows=${STORAGE_SQL_MAX_ROWS:1000} +app.storage.connect-timeout-seconds=${STORAGE_CONNECT_TIMEOUT_SECONDS:10} diff --git a/src/main/resources/storages.json b/src/main/resources/storages.json new file mode 100644 index 0000000..93044c0 --- /dev/null +++ b/src/main/resources/storages.json @@ -0,0 +1 @@ +{ "storages": [] } diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProviderTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProviderTest.java new file mode 100644 index 0000000..688688d --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/storage/StorageInstancesProviderTest.java @@ -0,0 +1,82 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.List; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; +import org.springframework.core.io.DefaultResourceLoader; + +import it.cnr.isti.workflow.manager.storage.postgres.PostgresStorageType; +import it.cnr.isti.workflow.manager.storage.s3.S3StorageType; +import tools.jackson.databind.ObjectMapper; + +/** The operator's catalog: loaded whole, or the service does not start. */ +class StorageInstancesProviderTest { + + @TempDir + Path folder; + + private final ObjectMapper mapper = new ObjectMapper(); + private final StorageEndpointGuard guard = new StorageEndpointGuard(false, ""); + private final StorageTypes types = new StorageTypes(List.of( + new S3StorageType(StorageSettings.defaults(), guard, mapper), + new PostgresStorageType(StorageSettings.defaults(), guard, mapper))); + + @Test + void withoutAFileTheCatalogIsEmpty() { + assertTrue(new StorageInstancesProvider(mapper, new DefaultResourceLoader(), types, "").all().isEmpty()); + } + + @Test + void aValidCatalogIsLoadedAndAnOperatorsPrivateAddressIsFine() throws Exception { + StorageInstancesProvider provider = load(""" + { "storages": [ + { "id": "flow-files", "name": "Flow files", "type": "S3", "allowInit": true, + "connection": { "endpoint": "http://127.0.0.1:9000", "bucket": "flow-files", + "accessKey": "a", "secretKey": "s" } }, + { "id": "reports-db", "type": "postgresql", "readOnly": true, + "connection": { "host": "10.0.0.5", "database": "reports", "user": "u", "password": "p" } } + ] }"""); + + assertEquals(List.of("flow-files", "reports-db"), provider.all().stream().map(StorageInstancesProvider.StorageInstance::id).toList()); + assertEquals(1, provider.ofType("PostgreSQL").size()); + assertTrue(provider.find("reports-db").orElseThrow().readOnly()); + assertFalse(provider.find("missing").isPresent()); + } + + @Test + void anEntryThatDoesNotHoldTogetherStopsTheServiceSayingWhich() { + IllegalStateException unknownType = assertThrows(IllegalStateException.class, () -> load(""" + { "storages": [ { "id": "x", "type": "Mongo", "connection": {} } ] }""")); + assertTrue(unknownType.getMessage().contains("(x) has unknown type Mongo"), unknownType.getMessage()); + + IllegalStateException badBucket = assertThrows(IllegalStateException.class, () -> load(""" + { "storages": [ { "id": "files", "type": "S3", "connection": + { "endpoint": "http://minio:9000", "bucket": "Bad_Bucket", "accessKey": "a", "secretKey": "s" } } ] }""")); + assertTrue(badBucket.getMessage().contains("(files)") && badBucket.getMessage().contains("not a valid bucket name"), + badBucket.getMessage()); + + IllegalStateException twice = assertThrows(IllegalStateException.class, () -> load(""" + { "storages": [ + { "id": "db", "type": "PostgreSQL", "connection": { "host": "db", "database": "d", "user": "u", "password": "p" } }, + { "id": "db", "type": "PostgreSQL", "connection": { "host": "db", "database": "d", "user": "u", "password": "p" } } ] }""")); + assertTrue(twice.getMessage().contains("id db is used twice"), twice.getMessage()); + } + + private StorageInstancesProvider load(String json) throws Exception { + Path file = folder.resolve("storages.json"); + Files.writeString(file, json); + return new StorageInstancesProvider(mapper, new DefaultResourceLoader(), types, file.toString()); + } +} diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageTypeTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageTypeTest.java new file mode 100644 index 0000000..c3f8c7c --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/PostgresStorageTypeTest.java @@ -0,0 +1,198 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage.postgres; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.List; +import java.util.Map; +import java.util.function.Function; + +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.postgresql.PostgreSQLContainer; + +import it.cnr.isti.workflow.manager.storage.StorageConnection; +import it.cnr.isti.workflow.manager.storage.StorageEndpointGuard; +import it.cnr.isti.workflow.manager.storage.StorageException; +import it.cnr.isti.workflow.manager.storage.StorageInit; +import it.cnr.isti.workflow.manager.storage.StorageReadResult; +import it.cnr.isti.workflow.manager.storage.StorageScope; +import it.cnr.isti.workflow.manager.storage.StorageSession; +import it.cnr.isti.workflow.manager.storage.StorageSettings; +import it.cnr.isti.workflow.manager.storage.StorageWriteResult; +import it.cnr.isti.workflow.manager.storage.ViewShape; +import tools.jackson.databind.JsonNode; +import tools.jackson.databind.ObjectMapper; + +/** Against a real PostgreSQL: values bound, reads kept read-only, and refusals said plainly. */ +@Testcontainers(disabledWithoutDocker = true) +class PostgresStorageTypeTest { + + @Container + static final PostgreSQLContainer POSTGRES = new PostgreSQLContainer("postgres:17-alpine"); + + private static final ObjectMapper MAPPER = new ObjectMapper(); + private static PostgresStorageType type; + + @BeforeAll + static void createType() { + type = new PostgresStorageType(StorageSettings.defaults(), new StorageEndpointGuard(false, ""), MAPPER); + try (StorageSession session = type.open(connection(false, true), StorageScope.shared())) { + session.initialize(new StorageInit(false, """ + CREATE TABLE IF NOT EXISTS rounds (id serial PRIMARY KEY, brief text NOT NULL, verdict text, score int, data jsonb); + CREATE TABLE IF NOT EXISTS notes (id serial PRIMARY KEY, body text);""")); + } + } + + @AfterAll + static void closePools() { + type.destroy(); + } + + @Test + void aWriteBindsTheFieldsOfItsValueAndAReadReturnsRowsAsJson() { + JsonNode round = MAPPER.readTree(""" + {"brief": "Build a list", "verdict": "rejected", "score": 3, "extra": {"k": [1, 2]}}"""); + try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) { + StorageWriteResult written = session.write( + "INSERT INTO rounds(brief, verdict, score, data) VALUES (${{value.brief}}, ${{value.verdict}}, ${{value.score}}, ${{value.extra}}) RETURNING id", + round, null, name -> null); + assertEquals(1, written.affected()); + assertTrue(written.rows().get(0).get("id").isNumber()); + + StorageReadResult read = session.read("SELECT brief, score, data FROM rounds WHERE verdict = ${{verdict}} ORDER BY id", + ViewShape.JSON, values(Map.of("verdict", "rejected"))); + JsonNode rows = (JsonNode) read.value(); + assertEquals("Build a list", rows.get(0).get("brief").asString()); + assertEquals(3, rows.get(0).get("score").asInt()); + assertEquals(2, rows.get(0).get("data").get("k").get(1).asInt(), "jsonb comes back as JSON, not a string"); + } + } + + @Test + void aValueThatLooksLikeSqlIsStoredAsTheCharactersItIs() { + String hostile = "x'); DROP TABLE notes; --"; + try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) { + session.write("INSERT INTO notes(body) VALUES (${{value}})", hostile, null, name -> null); + + JsonNode rows = (JsonNode) session.read("SELECT body FROM notes WHERE body = ${{body}}", ViewShape.JSON, + values(Map.of("body", hostile))).value(); + assertEquals(hostile, rows.get(0).get("body").asString()); + assertTrue(((List) session.list(null, name -> null).value()).contains("notes"), "the table is still there"); + } + } + + @Test + void aReadThatTriesToChangeDataIsStoppedByTheReadOnlyTransaction() { + try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) { + StorageException refusal = assertThrows(StorageException.class, () -> session.read( + "WITH gone AS (DELETE FROM notes RETURNING *) SELECT count(*) FROM gone", ViewShape.JSON, name -> null)); + assertEquals(StorageException.OPERATION_NOT_ALLOWED, refusal.getErrorCode()); + assertTrue(refusal.getMessage().contains("A read cannot change data"), refusal.getMessage()); + } + } + + @Test + void theVerbIsCheckedAtRunTimeTooAndTheServersReasonIsKept() { + try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) { + StorageException notASelect = assertThrows(StorageException.class, + () -> session.read("DELETE FROM notes", ViewShape.JSON, name -> null)); + assertEquals(StorageException.SQL_REJECTED, notASelect.getErrorCode()); + + StorageException missingTable = assertThrows(StorageException.class, + () -> session.read("SELECT * FROM nowhere", ViewShape.JSON, name -> null)); + assertEquals(StorageException.SQL_REJECTED, missingTable.getErrorCode()); + assertTrue(missingTable.getMessage().contains("relation \"nowhere\" does not exist"), missingTable.getMessage()); + + StorageException twoStatements = assertThrows(StorageException.class, + () -> session.read("SELECT 1; DROP TABLE notes", ViewShape.JSON, name -> null)); + assertEquals(StorageException.SQL_REJECTED, twoStatements.getErrorCode()); + assertTrue(twoStatements.getMessage().startsWith("Only one statement is allowed"), twoStatements.getMessage()); + + // The driver would have run both: it is the check before it that stops the second. + StorageException smuggled = assertThrows(StorageException.class, () -> session.write( + "INSERT INTO notes(body) VALUES (${{value}}); DROP TABLE notes", "x", null, name -> null)); + assertEquals(StorageException.SQL_REJECTED, smuggled.getErrorCode()); + assertTrue(((List) session.list(null, name -> null).value()).contains("notes")); + } + } + + @Test + void aDeleteReturnsHowManyRowsWentAndAReadOnlyStorageRefusesToChange() { + try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) { + session.write("INSERT INTO notes(body) VALUES (${{value}})", "to delete", null, name -> null); + StorageWriteResult deleted = session.delete("DELETE FROM notes WHERE body = ${{body}}", + values(Map.of("body", "to delete"))); + assertEquals(1, deleted.affected()); + assertNull(deleted.rows()); + } + try (StorageSession readOnly = type.open(connection(true, false), StorageScope.shared())) { + StorageException refusal = assertThrows(StorageException.class, + () -> readOnly.write("INSERT INTO notes(body) VALUES (${{value}})", "no", null, name -> null)); + assertEquals(StorageException.OPERATION_NOT_ALLOWED, refusal.getErrorCode()); + } + } + + @Test + void aMissingParameterIsNamed() { + try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) { + StorageException refusal = assertThrows(StorageException.class, + () -> session.read("SELECT * FROM notes WHERE body = ${{body}}", ViewShape.JSON, name -> null)); + assertEquals(StorageException.PLACEHOLDER_MISSING, refusal.getErrorCode()); + assertTrue(refusal.getMessage().contains("${{body}}"), refusal.getMessage()); + } + } + + @Test + void aPersonsOwnConnectionToAPrivateAddressIsRefusedAnOperatorsIsNot() { + StorageConnection personal = new StorageConnection("your PostgreSQL connection", PostgresStorageType.NAME, + connection(false, false).settings(), false, true, false); + StorageException refusal = assertThrows(StorageException.class, () -> type.open(personal, StorageScope.shared())); + assertEquals(StorageException.CONNECTION_INVALID, refusal.getErrorCode()); + assertTrue(refusal.getMessage().contains("private, local or otherwise disallowed"), refusal.getMessage()); + } + + @Test + void anInitScriptNeedsAStorageThatMayBePrepared() { + try (StorageSession session = type.open(connection(false, false), StorageScope.shared())) { + StorageException refusal = assertThrows(StorageException.class, + () -> session.initialize(new StorageInit(false, "CREATE TABLE x (id int)"))); + assertEquals(StorageException.INIT_FAILED, refusal.getErrorCode()); + } + } + + @Test + void editorChecksSayWhatIsWrong() { + assertTrue(type.validateQuery("INSERT INTO t VALUES (1)", ViewShape.JSON).getFirst().contains("must start with")); + assertTrue(type.validateQuery("SELECT 1", ViewShape.FILES).getFirst().contains("cannot return FILES")); + assertTrue(type.validateTemplate("UPDATE t SET a = '${{value}}'").getFirst().contains("inside a string literal")); + assertTrue(type.validateTemplate("INSERT INTO t(a) VALUES (${{value.a}})").isEmpty()); + assertTrue(type.validateDelete("DELETE FROM t WHERE id = ${{value}}").getFirst().contains("has no value")); + assertTrue(type.validateInit(new StorageInit(false, "CREATE TABLE ${{t}} (id int)")).getFirst().contains("cannot have placeholders")); + } + + private static Function values(Map values) { + return values::get; + } + + private static StorageConnection connection(boolean readOnly, boolean allowInit) { + JsonNode settings = MAPPER.valueToTree(Map.of( + "host", POSTGRES.getHost(), + "port", String.valueOf(POSTGRES.getFirstMappedPort()), + "database", POSTGRES.getDatabaseName(), + "user", POSTGRES.getUsername(), + "password", POSTGRES.getPassword(), + "sslMode", "disable")); + return new StorageConnection("test database" + (readOnly ? " (read-only)" : "") + (allowInit ? " (init)" : ""), + PostgresStorageType.NAME, settings, readOnly, allowInit, true); + } +} diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplateTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplateTest.java new file mode 100644 index 0000000..b8b5f9e --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/storage/postgres/SqlTemplateTest.java @@ -0,0 +1,81 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage.postgres; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.List; + +import org.junit.jupiter.api.Test; + +/** Placeholders become bound parameters, and only where a value can stand. */ +class SqlTemplateTest { + + @Test + void placeholdersBecomeParametersInOrderAndAUsedTwiceNameIsBoundTwice() { + SqlTemplate template = SqlTemplate.parse( + "INSERT INTO rounds(brief, verdict, again) VALUES (${{value.brief}}, ${{ value.verdict }}, ${{value.brief}})"); + + assertEquals("INSERT INTO rounds(brief, verdict, again) VALUES (?, ?, ?)", template.sql()); + assertEquals(List.of("value.brief", "value.verdict", "value.brief"), template.parameters()); + assertEquals("INSERT", template.verb()); + assertTrue(template.problems().isEmpty()); + } + + @Test + void aQuestionMarkTheAuthorWroteIsNotTakenForAParameter() { + SqlTemplate template = SqlTemplate.parse("SELECT * FROM docs WHERE data ? 'tag' AND id = ${{id}}"); + + assertEquals("SELECT * FROM docs WHERE data ?? 'tag' AND id = ?", template.sql()); + assertEquals(List.of("id"), template.parameters()); + } + + @Test + void aPlaceholderInsideAStringLiteralIsRefusedWithTheFix() { + SqlTemplate template = SqlTemplate.parse("SELECT * FROM t WHERE name = '${{name}}'"); + + assertEquals(1, template.problems().size()); + String problem = template.problems().getFirst(); + assertTrue(problem.contains("${{name}} is inside a string literal"), problem); + assertTrue(problem.contains("WHERE name = ${{name}}"), problem); + } + + @Test + void placeholdersInCommentsQuotedIdentifiersAndDollarBodiesAreRefusedToo() { + assertTrue(SqlTemplate.parse("SELECT 1 -- ${{x}}").problems().getFirst().contains("inside a comment")); + assertTrue(SqlTemplate.parse("SELECT 1 /* a /* nested ${{x}} */ */").problems().getFirst().contains("inside a comment")); + assertTrue(SqlTemplate.parse("SELECT \"${{x}}\" FROM t").problems().getFirst().contains("quoted identifier")); + assertTrue(SqlTemplate.parse("SELECT $body$ ${{x}} $body$").problems().getFirst().contains("dollar-quoted")); + } + + @Test + void quotesAreSkippedWithTheirEscapesSoWhatFollowsIsStillSeen() { + SqlTemplate template = SqlTemplate.parse("SELECT 'it''s', E'a\\'b' FROM t WHERE id = ${{id}}"); + + assertTrue(template.problems().isEmpty(), template.problems().toString()); + assertEquals(List.of("id"), template.parameters()); + assertEquals("SELECT 'it''s', E'a\\'b' FROM t WHERE id = ?", template.sql()); + } + + @Test + void theVerbIsTheFirstKeywordPastCommentsAndParentheses() { + assertEquals("SELECT", SqlTemplate.parse(" -- rows\n /* all */ (SELECT 1)").verb()); + assertEquals("WITH", SqlTemplate.parse("with x as (select 1) select * from x").verb()); + assertEquals("", SqlTemplate.parse(" ").verb()); + } + + @Test + void oneStatementOnlyThoughAClosingSemicolonIsFine() { + SqlTemplate closed = SqlTemplate.parse("SELECT * FROM t WHERE id = ${{id}}; -- done\n"); + assertTrue(closed.problems().isEmpty(), closed.problems().toString()); + assertEquals("SELECT * FROM t WHERE id = ?", closed.sql().trim()); + + SqlTemplate two = SqlTemplate.parse("INSERT INTO t VALUES (${{value}}); DROP TABLE t"); + assertTrue(two.problems().getFirst().startsWith("Only one statement is allowed"), two.problems().toString()); + + assertTrue(SqlTemplate.parse("SELECT ';' AS semicolon").problems().isEmpty(), "a ';' in a literal is text"); + } +} diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3KeysTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3KeysTest.java new file mode 100644 index 0000000..5bd741a --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3KeysTest.java @@ -0,0 +1,67 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage.s3; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.regex.Pattern; + +import org.junit.jupiter.api.Test; + +import it.cnr.isti.workflow.manager.storage.StorageException; + +class S3KeysTest { + + @Test + void aSingleStarStaysInsideOneFolderAndADoubleStarCrossesThem() { + Pattern pippo = S3Keys.globToRegex("pippo*"); + assertTrue(pippo.matcher("pippo.txt").matches()); + assertTrue(pippo.matcher("pippo").matches()); + assertFalse(pippo.matcher("pippo/inside.txt").matches()); + assertFalse(pippo.matcher("a-pippo.txt").matches()); + + Pattern deep = S3Keys.globToRegex("reports/**.json"); + assertTrue(deep.matcher("reports/2026/09/r.json").matches()); + assertFalse(deep.matcher("reports/r.txt").matches()); + + assertTrue(S3Keys.globToRegex("r?.md").matcher("r1.md").matches()); + assertFalse(S3Keys.globToRegex("r?.md").matcher("r/.md").matches()); + assertTrue(S3Keys.globToRegex("a+b(1).txt").matcher("a+b(1).txt").matches(), "regex characters are literal"); + } + + @Test + void aGlobIsListedFromItsFixedPart() { + assertEquals("reports/2026-", S3Keys.listingPrefix("reports/2026-*.json")); + assertEquals("", S3Keys.listingPrefix("*.json")); + assertEquals("exact.txt", S3Keys.listingPrefix("exact.txt")); + } + + @Test + void aKeyThatCouldLeaveItsPrefixIsRefusedAfterPlaceholdersAreIn() { + assertEquals("notes/1.txt", S3Keys.requireSafeKey("notes/1.txt")); + for (String key : new String[] { "../other/secret", "notes/../../x", "/absolute", "a//b", "" }) { + StorageException refusal = assertThrows(StorageException.class, () -> S3Keys.requireSafeKey(key), key); + assertEquals(StorageException.INVALID_LOCATION, refusal.getErrorCode()); + } + } + + @Test + void theKeyAsWrittenIsCheckedWithItsPlaceholdersStandingInForValues() { + assertTrue(S3Keys.problemsAsWritten("reports/${{context.iteration}}.json", "The key").isEmpty()); + assertFalse(S3Keys.problemsAsWritten("../${{x}}", "The key").isEmpty()); + assertFalse(S3Keys.problemsAsWritten("r/${{1bad}}", "The key").isEmpty()); + } + + @Test + void bucketNamesFollowTheS3Rules() { + assertTrue(S3Keys.isValidBucket("flow-data")); + assertFalse(S3Keys.isValidBucket("Flow_Data")); + assertFalse(S3Keys.isValidBucket("ab")); + assertFalse(S3Keys.isValidBucket("a..b")); + } +} diff --git a/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageTypeTest.java b/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageTypeTest.java new file mode 100644 index 0000000..1400740 --- /dev/null +++ b/src/test/java/it/cnr/isti/workflow/manager/storage/s3/S3StorageTypeTest.java @@ -0,0 +1,162 @@ +// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii - ISTI-CNR +// SPDX-License-Identifier: AGPL-3.0-or-later +// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM. + +package it.cnr.isti.workflow.manager.storage.s3; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.File; +import java.nio.file.Files; +import java.util.List; +import java.util.Map; + +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.testcontainers.containers.MinIOContainer; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; + +import it.cnr.isti.workflow.manager.storage.StorageConnection; +import it.cnr.isti.workflow.manager.storage.StorageEndpointGuard; +import it.cnr.isti.workflow.manager.storage.StorageException; +import it.cnr.isti.workflow.manager.storage.StorageInit; +import it.cnr.isti.workflow.manager.storage.StorageReadResult; +import it.cnr.isti.workflow.manager.storage.StorageScope; +import it.cnr.isti.workflow.manager.storage.StorageSession; +import it.cnr.isti.workflow.manager.storage.StorageSettings; +import it.cnr.isti.workflow.manager.storage.ViewShape; +import tools.jackson.databind.JsonNode; +import tools.jackson.databind.ObjectMapper; + +/** Against a real MinIO: keys, globs, scopes and every shape a read can take. */ +@Testcontainers(disabledWithoutDocker = true) +class S3StorageTypeTest { + + @Container + static final MinIOContainer MINIO = new MinIOContainer("minio/minio:latest"); + + private static final ObjectMapper MAPPER = new ObjectMapper(); + private static S3StorageType type; + + @BeforeAll + static void createBucket() { + type = new S3StorageType(StorageSettings.defaults(), new StorageEndpointGuard(false, ""), MAPPER); + try (StorageSession session = type.open(connection("flow-data", false, true), StorageScope.shared())) { + session.initialize(new StorageInit(true, null)); + } + } + + @Test + void whatIsWrittenCanBeReadBackAsFilesTextsJsonOrKeys() throws Exception { + File report = Files.createTempFile("report", ".md").toFile(); + Files.writeString(report.toPath(), "# Round 1"); + try (StorageSession session = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-shapes"))) { + assertEquals("pippo-1.md", session.write("pippo-${{context.iteration}}.md", report, null, + Map.of("context.iteration", 1)::get).location()); + session.write("pippo-2.json", MAPPER.readTree("{\"round\": 2}"), null, name -> null); + session.write("notes/pippo-3.txt", "not at the top level", null, name -> null); + session.write("other.txt", "not pippo", null, name -> null); + + StorageReadResult keys = session.read("pippo*", ViewShape.KEYS, name -> null); + assertEquals(List.of("pippo-1.md", "pippo-2.json"), keys.value(), "one * stays in its folder"); + + StorageReadResult deep = session.read("**pippo*", ViewShape.KEYS, name -> null); + assertEquals(3, deep.count()); + + @SuppressWarnings("unchecked") + List files = (List) session.read("pippo-1.md", ViewShape.FILES, name -> null).value(); + assertEquals("pippo-1.md", files.getFirst().getName()); + assertEquals("# Round 1", Files.readString(files.getFirst().toPath())); + + assertEquals(List.of("not at the top level"), session.read("notes/*", ViewShape.TEXTS, name -> null).value()); + + JsonNode json = (JsonNode) session.read("pippo-2.json", ViewShape.JSON, name -> null).value(); + assertEquals(2, json.get(0).get("round").asInt()); + } + } + + @Test + void aScopePerExecutionSeesOnlyWhatThatExecutionWrote() { + try (StorageSession first = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-a"))) { + first.write("brief.txt", "from a", null, name -> null); + } + try (StorageSession second = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-b"))) { + assertEquals(0, second.read("brief.txt", ViewShape.TEXTS, name -> null).count()); + second.write("brief.txt", "from b", null, name -> null); + assertEquals(List.of("from b"), second.read("*", ViewShape.TEXTS, name -> null).value()); + } + try (StorageSession shared = type.open(connection("flow-data", false, false), StorageScope.shared())) { + assertTrue(((List) shared.list("run-*/brief.txt", name -> null).value()).containsAll(List.of("run-a/brief.txt", "run-b/brief.txt"))); + } + } + + @Test + void aReadThatMatchesNothingIsEmptyRatherThanAFailure() { + try (StorageSession session = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-empty"))) { + assertEquals(0, session.read("notes/*.json", ViewShape.JSON, name -> null).count()); + assertEquals(0, session.read("first-round.txt", ViewShape.FILES, name -> null).count()); + } + } + + @Test + void aKeyFromAValueCannotClimbOutOfItsScope() { + try (StorageSession session = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-safe"))) { + StorageException refusal = assertThrows(StorageException.class, () -> session.write("${{name}}", "x", null, + Map.of("name", "../run-a/brief.txt")::get)); + assertEquals(StorageException.INVALID_LOCATION, refusal.getErrorCode()); + } + } + + @Test + void deleteRemovesOneObjectAndSaysWhetherItWasThere() { + try (StorageSession session = type.open(connection("flow-data", false, false), StorageScope.perExecution("run-delete"))) { + session.write("gone.txt", "bye", null, name -> null); + assertEquals(1, session.delete("gone.txt", name -> null).affected()); + assertEquals(0, session.delete("gone.txt", name -> null).affected()); + } + } + + @Test + void aReadOnlyStorageRefusesWritesAndAMissingBucketIsSaidAtInit() { + try (StorageSession readOnly = type.open(connection("flow-data", true, false), StorageScope.shared())) { + assertEquals(StorageException.OPERATION_NOT_ALLOWED, + assertThrows(StorageException.class, () -> readOnly.write("x.txt", "x", null, name -> null)).getErrorCode()); + } + try (StorageSession missing = type.open(connection("not-created", false, false), StorageScope.shared())) { + StorageException refusal = assertThrows(StorageException.class, () -> missing.initialize(new StorageInit(true, null))); + assertEquals(StorageException.INIT_FAILED, refusal.getErrorCode()); + assertTrue(refusal.getMessage().contains("may not be prepared by flows"), refusal.getMessage()); + } + } + + @Test + void wrongCredentialsAreSaidAsARefusal() { + StorageConnection wrong = new StorageConnection("test store", S3StorageType.NAME, MAPPER.valueToTree(Map.of( + "endpoint", MINIO.getS3URL(), "bucket", "flow-data", "accessKey", "nobody", "secretKey", "wrong-secret")), + false, false, true); + try (StorageSession session = type.open(wrong, StorageScope.shared())) { + StorageException refusal = assertThrows(StorageException.class, () -> session.read("*", ViewShape.KEYS, name -> null)); + assertEquals(StorageException.OPERATION_NOT_ALLOWED, refusal.getErrorCode()); + } + } + + @Test + void editorChecksSayWhatIsWrong() { + assertTrue(type.validateTemplate("reports/*.json").stream().anyMatch(p -> p.contains("cannot contain '*'"))); + assertTrue(type.validateTemplate("${{value}}.txt").stream().anyMatch(p -> p.contains("cannot be part of its key"))); + assertTrue(type.validateTemplate("reports/${{value.id}}.json").isEmpty()); + assertTrue(type.validateQuery("../x", ViewShape.TEXTS).stream().anyMatch(p -> p.contains("'..'"))); + assertTrue(type.validateDelete("old/*").stream().anyMatch(p -> p.contains("names one object"))); + assertTrue(type.validateInit(new StorageInit(false, "CREATE TABLE t (id int)")).getFirst().contains("no init script")); + } + + private static StorageConnection connection(String bucket, boolean readOnly, boolean allowInit) { + return new StorageConnection("test store", S3StorageType.NAME, MAPPER.valueToTree(Map.of( + "endpoint", MINIO.getS3URL(), "bucket", bucket, "rootPrefix", "flows", + "accessKey", MINIO.getUserName(), "secretKey", MINIO.getPassword())), + readOnly, allowInit, true); + } +}