From 1e7eecb1965606ddfb81c38df8422f9def3a930f Mon Sep 17 00:00:00 2001 From: Hendrik Ebbers Date: Fri, 25 Sep 2026 08:46:49 +0200 Subject: [PATCH] feat(storage): add the spring-services-storage module MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Moves the ObjectStore API and its three implementations — S3, file system and in-memory — into a new reactor module, so the code that so far lived inside one application becomes part of the library. This is the first step and deliberately only that: the sources are in, the reactor builds them, and nothing about their Spring-facing design has been decided yet. What that leaves open is recorded in docs/TODO.md. - Package renamed from com.openelements.api.storage to the reactor's com.openelements.spring.base.services.storage. - Module wired into the reactor, spring-services-all and the BOM. - software.amazon.awssdk:s3 added, version pinned in the root POM next to the other third-party versions java-parent does not manage. - Javadoc completed on every public member: CI builds with -Dmaven.javadoc.failOnWarnings=true, and the moved sources carried 17 missing-comment warnings. - ObjectStoreException's cause and InMemoryObjectStore.contentOf's return marked @Nullable — both are null in practice, which the packages' @NullMarked contract otherwise forbids. Co-Authored-By: Claude Opus 5 (1M context) --- README.md | 7 +- docs/TODO.md | 31 ++ pom.xml | 2 + spring-services-all/pom.xml | 5 + spring-services-bom/pom.xml | 5 + spring-services-storage/pom.xml | 37 ++ .../storage/ObjectNotFoundException.java | 25 ++ .../base/services/storage/ObjectStore.java | 84 ++++ .../storage/ObjectStoreException.java | 17 + .../base/services/storage/StoredObject.java | 28 ++ .../storage/file/FileObjectStore.java | 378 ++++++++++++++++++ .../services/storage/file/package-info.java | 7 + .../storage/memory/InMemoryObjectStore.java | 146 +++++++ .../services/storage/memory/package-info.java | 7 + .../base/services/storage/package-info.java | 7 + .../base/services/storage/s3/S3Config.java | 100 +++++ .../services/storage/s3/S3ObjectStore.java | 207 ++++++++++ .../services/storage/s3/package-info.java | 7 + 18 files changed, 1099 insertions(+), 1 deletion(-) create mode 100644 spring-services-storage/pom.xml create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectNotFoundException.java create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectStore.java create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectStoreException.java create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/StoredObject.java create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/file/FileObjectStore.java create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/file/package-info.java create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/memory/InMemoryObjectStore.java create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/memory/package-info.java create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/package-info.java create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/S3Config.java create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/S3ObjectStore.java create mode 100644 spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/package-info.java diff --git a/README.md b/README.md index d21777e..14cace6 100644 --- a/README.md +++ b/README.md @@ -32,7 +32,7 @@ coordinate is the reactor parent (a `pom`, no classes) — depend on one of the **À la carte:** import the BOM once, then declare `spring-services-core` plus only the feature modules you need (`spring-services-slack`, `spring-services-mcp`, `spring-services-email`, `spring-services-search`, `spring-services-dbbackup`, `spring-services-scim`, -`spring-services-tenant`) without versions: +`spring-services-tenant`, `spring-services-storage`) without versions: ```xml @@ -353,6 +353,7 @@ spring-services/ — reactor parent (packaging=pom) ├── spring-services-dbbackup — db-backup sidecar client (RestClient, no extra dep) ├── spring-services-scim — SCIM 2.0 Users provider (opt-in via openelements.scim.token) ├── spring-services-tenant — row-level multi-tenancy (self-activates on the classpath) +├── spring-services-storage — object store: S3, file system, in-memory (→ AWS SDK v2) ├── spring-services-all — everything bundle (depends on all modules; no config of its own) └── spring-services-bom — bill of materials for lockstep versioning ``` @@ -360,6 +361,10 @@ spring-services/ — reactor parent (packaging=pom) Each optional feature module ships its own `@AutoConfiguration` guarded by `@ConditionalOnClass`, so it self-activates when present and never pulls its heavy dependency into a consumer that skips it. +`spring-services-storage` is the one exception so far: it ships the `ObjectStore` implementations but +no auto-configuration, because an application has to choose one of them — declaring the implementation +it wants as a bean is that choice. Picking by classpath order would make it by accident. + ## Release Process ### SNAPSHOT Publishing (automatic) diff --git a/docs/TODO.md b/docs/TODO.md index 566a652..caea1a9 100644 --- a/docs/TODO.md +++ b/docs/TODO.md @@ -1,5 +1,36 @@ # TODO +## Finish the storage module (spring-services-storage) + +The module currently holds the `ObjectStore` API and its three implementations, lifted from an +application that used them, and nothing more: it compiles and is part of the reactor, but it is not +yet a spring-services feature module in the sense the other nine are. Open work, roughly in order: + +- **Auto-configuration.** `S3Config` is a plain `@Configuration` that no + `AutoConfiguration.imports` file names, so it is inert unless a consumer component-scans it. It + needs a `StorageAutoConfiguration` like every other module, and that raises the question the + README note already records: which implementation activates, and on what condition. `s3` vs + `file` cannot be decided by `@ConditionalOnClass` — both are always on the classpath. +- **Property namespace.** The `@Value` placeholders are `storage.s3.endpoint`, `.region`, + `.access-key`, `.secret-key` and `.bucket` — the origin application's namespace. The repo's is + `openelements.*`, and the values belong in a `@ConfigurationProperties` record rather than five + `@Value` parameters. +- **Tests.** The sources arrived without any. `FileObjectStore` alone justifies several: the + traversal guard on keys, the scratch-then-move visibility guarantee, ranged reads past the end, + and the incomplete-upload sweep. +- **`InMemoryObjectStore` is public API here, not a test fixture.** Its `failDeletes`, `failPuts`, + `lastGetOffset` and `lastGetLength` are public mutable fields — fine inside one application, not + as a published surface. Either give it a proper test-control API or move it to a test artifact. +- **The AWS SDK is a hard dependency.** A consumer that only wants `FileObjectStore` still pulls + `software.amazon.awssdk:s3`. `optional` would stop that, at the cost of making the S3 path require + an explicit declaration. +- **Javadoc still describes the origin application.** It refers to `OrphanSweep`, to "audio" as the + payload, and to Record Store as the target — none of which mean anything to a reader of this + library. The reasoning behind the prose is worth keeping; the nouns are not. + +**Context:** The module was created by moving the sources in as a deliberate first step — get +everything into the reactor compiling, decide the Spring-facing design afterwards. + ## Property toggles and consumer overridability for core security beans Per-feature `@ConditionalOnMissingBean` / `@ConditionalOnProperty` for all library beans, so diff --git a/pom.xml b/pom.xml index abaf1ad..bf3b257 100644 --- a/pom.xml +++ b/pom.xml @@ -47,6 +47,7 @@ spring-services-mcp spring-services-scim spring-services-tenant + spring-services-storage spring-services-all spring-services-bom @@ -60,6 +61,7 @@ 1.45.3 3.10.0 0.18.3 + 2.55.5 diff --git a/spring-services-all/pom.xml b/spring-services-all/pom.xml index c0f26c3..5899acd 100644 --- a/spring-services-all/pom.xml +++ b/spring-services-all/pom.xml @@ -57,6 +57,11 @@ spring-services-tenant ${project.version} + + com.open-elements + spring-services-storage + ${project.version} + diff --git a/spring-services-bom/pom.xml b/spring-services-bom/pom.xml index d5cd9b7..c50fad6 100644 --- a/spring-services-bom/pom.xml +++ b/spring-services-bom/pom.xml @@ -58,6 +58,11 @@ spring-services-tenant ${project.version} + + com.open-elements + spring-services-storage + ${project.version} + com.open-elements spring-services-all diff --git a/spring-services-storage/pom.xml b/spring-services-storage/pom.xml new file mode 100644 index 0000000..627ad5d --- /dev/null +++ b/spring-services-storage/pom.xml @@ -0,0 +1,37 @@ + + 4.0.0 + + + com.open-elements + spring-services + 1.5.0-SNAPSHOT + + + spring-services-storage + + Spring Services Storage + Optional object-store feature module for spring-services + https://github.com/OpenElementsLabs/spring-services + + + + com.open-elements + spring-services-core + ${project.version} + + + software.amazon.awssdk + s3 + ${awssdk.version} + + + + + org.springframework.boot + spring-boot-starter-test + test + + + + diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectNotFoundException.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectNotFoundException.java new file mode 100644 index 0000000..055c024 --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectNotFoundException.java @@ -0,0 +1,25 @@ +package com.openelements.spring.base.services.storage; + +/** Thrown when an object is read but no object exists under the key. */ +public class ObjectNotFoundException extends RuntimeException { + + /** + * Creates an exception naming the key that has no object. + * + * @param key the key that was read + */ + public ObjectNotFoundException(String key) { + super("No object found for key: " + key); + } + + /** + * Creates an exception naming the key that has no object, keeping the store's own report of the + * miss as the cause. + * + * @param key the key that was read + * @param cause the store's exception reporting the missing object + */ + public ObjectNotFoundException(String key, Throwable cause) { + super("No object found for key: " + key, cause); + } +} diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectStore.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectStore.java new file mode 100644 index 0000000..a1ce948 --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectStore.java @@ -0,0 +1,84 @@ +package com.openelements.spring.base.services.storage; + + +import java.io.InputStream; +import java.time.Instant; +import java.util.OptionalLong; +import java.util.stream.Stream; + +/** + * Abstraction over an S3-compatible object store. The single implementation targets AWS S3, Hetzner + * Object Storage and Record Store through an explicit endpoint and path-style addressing. + * + *

Writes stream without buffering the whole payload; reads return the store's response stream + * directly. Audio of any length must move through here without being materialised in the backend's + * heap or on its disk. + */ +public interface ObjectStore { + + /** + * Writes a stream under a key. The length need not be known in advance; the data is uploaded in + * bounded parts so an hour of audio never lands in memory. + * + * @param key the object key + * @param data the payload; fully consumed and closed by the store + * @param contentType the object's content type + */ + void put(String key, InputStream data, String contentType); + + /** + * Opens the object for reading. + * + * @param key the object key + * @return the object's content stream (caller closes it) + * @throws ObjectNotFoundException if no object exists under the key + */ + InputStream get(String key); + + /** + * Opens a byte range of the object. A range extending past the end returns the available bytes + * rather than failing. + * + * @param key the object key + * @param offset the first byte to return + * @param length the maximum number of bytes to return + * @return the requested slice (caller closes it) + * @throws ObjectNotFoundException if no object exists under the key + */ + InputStream get(String key, long offset, long length); + + /** + * The size of the object, or empty if it does not exist — distinguishable from an object of + * length zero. + * + * @param key the object key + * @return the size in bytes, or empty if there is no such object + */ + OptionalLong size(String key); + + /** + * Deletes the object. Idempotent: deleting a missing key succeeds. + * + * @param key the object key + * @return {@code true} if an object was removed, {@code false} if there was nothing to remove + */ + boolean delete(String key); + + /** + * Lists objects under a prefix. + * + * @param prefix the key prefix + * @return the matching objects + */ + Stream list(String prefix); + + /** + * Aborts incomplete multipart uploads initiated before the given moment. Such uploads are + * invisible to {@link #list(String)} — only reachable through the multipart-uploads listing — so + * the orphan sweep must reclaim them separately or they accumulate unreported. + * + * @param threshold abort uploads initiated before this instant + * @return the number of uploads aborted + */ + int abortIncompleteUploadsOlderThan(Instant threshold); +} diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectStoreException.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectStoreException.java new file mode 100644 index 0000000..bdd3103 --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/ObjectStoreException.java @@ -0,0 +1,17 @@ +package com.openelements.spring.base.services.storage; + +import org.jspecify.annotations.Nullable; + +/** Thrown when an object-store operation fails for a reason other than a missing object. */ +public class ObjectStoreException extends RuntimeException { + + /** + * Creates an exception describing a failed operation. + * + * @param message what the store could not do + * @param cause the underlying failure, or {@code null} where the store itself is the origin + */ + public ObjectStoreException(String message, @Nullable Throwable cause) { + super(message, cause); + } +} diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/StoredObject.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/StoredObject.java new file mode 100644 index 0000000..68eedf9 --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/StoredObject.java @@ -0,0 +1,28 @@ +package com.openelements.spring.base.services.storage; + +import java.time.Instant; +import java.util.Objects; + +/** + * A stored object as listed by {@link ObjectStore#list(String)}. + * + * @param key the object key + * @param size the object size in bytes + * @param lastModified when the object was last written + */ +public record StoredObject(String key, long size, Instant lastModified) { + + /** + * Validates the components. + * + * @throws NullPointerException if {@code key} or {@code lastModified} is {@code null} + * @throws IllegalArgumentException if {@code size} is negative + */ + public StoredObject { + Objects.requireNonNull(key, "key must not be null"); + if (size < 0) { + throw new IllegalArgumentException("size must not be negative, but was " + size); + } + Objects.requireNonNull(lastModified, "lastModified must not be null"); + } +} diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/file/FileObjectStore.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/file/FileObjectStore.java new file mode 100644 index 0000000..5a9487f --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/file/FileObjectStore.java @@ -0,0 +1,378 @@ +package com.openelements.spring.base.services.storage.file; + +import com.openelements.spring.base.services.storage.ObjectNotFoundException; +import com.openelements.spring.base.services.storage.ObjectStore; +import com.openelements.spring.base.services.storage.ObjectStoreException; +import com.openelements.spring.base.services.storage.StoredObject; +import java.io.IOException; +import java.io.InputStream; +import java.nio.channels.Channels; +import java.nio.channels.SeekableByteChannel; +import java.nio.file.AtomicMoveNotSupportedException; +import java.nio.file.DirectoryStream; +import java.nio.file.FileVisitResult; +import java.nio.file.Files; +import java.nio.file.NoSuchFileException; +import java.nio.file.Path; +import java.nio.file.SimpleFileVisitor; +import java.nio.file.StandardCopyOption; +import java.nio.file.StandardOpenOption; +import java.nio.file.attribute.BasicFileAttributes; +import java.time.Instant; +import java.util.ArrayList; +import java.util.List; +import java.util.Objects; +import java.util.OptionalLong; +import java.util.UUID; +import java.util.stream.Stream; +import org.jspecify.annotations.Nullable; + +/** + * An {@link ObjectStore} backed by a directory on the local file system. + * + *

It exists for deployments that have no object store to point at — a single machine with a disk, + * or a developer who does not want a container running. It keeps the interface's promises rather than + * approximating them, and the three that take work are these: + * + *

    + *
  • A key never becomes visible half-written. A {@link #put} streams into a + * scratch file under {@code uploads/} and then moves it into place, so a reader sees + * either the previous object or the complete new one. This is what makes a failed upload + * equivalent to S3's aborted multipart upload rather than a truncated object that later passes + * for audio. + *
  • Nothing is materialised. Writes copy stream-to-file, reads hand back a file + * stream, and a ranged read positions a {@link SeekableByteChannel} instead of skipping bytes — + * so seeking into an hour of audio costs one seek, not an hour of reads. + *
  • A key cannot escape the root. Keys are paths here, unlike in S3 where they + * are opaque strings, so {@code ../} would be a directory traversal. Every segment is validated + * and the resolved path is checked against the root. + *
+ * + *

The scratch directory also gives {@link #abortIncompleteUploadsOlderThan} something real to do: + * a crash mid-{@code put} leaves a file there, invisible to {@link #list} exactly as an incomplete + * multipart upload is invisible to S3's listing, and reclaimed by the same sweep. + * + *

The content type passed to {@link #put} is ignored. A file system has nowhere to + * keep it, and the interface offers no way to read it back — the application stores an object's + * content type in its own database. Inventing a sidecar file for a value nobody reads from here would + * be a second source of truth for it. + * + *

This class registers no bean of its own: an application chooses one {@code ObjectStore} + * implementation, and making two of them self-registering would decide that by classpath order. + */ +public class FileObjectStore implements ObjectStore { + + /** Objects live here, one file per key, directories mirroring the key's slashes. */ + private static final String OBJECTS_DIR = "objects"; + + /** In-flight writes live here until they are moved into {@link #OBJECTS_DIR}. */ + private static final String UPLOADS_DIR = "uploads"; + + private final Path objects; + private final Path uploads; + + /** + * Opens (and creates) a store rooted at {@code root}. + * + * @param root the directory the store owns; created if missing + * @throws ObjectStoreException if the directories cannot be created + */ + public FileObjectStore(final Path root) { + Objects.requireNonNull(root, "root must not be null"); + final Path base = root.toAbsolutePath().normalize(); + this.objects = base.resolve(OBJECTS_DIR); + this.uploads = base.resolve(UPLOADS_DIR); + try { + Files.createDirectories(objects); + Files.createDirectories(uploads); + } catch (final IOException e) { + throw new ObjectStoreException("Could not create the object store under " + root, e); + } + } + + @Override + public void put(final String key, final InputStream data, final String contentType) { + Objects.requireNonNull(key, "key must not be null"); + Objects.requireNonNull(data, "data must not be null"); + Objects.requireNonNull(contentType, "contentType must not be null"); + final Path target = pathOf(key); + final Path scratch = uploads.resolve(UUID.randomUUID() + ".part"); + try (data) { + // Stream to disk first. A failure here — including a source that throws mid-stream — must + // leave the key untouched, which is why nothing is written at the target path yet. + Files.copy(data, scratch, StandardCopyOption.REPLACE_EXISTING); + Files.createDirectories(target.getParent()); + moveIntoPlace(scratch, target); + } catch (final IOException e) { + deleteQuietly(scratch); + throw new ObjectStoreException("Could not write the object for key " + key, e); + } catch (final RuntimeException e) { + // A source stream that throws its own failure (a size limit, say) keeps that exception: + // the caller distinguishes those, and wrapping would flatten them into one store error. + deleteQuietly(scratch); + throw e; + } + } + + @Override + public InputStream get(final String key) { + Objects.requireNonNull(key, "key must not be null"); + try { + return Files.newInputStream(pathOf(key), StandardOpenOption.READ); + } catch (final NoSuchFileException e) { + throw new ObjectNotFoundException(key, e); + } catch (final IOException e) { + throw new ObjectStoreException("Could not read the object for key " + key, e); + } + } + + @Override + public InputStream get(final String key, final long offset, final long length) { + Objects.requireNonNull(key, "key must not be null"); + if (offset < 0) { + throw new IllegalArgumentException("offset must not be negative"); + } + if (length < 0) { + throw new IllegalArgumentException("length must not be negative"); + } + final Path path = pathOf(key); + if (length == 0) { + // Empty, but only for an object that exists: a zero-length read of a missing key is a + // wrong key, not "no bytes". + if (!Files.isRegularFile(path)) { + throw new ObjectNotFoundException(key); + } + return InputStream.nullInputStream(); + } + SeekableByteChannel channel = null; + try { + channel = Files.newByteChannel(path, StandardOpenOption.READ); + if (offset >= channel.size()) { + // A range starting past the end yields the available bytes — none — as the interface + // specifies, rather than failing the way an out-of-range seek would. + channel.close(); + return InputStream.nullInputStream(); + } + channel.position(offset); + return new BoundedInputStream(Channels.newInputStream(channel), length); + } catch (final NoSuchFileException e) { + throw new ObjectNotFoundException(key, e); + } catch (final IOException e) { + closeQuietly(channel); + throw new ObjectStoreException("Could not read the object for key " + key, e); + } + } + + @Override + public OptionalLong size(final String key) { + Objects.requireNonNull(key, "key must not be null"); + try { + return OptionalLong.of(Files.size(pathOf(key))); + } catch (final NoSuchFileException e) { + return OptionalLong.empty(); + } catch (final IOException e) { + throw new ObjectStoreException("Could not stat the object for key " + key, e); + } + } + + @Override + public boolean delete(final String key) { + Objects.requireNonNull(key, "key must not be null"); + final Path path = pathOf(key); + try { + final boolean removed = Files.deleteIfExists(path); + // The directories a key's slashes created are an artefact of this implementation — S3 has + // none — so an emptied one is swept away rather than left to accumulate per deleted object. + pruneEmptyParents(path.getParent()); + return removed; + } catch (final IOException e) { + throw new ObjectStoreException("Could not delete the object for key " + key, e); + } + } + + @Override + public Stream list(final String prefix) { + Objects.requireNonNull(prefix, "prefix must not be null"); + final List snapshot = new ArrayList<>(); + try { + Files.walkFileTree(objects, new SimpleFileVisitor() { + @Override + public FileVisitResult visitFile(final Path file, final BasicFileAttributes attrs) { + final String key = keyOf(file); + if (key.startsWith(prefix)) { + snapshot.add(new StoredObject(key, attrs.size(), + attrs.lastModifiedTime().toInstant())); + } + return FileVisitResult.CONTINUE; + } + + @Override + public FileVisitResult visitFileFailed(final Path file, final IOException e) { + // A file that vanished between the directory read and the stat was deleted by a + // concurrent caller. The orphan sweep walks this store while uploads run, and a + // listing that dies on that race would be worse than one that omits the file. + return FileVisitResult.CONTINUE; + } + }); + } catch (final IOException e) { + throw new ObjectStoreException("Could not list objects under " + prefix, e); + } + // A snapshot rather than a lazy walk, for the same reason: the caller iterates while other + // threads write, and the walk's own cursor would be the thing that breaks. + return snapshot.stream(); + } + + @Override + public int abortIncompleteUploadsOlderThan(final Instant threshold) { + Objects.requireNonNull(threshold, "threshold must not be null"); + int aborted = 0; + try (DirectoryStream stale = Files.newDirectoryStream(uploads)) { + for (final Path scratch : stale) { + try { + if (Files.getLastModifiedTime(scratch).toInstant().isBefore(threshold) + && Files.deleteIfExists(scratch)) { + aborted++; + } + } catch (final NoSuchFileException e) { + // Another sweep won the race; the upload is reclaimed either way. + } + } + } catch (final IOException e) { + throw new ObjectStoreException("Could not reclaim incomplete uploads", e); + } + return aborted; + } + + /** + * Resolves a key to a file below {@code objects/}, refusing anything that would leave it. + * + *

In S3 a key is an opaque string and {@code a/../b} is simply a key. Here it is a path, so the + * two readings differ and only one of them stays inside the store. + */ + private Path pathOf(final String key) { + if (key.isEmpty()) { + throw new IllegalArgumentException("key must not be empty"); + } + if (key.startsWith("/") || key.endsWith("/")) { + throw new IllegalArgumentException("key must not start or end with '/': " + key); + } + Path resolved = objects; + for (final String segment : key.split("/", -1)) { + if (segment.isEmpty() || ".".equals(segment) || "..".equals(segment) + || segment.indexOf('\\') >= 0) { + throw new IllegalArgumentException("key contains an unusable path segment: " + key); + } + resolved = resolved.resolve(segment); + } + final Path normalised = resolved.normalize(); + if (!normalised.startsWith(objects)) { + throw new IllegalArgumentException("key escapes the object store: " + key); + } + return normalised; + } + + /** The key a stored file represents — the inverse of {@link #pathOf}. */ + private String keyOf(final Path file) { + final StringBuilder key = new StringBuilder(); + for (final Path segment : objects.relativize(file)) { + if (!key.isEmpty()) { + key.append('/'); + } + key.append(segment); + } + return key.toString(); + } + + private static void moveIntoPlace(final Path scratch, final Path target) throws IOException { + try { + Files.move(scratch, target, StandardCopyOption.REPLACE_EXISTING, + StandardCopyOption.ATOMIC_MOVE); + } catch (final AtomicMoveNotSupportedException e) { + // Some file systems, and some container volume drivers, refuse the atomic flag. The plain + // move is still a rename within one tree; the guarantee is weaker, but a half-written + // target is still not reachable under the key. + Files.move(scratch, target, StandardCopyOption.REPLACE_EXISTING); + } + } + + /** Removes directories the key's slashes created, stopping at the first non-empty one. */ + private void pruneEmptyParents(final @Nullable Path from) throws IOException { + Path directory = from; + while (directory != null && directory.startsWith(objects) && !directory.equals(objects)) { + try (DirectoryStream entries = Files.newDirectoryStream(directory)) { + if (entries.iterator().hasNext()) { + return; + } + } catch (final NoSuchFileException e) { + return; + } + try { + Files.deleteIfExists(directory); + } catch (final IOException e) { + // A concurrent write repopulated it between the check and the delete. Leaving it is + // correct, and the next deletion below it will try again. + return; + } + directory = directory.getParent(); + } + } + + private static void deleteQuietly(final Path path) { + try { + Files.deleteIfExists(path); + } catch (final IOException ignored) { + // Already unreachable by key; a stale scratch file is what the upload sweep is for. + } + } + + private static void closeQuietly(final @Nullable SeekableByteChannel channel) { + if (channel != null) { + try { + channel.close(); + } catch (final IOException ignored) { + // Nothing left to report: the caller is already being handed a failure. + } + } + } + + /** Stops a file stream after {@code limit} bytes, so a ranged read cannot run past its range. */ + private static final class BoundedInputStream extends InputStream { + + private final InputStream delegate; + private long remaining; + + private BoundedInputStream(final InputStream delegate, final long limit) { + this.delegate = delegate; + this.remaining = limit; + } + + @Override + public int read() throws IOException { + if (remaining <= 0) { + return -1; + } + final int value = delegate.read(); + if (value >= 0) { + remaining--; + } + return value; + } + + @Override + public int read(final byte[] buffer, final int off, final int len) throws IOException { + if (remaining <= 0) { + return -1; + } + final int read = delegate.read(buffer, off, (int) Math.min(len, remaining)); + if (read > 0) { + remaining -= read; + } + return read; + } + + @Override + public void close() throws IOException { + delegate.close(); + } + } +} diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/file/package-info.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/file/package-info.java new file mode 100644 index 0000000..b5d734b --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/file/package-info.java @@ -0,0 +1,7 @@ +/** + * Null-marked package (JSpecify): every type is non-null unless annotated {@link org.jspecify.annotations.Nullable}. + */ +@NullMarked +package com.openelements.spring.base.services.storage.file; + +import org.jspecify.annotations.NullMarked; diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/memory/InMemoryObjectStore.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/memory/InMemoryObjectStore.java new file mode 100644 index 0000000..7307d12 --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/memory/InMemoryObjectStore.java @@ -0,0 +1,146 @@ +package com.openelements.spring.base.services.storage.memory; + +import com.openelements.spring.base.services.storage.ObjectStoreException; +import com.openelements.spring.base.services.storage.ObjectNotFoundException; +import com.openelements.spring.base.services.storage.ObjectStore; +import com.openelements.spring.base.services.storage.StoredObject; +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.time.Instant; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.OptionalLong; +import java.util.concurrent.ConcurrentHashMap; +import java.util.stream.Stream; +import org.jspecify.annotations.Nullable; + +/** + * In-memory {@link ObjectStore} for service-level tests (idempotency, deletion). This is a fake of + * the interface, not a mock of the AWS SDK — the S3 implementation itself is tested against a real + * Record Store. {@link #failDeletes} lets a test force a delete failure. + * + *

The backing map is concurrent and {@link #list} snapshots it, because the startup + * {@code OrphanSweep} iterates the store on a background thread while tests write to it. + */ +public class InMemoryObjectStore implements ObjectStore { + + private final Map objects = new ConcurrentHashMap<>(); + /** + * When each object was written. + * + *

{@link #list} used to report {@code Instant.now()} as every object's last-modified time, + * which made an object look freshly written no matter how long it had been there — and so made + * {@code OrphanSweep}'s age filter impossible to exercise against this fake. Recording the write + * time makes the fake able to express age, which is the only thing that filter reacts to. + */ + private final Map writtenAt = new ConcurrentHashMap<>(); + /** When set, every {@link #delete(String)} fails the way the S3 store reports a refused delete. */ + public volatile boolean failDeletes = false; + /** When set, every {@link #put(String, InputStream, String)} fails after the source is drained. */ + public volatile boolean failPuts = false; + /** The offset of the most recent ranged {@link #get(String, long, long)} call, for assertions. */ + public volatile long lastGetOffset = -1; + /** The length of the most recent ranged {@link #get(String, long, long)} call, for assertions. */ + public volatile long lastGetLength = -1; + + /** Creates an empty store. */ + public InMemoryObjectStore() { + } + + @Override + public void put(String key, InputStream data, String contentType) { + try { + // Drain the stream first so a size-limited source still throws its own overflow. + byte[] bytes = data.readAllBytes(); + if (failPuts) { + // Mirror the real S3 store, which reports a failed part as ObjectStoreException. + throw new ObjectStoreException("put rejected by test", null); + } + objects.put(key, bytes); + writtenAt.put(key, Instant.now()); + } catch (IOException e) { + throw new IllegalStateException(e); + } + } + + @Override + public InputStream get(String key) { + byte[] bytes = objects.get(key); + if (bytes == null) { + throw new ObjectNotFoundException(key); + } + return new ByteArrayInputStream(bytes); + } + + @Override + public InputStream get(String key, long offset, long length) { + lastGetOffset = offset; + lastGetLength = length; + byte[] bytes = objects.get(key); + if (bytes == null) { + throw new ObjectNotFoundException(key); + } + int from = (int) Math.min(offset, bytes.length); + int to = (int) Math.min(offset + length, bytes.length); + byte[] slice = new byte[to - from]; + System.arraycopy(bytes, from, slice, 0, to - from); + return new ByteArrayInputStream(slice); + } + + @Override + public OptionalLong size(String key) { + byte[] bytes = objects.get(key); + return bytes == null ? OptionalLong.empty() : OptionalLong.of(bytes.length); + } + + @Override + public boolean delete(String key) { + if (failDeletes) { + // Mirror the real S3 store, which surfaces delete failures as ObjectStoreException. + throw new ObjectStoreException("delete rejected by test", null); + } + writtenAt.remove(key); + return objects.remove(key) != null; + } + + @Override + public Stream list(String prefix) { + // Snapshot so a concurrent write (e.g. a test upload) cannot disturb the iteration. + List snapshot = new ArrayList<>(); + objects.forEach((key, value) -> { + if (key.startsWith(prefix)) { + snapshot.add(new StoredObject(key, value.length, + writtenAt.getOrDefault(key, Instant.now()))); + } + }); + return snapshot.stream(); + } + + @Override + public int abortIncompleteUploadsOlderThan(java.time.Instant threshold) { + // No multipart uploads in the in-memory fake; nothing to reclaim. + return 0; + } + + /** + * Whether an object is stored under a key. + * + * @param key the object key + * @return {@code true} if the key has an object + */ + public boolean contains(String key) { + return objects.containsKey(key); + } + + /** + * The stored bytes for a key. + * + * @param key the object key + * @return the stored bytes, or {@code null} if the key has no object + */ + public byte @Nullable [] contentOf(String key) { + return objects.get(key); + } +} diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/memory/package-info.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/memory/package-info.java new file mode 100644 index 0000000..b8ce6f5 --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/memory/package-info.java @@ -0,0 +1,7 @@ +/** + * Null-marked package (JSpecify): every type is non-null unless annotated {@link org.jspecify.annotations.Nullable}. + */ +@NullMarked +package com.openelements.spring.base.services.storage.memory; + +import org.jspecify.annotations.NullMarked; diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/package-info.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/package-info.java new file mode 100644 index 0000000..2458dad --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/package-info.java @@ -0,0 +1,7 @@ +/** + * Null-marked package (JSpecify): every type is non-null unless annotated {@link org.jspecify.annotations.Nullable}. + */ +@NullMarked +package com.openelements.spring.base.services.storage; + +import org.jspecify.annotations.NullMarked; diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/S3Config.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/S3Config.java new file mode 100644 index 0000000..4c8e978 --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/S3Config.java @@ -0,0 +1,100 @@ +package com.openelements.spring.base.services.storage.s3; + +import java.net.URI; + +import com.openelements.spring.base.services.storage.ObjectStore; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import software.amazon.awssdk.auth.credentials.AwsBasicCredentials; +import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider; +import software.amazon.awssdk.regions.Region; +import software.amazon.awssdk.services.s3.S3Client; +import software.amazon.awssdk.services.s3.S3ClientBuilder; +import software.amazon.awssdk.services.s3.S3Configuration; + +/** + * Wires the AWS SDK v2 {@link S3Client} and the {@link ObjectStore}. + * + *

The bucket and credentials are read through {@code @Value}, which fails start-up when the + * underlying {@code S3_BUCKET} / {@code S3_ACCESS_KEY} / {@code S3_SECRET_KEY} placeholder cannot be + * resolved (unlike {@code @ConfigurationProperties}, which would silently keep the unresolved + * string). A backend that starts without a store would accept uploads it cannot keep. + */ +@Configuration +public class S3Config { + + /** Creates the configuration; Spring instantiates it. */ + public S3Config() { + } + + /** + * The SDK client every request goes through, closed when the context shuts down. + * + * @param endpoint the store's endpoint + * @param region the region to sign requests for + * @param accessKey the access key + * @param secretKey the secret key + * @return the client + */ + @Bean(destroyMethod = "close") + public S3Client s3Client( + @Value("${storage.s3.endpoint}") String endpoint, + @Value("${storage.s3.region}") String region, + @Value("${storage.s3.access-key}") String accessKey, + @Value("${storage.s3.secret-key}") String secretKey) { + return withoutChunkedEncoding(S3Client.builder() + .endpointOverride(URI.create(endpoint)) + .region(Region.of(region)) + .credentialsProvider(StaticCredentialsProvider.create( + AwsBasicCredentials.create(accessKey, secretKey)))) + // Several S3-compatible providers do not support virtual-host addressing. + .forcePathStyle(true) + .build(); + } + + /** + * Makes the Java SDK send each request body in one piece instead of {@code aws-chunked}. + * + *

By default the SDK frames a body as {@code aws-chunked} with a trailing CRC32 + * ({@code x-amz-trailer: x-amz-checksum-crc32}). Record Store does not implement trailing + * checksums and answers {@code 501 NotImplemented} (since 0.1.2; earlier versions returned a + * generic {@code 400 InvalidRequest}), so chunked encoding is switched off and the body goes + * out whole, with its checksum in an ordinary header the store verifies. + * + *

This costs one extra pass over each part and no memory: {@code S3ObjectStore} already + * hands the SDK a {@link java.io.ByteArrayInputStream} over a bounded {@code PART_SIZE_BYTES} + * buffer, which is both resettable and already resident, so nothing is buffered that was not + * already there. + * + *

The setting is not specific to Record Store in the sense of breaking anything else: + * {@code aws-chunked} is an optimisation, so AWS and other S3-compatible providers accept a + * body sent in one piece just as well. + * + *

A second departure from the SDK's defaults used to live here: payload signing, because the + * SDK's {@code x-amz-content-sha256: UNSIGNED-PAYLOAD} was refused. Record Store 0.1.2 accepts + * it (record-store#74), so that one is gone and the SDK's default signing applies. + * + * @param builder the builder to configure + * @return the same builder, for chaining + */ + // Public, not package-private: the sweep's integration test lives in the orphan package and + // builds its client through this very helper on purpose, so a test can never pass against a + // client configured more leniently than the one the application bean gets. + public static S3ClientBuilder withoutChunkedEncoding(S3ClientBuilder builder) { + return builder.serviceConfiguration( + S3Configuration.builder().chunkedEncodingEnabled(false).build()); + } + + /** + * The store the application talks to. + * + * @param s3Client the configured client + * @param bucket the bucket every key lives in + * @return the store + */ + @Bean + public ObjectStore objectStore(S3Client s3Client, @Value("${storage.s3.bucket}") String bucket) { + return new S3ObjectStore(s3Client, bucket); + } +} diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/S3ObjectStore.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/S3ObjectStore.java new file mode 100644 index 0000000..6b6cb77 --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/S3ObjectStore.java @@ -0,0 +1,207 @@ +package com.openelements.spring.base.services.storage.s3; + +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.time.Instant; +import java.util.ArrayList; +import java.util.List; +import java.util.Objects; +import java.util.OptionalLong; +import java.util.stream.Stream; + +import com.openelements.spring.base.services.storage.ObjectNotFoundException; +import com.openelements.spring.base.services.storage.ObjectStore; +import com.openelements.spring.base.services.storage.ObjectStoreException; +import com.openelements.spring.base.services.storage.StoredObject; +import software.amazon.awssdk.core.sync.RequestBody; +import software.amazon.awssdk.services.s3.S3Client; +import software.amazon.awssdk.services.s3.model.CompletedMultipartUpload; +import software.amazon.awssdk.services.s3.model.CompletedPart; +import software.amazon.awssdk.services.s3.model.CreateMultipartUploadResponse; +import software.amazon.awssdk.services.s3.model.NoSuchKeyException; +import software.amazon.awssdk.services.s3.model.UploadPartResponse; + +/** + * {@link ObjectStore} on the AWS SDK v2 synchronous client. + * + *

Uploads stream in bounded {@value #PART_SIZE_BYTES}-byte parts: a small payload becomes a single + * {@code PutObject} once its (short) length is known after reading, a large one becomes a multipart + * upload whose parts are uploaded and discarded one at a time, so heap use stays bounded no matter + * how long the audio is. + */ +public class S3ObjectStore implements ObjectStore { + + /** Part size for multipart uploads (8 MiB, above S3's 5 MiB minimum for non-final parts). */ + static final int PART_SIZE_BYTES = 8 * 1024 * 1024; + + private final S3Client s3; + private final String bucket; + + /** + * Creates a store over one bucket. + * + * @param s3 the client to issue requests through + * @param bucket the bucket every key lives in + */ + public S3ObjectStore(final S3Client s3, final String bucket) { + this.s3 = Objects.requireNonNull(s3, "s3 must not be null"); + this.bucket = Objects.requireNonNull(bucket, "bucket must not be null"); + } + + @Override + public void put(final String key, final InputStream data, final String contentType) { + Objects.requireNonNull(key, "key must not be null"); + Objects.requireNonNull(data, "data must not be null"); + Objects.requireNonNull(contentType, "contentType must not be null"); + try (data) { + final byte[] firstPart = new byte[PART_SIZE_BYTES]; + final int firstLen = readFully(data, firstPart); + + if (firstLen < PART_SIZE_BYTES) { + // The whole stream (including an empty one) fits in one buffer: a single PutObject + // with the now-known length. No multipart machinery for small objects. + s3.putObject(b -> b.bucket(bucket).key(key).contentType(contentType) + .contentLength((long) firstLen), + RequestBody.fromInputStream(new ByteArrayInputStream(firstPart, 0, firstLen), + firstLen)); + return; + } + putMultipart(key, contentType, data, firstPart); + } catch (final IOException e) { + throw new ObjectStoreException("Could not read the stream for key " + key, e); + } + } + + private void putMultipart(final String key, final String contentType, final InputStream data, final byte[] firstPart) { + Objects.requireNonNull(key, "key must not be null"); + Objects.requireNonNull(data, "data must not be null"); + Objects.requireNonNull(contentType, "contentType must not be null"); + Objects.requireNonNull(firstPart, "firstPart must not be null"); + final CreateMultipartUploadResponse created = s3.createMultipartUpload( + b -> b.bucket(bucket).key(key).contentType(contentType)); + final String uploadId = created.uploadId(); + final List parts = new ArrayList<>(); + try { + int partNumber = 1; + byte[] buffer = firstPart; + int len = PART_SIZE_BYTES; + while (len == PART_SIZE_BYTES) { + parts.add(uploadPart(key, uploadId, partNumber, buffer, len)); + partNumber++; + buffer = new byte[PART_SIZE_BYTES]; + len = readFully(data, buffer); + if (len == 0) { + break; + } + if (len < PART_SIZE_BYTES) { + parts.add(uploadPart(key, uploadId, partNumber, buffer, len)); + break; + } + } + final List finalParts = parts; + s3.completeMultipartUpload(b -> b.bucket(bucket).key(key).uploadId(uploadId) + .multipartUpload(CompletedMultipartUpload.builder().parts(finalParts).build())); + } catch (final RuntimeException | IOException e) { + // Includes a failure reading the source stream mid-upload: abort so no incomplete + // multipart upload is left accruing storage (the sweep lists objects, not uploads). + s3.abortMultipartUpload(b -> b.bucket(bucket).key(key).uploadId(uploadId)); + throw new ObjectStoreException("Multipart upload failed for key " + key, e); + } + } + + private CompletedPart uploadPart(final String key, final String uploadId, final int partNumber, final byte[] buffer, + final int length) { + Objects.requireNonNull(key, "key must not be null"); + Objects.requireNonNull(uploadId, "uploadId must not be null"); + Objects.requireNonNull(buffer, "buffer must not be null"); + Objects.requireNonNull(buffer, "buffer must not be null"); + final UploadPartResponse response = s3.uploadPart( + b -> b.bucket(bucket).key(key).uploadId(uploadId).partNumber(partNumber), + RequestBody.fromInputStream(new ByteArrayInputStream(buffer, 0, length), length)); + return CompletedPart.builder().partNumber(partNumber).eTag(response.eTag()).build(); + } + + @Override + public InputStream get(final String key) { + Objects.requireNonNull(key, "key must not be null"); + try { + return s3.getObject(b -> b.bucket(bucket).key(key)); + } catch (final NoSuchKeyException e) { + throw new ObjectNotFoundException(key, e); + } + } + + @Override + public InputStream get(final String key, final long offset, final long length) { + Objects.requireNonNull(key, "key must not be null"); + if (offset < 0) { + throw new IllegalArgumentException("offset must not be negative"); + } + if(length < 0) { + throw new IllegalArgumentException("length must not be negative"); + } + if (length == 0) { + return new ByteArrayInputStream(new byte[0]); + } + long last = offset + length - 1; + try { + return s3.getObject(b -> b.bucket(bucket).key(key).range("bytes=" + offset + "-" + last)); + } catch (final NoSuchKeyException e) { + throw new ObjectNotFoundException(key, e); + } + } + + @Override + public OptionalLong size(final String key) { + try { + return OptionalLong.of(s3.headObject(b -> b.bucket(bucket).key(key)).contentLength()); + } catch (final NoSuchKeyException e) { + return OptionalLong.empty(); + } + } + + @Override + public boolean delete(final String key) { + Objects.requireNonNull(key, "key must not be null"); + final boolean existed = size(key).isPresent(); + s3.deleteObject(b -> b.bucket(bucket).key(key)); + return existed; + } + + @Override + public Stream list(final String prefix) { + Objects.requireNonNull(prefix, "prefix must not be null"); + return s3.listObjectsV2Paginator(b -> b.bucket(bucket).prefix(prefix)).contents().stream() + .map(o -> new StoredObject(o.key(), o.size(), o.lastModified())); + } + + @Override + public int abortIncompleteUploadsOlderThan(final Instant threshold) { + Objects.requireNonNull(threshold, "threshold must not be null"); + int aborted = 0; + for (var upload : s3.listMultipartUploadsPaginator(b -> b.bucket(bucket)).uploads()) { + if (upload.initiated().isBefore(threshold)) { + s3.abortMultipartUpload(b -> b.bucket(bucket).key(upload.key()) + .uploadId(upload.uploadId())); + aborted++; + } + } + return aborted; + } + + /** Reads until the buffer is full or the stream ends; returns the number of bytes read. */ + private static int readFully(final InputStream in, final byte[] buffer) throws IOException { + Objects.requireNonNull(in, "in must not be null"); + Objects.requireNonNull(buffer, "buffer must not be null"); + int total = 0; + while (total < buffer.length) { + int read = in.read(buffer, total, buffer.length - total); + if (read < 0) { + break; + } + total += read; + } + return total; + } +} diff --git a/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/package-info.java b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/package-info.java new file mode 100644 index 0000000..e491465 --- /dev/null +++ b/spring-services-storage/src/main/java/com/openelements/spring/base/services/storage/s3/package-info.java @@ -0,0 +1,7 @@ +/** + * Null-marked package (JSpecify): every type is non-null unless annotated {@link org.jspecify.annotations.Nullable}. + */ +@NullMarked +package com.openelements.spring.base.services.storage.s3; + +import org.jspecify.annotations.NullMarked;