From 00d2af2fb05d38660875d8307394233dfa9ed12a Mon Sep 17 00:00:00 2001 From: LunarSkyOSS Date: Sat, 10 Oct 2026 08:48:15 +0200 Subject: [PATCH] Add cluster multipart uploads and listings --- README.md | 8 +- TESTS.md | 4 +- scripts/test-cluster-http.py | 30 +- scripts/test-cluster.sh | 14 + src/cloud/lunarsky/store/ClusterStore.java | 407 +++++++++++++++++- src/cloud/lunarsky/store/Main.java | 89 +++- .../lunarsky/store/MultipartStorage.java | 5 + src/cloud/lunarsky/store/MultipartStore.java | 42 ++ src/cloud/lunarsky/store/SchemaMigrator.java | 8 +- .../lunarsky/store/UnavailableMultipart.java | 17 - .../store/ClusterIntegrationTest.java | 62 +++ test/cloud/lunarsky/store/HttpTest.java | 12 + 12 files changed, 660 insertions(+), 38 deletions(-) delete mode 100644 src/cloud/lunarsky/store/UnavailableMultipart.java diff --git a/README.md b/README.md index 4884f2f..6ea1803 100644 --- a/README.md +++ b/README.md @@ -43,10 +43,10 @@ New objects retain their content type and key. Objects written by the earlier si - ✅ `PutObject`, `GetObject`, `HeadObject`, and `DeleteObject` in both modes - ✅ Single-range GET and `ListObjectsV2` in both modes - ✅ SHA-256 payload verification and `x-amz-checksum-sha256` in both modes -- ✅ `CreateMultipartUpload`, `UploadPart`, `CompleteMultipartUpload`, and `AbortMultipartUpload` in single-node mode -- ⬜ Multipart uploads in cluster mode +- ✅ `CreateMultipartUpload`, `UploadPart`, `CompleteMultipartUpload`, and `AbortMultipartUpload` in both modes +- ✅ `ListParts` and `ListMultipartUploads` in both modes - ⬜ Presigned URLs and streaming Signature V4 uploads -- ⬜ `CopyObject`, `ListParts`, and `ListMultipartUploads` +- ⬜ `CopyObject` - ⬜ `Content-MD5` and checksum algorithms other than SHA-256 - ⬜ Bucket creation and listing, object versioning, ACLs, tags, and user metadata @@ -92,6 +92,8 @@ docker compose --env-file /path/to/cluster.env -f compose.cluster.yaml run --rm The cluster S3 endpoint binds to `127.0.0.1:9001`; storage nodes and PostgreSQL have no published ports. The separate repair container holds the repair credential and restores missing or corrupt replicas. +Multipart parts are stored on cluster nodes and indexed in PostgreSQL. Incomplete uploads count toward the logical capacity limit; abort them to release that capacity. Repair includes staged parts. The gateway upgrades the metadata schema when it starts, so back up the database before upgrading an existing cluster. + Node UUIDs persist on their volumes, and replica manifests use those UUIDs so reordering configured URLs cannot move an existing replica. Each node also has an operator-assigned physical host UUID. New writes require acknowledgements from two different host UUIDs. The optional `CLUSTER_TEST_NODE_DOMAINS=true` override counts containers instead, solely for local process tests; all containers in this Compose file share one physical host. ## Migrating a local cluster diff --git a/TESTS.md b/TESTS.md index 2aa1f66..5a463be 100644 --- a/TESTS.md +++ b/TESTS.md @@ -18,7 +18,7 @@ The script compiles the source and test programs into `out/classes`, then runs: | --- | --- | | `StoreTest` | Signature V4 test vector and tampering, local writes and reads, quotas, restart persistence, multipart recovery, legacy reads, locking, and corruption rejection. | | `ConcurrencyTest` | Atomic local overwrites and consistent reads, listings, and deletes during concurrent access. | -| `HttpTest` | Signed HTTP requests, object operations, ranges, listing, and single-node multipart uploads. | +| `HttpTest` | Signed HTTP requests, object operations, ranges, listing, multipart uploads, and multipart listings in single-node mode. | | `ClusterNodeTest` | Node identity and locking, authenticated segment transfers, checksum rejection, repair authorization, and restart cleanup. | | `CliTest` | Version, status, verification, and a nonzero result for corrupt data. | @@ -47,7 +47,7 @@ COMPOSE_PROJECT_NAME=objectstore-tests docker compose --env-file /tmp/objectstor If port 9001 is occupied, set `CLUSTER_HOST_PORT` to the same free port in both the environment file and the shell before running the script. The script reads that port from the shell; Compose reads it from the file. -The Docker suite checks signed S3 operations, multi-segment objects, concurrent overwrites, reads and writes with a node stopped, refusal to write without a storage quorum, restart recovery, corrupt-replica repair, metadata unavailability, and placement on a newly joined node. It also checks that containers labeled as one physical host cannot satisfy the normal host quorum. Its local-only override permits the remaining phases to use containers as separate test domains. +The Docker suite checks signed S3 operations, multi-segment objects, concurrent overwrites, multipart staging and listings, completion after a gateway restart and node loss, reads and writes with a node stopped, refusal to write without a storage quorum, restart recovery, corrupt-replica repair including staged parts, metadata unavailability, and placement on a newly joined node. It also checks that containers labeled as one physical host cannot satisfy the normal host quorum. Its local-only override permits the remaining phases to use containers as separate test domains. `ClusterMigrationTest` is a separate legacy-format fixture and is **not** run by either test script. Do not run its `create` phase against a populated metadata database. The migration procedure is in the [README](README.md#migrating-a-local-cluster). diff --git a/scripts/test-cluster-http.py b/scripts/test-cluster-http.py index 3e13e5c..c7c5bdf 100644 --- a/scripts/test-cluster-http.py +++ b/scripts/test-cluster-http.py @@ -6,6 +6,7 @@ import http.client import pathlib import sys import urllib.parse +import xml.etree.ElementTree as ET values = dict(line.strip().split("=", 1) for line in pathlib.Path(sys.argv[1]).read_text().splitlines() @@ -43,7 +44,7 @@ def request(method, path, body=b"", extra=None): f"SignedHeaders={signed_names},Signature={signature}") connection = http.client.HTTPConnection("127.0.0.1", port, timeout=30) try: - connection.request(method, path, body=body if method == "PUT" else None, headers=headers) + connection.request(method, path, body=body if method in ("PUT", "POST") else None, headers=headers) response = connection.getresponse() return response.status, response.read(), response.headers finally: @@ -81,4 +82,29 @@ status, _, _ = request("DELETE", key) assert status == 204, status status, _, _ = request("GET", key) assert status == 404, status -print("Cluster HTTP tests passed: signed PUT, GET, range, HEAD, LIST, DELETE") + +multipart_key = f"/{bucket}/cluster-test/http-multipart.txt" +status, content, _ = request("POST", multipart_key + "?uploads") +assert status == 200, (status, content) +upload_id = ET.fromstring(content).findtext("UploadId") +assert upload_id, content +part_etags = [] +for number, part in enumerate((b"hello ", b"world"), start=1): + status, _, headers = request("PUT", multipart_key + f"?partNumber={number}&uploadId={upload_id}", part) + assert status == 200, status + part_etags.append(headers["etag"]) +status, content, _ = request("GET", multipart_key + f"?uploadId={upload_id}&max-parts=1") +assert status == 200 and b"true" in content, (status, content) +status, content, _ = request("GET", f"/{bucket}?uploads&prefix=cluster-test%2Fhttp-multipart") +assert status == 200 and upload_id.encode() in content, (status, content) +completion = "" + "".join( + f"{number}{etag}" + for number, etag in enumerate(part_etags, start=1)) + "" +status, content, _ = request("POST", multipart_key + f"?uploadId={upload_id}", completion.encode(), + {"content-type": "application/xml"}) +assert status == 200 and b"" in content, (status, content) +status, content, _ = request("GET", multipart_key) +assert status == 200 and content == b"hello world", (status, content) +status, content, _ = request("GET", f"/{bucket}?uploads&prefix=cluster-test%2Fhttp-multipart") +assert status == 200 and upload_id.encode() not in content, (status, content) +print("Cluster HTTP tests passed: signed object and multipart operations") diff --git a/scripts/test-cluster.sh b/scripts/test-cluster.sh index 2f077b2..f67dfa7 100644 --- a/scripts/test-cluster.sh +++ b/scripts/test-cluster.sh @@ -20,10 +20,24 @@ wait_ready() { done } run_phase basic +run_phase multipart-stage +compose restart gateway +wait_ready +part_segment=$(compose exec -T metadata psql -U objectstore -d objectstore -At -c \ + "SELECT segment_id FROM cluster_upload_segments LIMIT 1") +printf '%s\n' "$part_segment" | grep -Eq '^[0-9a-f-]{36}$' +part_shard=$(printf '%s' "$part_segment" | cut -c1-2) +compose exec -T node-a sh -c 'printf corrupted > "/data/segments/$1/$2"' _ "$part_shard" "$part_segment" +compose run --rm -T repair +part_expected=$(compose exec -T metadata psql -U objectstore -d objectstore -At -c \ + "SELECT encode(sha256,'hex') FROM cluster_upload_segments WHERE segment_id='$part_segment'") +part_actual=$(compose exec -T node-a sha256sum "/data/segments/$part_shard/$part_segment" | cut -d' ' -f1) +[ "$part_expected" = "$part_actual" ] run_phase same-host run_phase concurrent compose stop node-a run_phase degraded +run_phase multipart-complete compose stop node-b run_phase quorum-lost compose start node-a node-b diff --git a/src/cloud/lunarsky/store/ClusterStore.java b/src/cloud/lunarsky/store/ClusterStore.java index c6ebafd..7a6748e 100644 --- a/src/cloud/lunarsky/store/ClusterStore.java +++ b/src/cloud/lunarsky/store/ClusterStore.java @@ -18,15 +18,18 @@ import java.sql.SQLException; import java.time.Instant; import java.util.ArrayList; import java.util.Base64; +import java.util.Comparator; import java.util.HexFormat; import java.util.HashSet; import java.util.List; import java.util.Set; import java.util.UUID; -final class ClusterStore implements ObjectStorage { +final class ClusterStore implements ObjectStorage, MultipartStorage { private record Segment(UUID id, int length, byte[] hash, List replicas) {} - private record RepairTarget(UUID generation, int ordinal, long version, Segment segment) {} + private record RepairTarget(UUID id, int part, int ordinal, long version, Segment segment) {} + private record Upload(String contentType) {} + private record StoredPart(long length, String etag, List segments) {} record RepairReport(int scanned, int restored, int underReplicated, int unrecoverable) {} private final String jdbcUrl, user, password, configuredBucket; private final NodeClient nodes; @@ -116,7 +119,8 @@ final class ClusterStore implements ObjectStorage { query.setString(1, bucket); try (ResultSet result = query.executeQuery()) { if (!result.next()) throw new SQLException("Bucket quota row is missing"); - if (result.getLong(1) - Math.max(0, previous) > maxTotal - length) + if (maxTotal - (result.getLong(1) - Math.max(0, previous)) - + stagedBytes(connection, bucket) < length) throw new StoreException(507, "InsufficientStorage", "Store capacity limit reached"); } } @@ -168,7 +172,8 @@ final class ClusterStore implements ObjectStorage { long previous = currentLength(connection, bucket, key); if (createOnly && previous >= 0) throw new StoreException(412, "PreconditionFailed", "Object already exists"); - if (used - Math.max(0, previous) > maxTotal - length) + if (maxTotal - (used - Math.max(0, previous)) - + stagedBytes(connection, bucket) < length) throw new StoreException(507, "InsufficientStorage", "Store capacity limit reached"); try (PreparedStatement insert = connection.prepareStatement( "INSERT INTO cluster_segments (generation, ordinal, segment_id, length, sha256, replicas, replica_ids) VALUES (?, ?, ?, ?, ?, 'v2', ?)")) { @@ -208,6 +213,276 @@ final class ClusterStore implements ObjectStorage { } catch (SQLException error) { throw databaseError(error); } } + @Override public String create(String bucket, String key, String contentType) throws IOException { + if (!configuredBucket.equals(bucket)) throw new StoreException(404, "NoSuchBucket", "Bucket not found"); + if (contentType == null || contentType.getBytes(java.nio.charset.StandardCharsets.UTF_8).length > 255) + throw new StoreException(400, "InvalidArgument", "Invalid Content-Type"); + UUID id = UUID.randomUUID(); + try (Connection connection = connect()) { + connection.setAutoCommit(false); + try { + lockUsage(connection, bucket); + try (PreparedStatement count = connection.prepareStatement("SELECT count(*) FROM cluster_uploads WHERE bucket=?")) { + count.setString(1, bucket); + try (ResultSet result = count.executeQuery()) { + result.next(); + if (result.getLong(1) >= 32) + throw new StoreException(503, "SlowDown", "Too many active uploads"); + } + } + try (PreparedStatement insert = connection.prepareStatement( + "INSERT INTO cluster_uploads VALUES (?, ?, ?, ?, ?)")) { + insert.setObject(1, id); + insert.setString(2, bucket); + insert.setString(3, key); + insert.setString(4, contentType); + insert.setLong(5, Instant.now().toEpochMilli()); + insert.executeUpdate(); + } + connection.commit(); + } catch (SQLException | RuntimeException error) { + connection.rollback(); + if (error instanceof SQLException sql) throw databaseError(sql); + throw error; + } + } catch (SQLException error) { throw databaseError(error); } + return id.toString(); + } + + @Override public String putPart(String id, String bucket, String key, int number, InputStream input, + long length, String expectedHash, String checksum) throws IOException { + if (number < 1 || number > 10000) throw new StoreException(400, "InvalidArgument", "Invalid part number"); + validatePut(bucket, length, "application/octet-stream"); + UUID uploadId = uploadId(id); + try (Connection connection = connect()) { upload(connection, uploadId, bucket, key, false); } + catch (SQLException error) { throw databaseError(error); } + Path staged = Files.createTempFile("objectstore-part-", ".pending"); + MessageDigest md5 = digest("MD5"); + List segments; + try { + stageInput(staged, input, length, expectedHash, checksum, md5); + segments = uploadSegments(staged, length); + } finally { Files.deleteIfExists(staged); } + String etag = HexFormat.of().formatHex(md5.digest()); + try (Connection connection = connect()) { + connection.setAutoCommit(false); + try { + long used = lockUsage(connection, bucket); + upload(connection, uploadId, bucket, key, true); + long previous = partLength(connection, uploadId, number); + if (maxTotal - used - (stagedBytes(connection, bucket) - previous) < length) + throw new StoreException(507, "InsufficientStorage", "Multipart staging limit reached"); + try (PreparedStatement insert = connection.prepareStatement( + "INSERT INTO cluster_upload_parts VALUES (?, ?, ?, ?, ?) ON CONFLICT (upload_id, part_number) DO UPDATE SET length=EXCLUDED.length, etag=EXCLUDED.etag, modified=EXCLUDED.modified")) { + insert.setObject(1, uploadId); + insert.setInt(2, number); + insert.setLong(3, length); + insert.setString(4, etag); + insert.setLong(5, Instant.now().toEpochMilli()); + insert.executeUpdate(); + } + try (PreparedStatement delete = connection.prepareStatement( + "DELETE FROM cluster_upload_segments WHERE upload_id=? AND part_number=?")) { + delete.setObject(1, uploadId); + delete.setInt(2, number); + delete.executeUpdate(); + } + try (PreparedStatement insert = connection.prepareStatement( + "INSERT INTO cluster_upload_segments VALUES (?, ?, ?, ?, ?, ?, ?)")) { + for (int ordinal = 0; ordinal < segments.size(); ordinal++) { + Segment segment = segments.get(ordinal); + insert.setObject(1, uploadId); + insert.setInt(2, number); + insert.setInt(3, ordinal); + insert.setObject(4, segment.id()); + insert.setInt(5, segment.length()); + insert.setBytes(6, segment.hash()); + insert.setArray(7, connection.createArrayOf("uuid", segment.replicas().toArray())); + insert.addBatch(); + } + insert.executeBatch(); + } + connection.commit(); + } catch (SQLException | RuntimeException error) { + connection.rollback(); + if (error instanceof SQLException sql) throw databaseError(sql); + throw error; + } + } catch (SQLException error) { throw databaseError(error); } + return etag; + } + + @Override public Metadata complete(String id, String bucket, String key, List parts) throws IOException { + if (parts.isEmpty() || parts.size() > 10000) + throw new StoreException(400, "InvalidPart", "No valid parts supplied"); + UUID uploadId = uploadId(id); + try (Connection connection = connect()) { + connection.setAutoCommit(false); + try { + long used = lockUsage(connection, bucket); + Upload upload = upload(connection, uploadId, bucket, key, true); + long staged = stagedBytes(connection, bucket); + long uploadBytes = uploadLength(connection, uploadId); + MessageDigest fullHash = digest("SHA-256"); + MessageDigest etagHash = digest("MD5"); + List selected = new ArrayList<>(); + long total = 0; + int last = 0; + for (MultipartStorage.Part requested : parts) { + if (requested.number() <= last || requested.number() > 10000) + throw new StoreException(400, "InvalidPartOrder", "Parts must be in ascending order"); + last = requested.number(); + StoredPart part = storedPart(connection, uploadId, requested.number()); + if (part == null || !part.etag().equals(requested.etag().replace("\"", ""))) + throw new StoreException(400, "InvalidPart", "Part ETag mismatch"); + if (part.length() > maxObject - total) + throw new StoreException(413, "EntityTooLarge", "Object exceeds the configured size limit"); + total += part.length(); + MessageDigest partHash = digest("MD5"); + for (Segment segment : part.segments()) { + byte[] bytes = readableSegment(segment); + if (bytes == null) throw new StoreException(503, "SlowDown", "A part has no verified replica"); + fullHash.update(bytes); + partHash.update(bytes); + selected.add(segment); + } + byte[] md5 = partHash.digest(); + if (!part.etag().equals(HexFormat.of().formatHex(md5))) + throw new StoreException(503, "SlowDown", "A part failed integrity verification"); + etagHash.update(md5); + } + long previous = currentLength(connection, bucket, key); + if (maxTotal - (used - Math.max(0, previous)) - (staged - uploadBytes) < total) + throw new StoreException(507, "InsufficientStorage", "Store capacity limit reached"); + Metadata metadata = new Metadata(total, Instant.now().toEpochMilli(), + HexFormat.of().formatHex(etagHash.digest()) + "-" + parts.size(), fullHash.digest(), + bucket, key, upload.contentType()); + UUID generation = UUID.randomUUID(); + try (PreparedStatement insert = connection.prepareStatement( + "INSERT INTO cluster_segments (generation, ordinal, segment_id, length, sha256, replicas, replica_ids) VALUES (?, ?, ?, ?, ?, 'v2', ?)")) { + for (int ordinal = 0; ordinal < selected.size(); ordinal++) { + Segment segment = selected.get(ordinal); + insert.setObject(1, generation); + insert.setInt(2, ordinal); + insert.setObject(3, segment.id()); + insert.setInt(4, segment.length()); + insert.setBytes(5, segment.hash()); + insert.setArray(6, connection.createArrayOf("uuid", segment.replicas().toArray())); + insert.addBatch(); + } + insert.executeBatch(); + } + try (PreparedStatement update = connection.prepareStatement( + "INSERT INTO cluster_objects VALUES (?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT (bucket, object_key) DO UPDATE SET generation=EXCLUDED.generation, length=EXCLUDED.length, modified=EXCLUDED.modified, etag=EXCLUDED.etag, sha256=EXCLUDED.sha256, content_type=EXCLUDED.content_type")) { + bindObject(update, metadata, generation); + update.executeUpdate(); + } + try (PreparedStatement update = connection.prepareStatement("UPDATE cluster_usage SET used_bytes=? WHERE bucket=?")) { + update.setLong(1, used - Math.max(0, previous) + total); + update.setString(2, bucket); + update.executeUpdate(); + } + try (PreparedStatement delete = connection.prepareStatement("DELETE FROM cluster_tombstones WHERE bucket=? AND object_key=?")) { + delete.setString(1, bucket); + delete.setString(2, key); + delete.executeUpdate(); + } + try (PreparedStatement delete = connection.prepareStatement("DELETE FROM cluster_uploads WHERE upload_id=?")) { + delete.setObject(1, uploadId); + delete.executeUpdate(); + } + connection.commit(); + return metadata; + } catch (SQLException | IOException | RuntimeException error) { + connection.rollback(); + if (error instanceof SQLException sql) throw databaseError(sql); + if (error instanceof IOException io) throw io; + throw error; + } + } catch (SQLException error) { throw databaseError(error); } + } + + @Override public void abort(String id, String bucket, String key) throws IOException { + UUID uploadId = uploadId(id); + try (Connection connection = connect()) { + connection.setAutoCommit(false); + try { + lockUsage(connection, bucket); + upload(connection, uploadId, bucket, key, true); + try (PreparedStatement delete = connection.prepareStatement("DELETE FROM cluster_uploads WHERE upload_id=?")) { + delete.setObject(1, uploadId); + delete.executeUpdate(); + } + connection.commit(); + } catch (SQLException | RuntimeException error) { + connection.rollback(); + if (error instanceof SQLException sql) throw databaseError(sql); + throw error; + } + } catch (SQLException error) { throw databaseError(error); } + } + + @Override public PartPage listParts(String id, String bucket, String key, int marker, int maxParts) + throws IOException { + UUID uploadId = uploadId(id); + try (Connection connection = connect()) { + upload(connection, uploadId, bucket, key, false); + List parts = new ArrayList<>(); + boolean truncated = false; + try (PreparedStatement query = connection.prepareStatement( + "SELECT part_number, length, etag, modified FROM cluster_upload_parts WHERE upload_id=? AND part_number>? ORDER BY part_number LIMIT ?")) { + query.setObject(1, uploadId); + query.setInt(2, marker); + query.setInt(3, maxParts + 1); + try (ResultSet result = query.executeQuery()) { + while (result.next()) { + if (parts.size() == maxParts) { + truncated = true; + break; + } + parts.add(new PartInfo(result.getInt(1), result.getLong(2), result.getString(3), result.getLong(4))); + } + } + } + int next = parts.isEmpty() ? marker : parts.getLast().number(); + return new PartPage(parts, next, truncated); + } catch (SQLException error) { throw databaseError(error); } + } + + @Override public List listUploads(String bucket, String prefix) throws IOException { + if (!configuredBucket.equals(bucket)) throw new StoreException(404, "NoSuchBucket", "Bucket not found"); + List uploads = new ArrayList<>(); + try (Connection connection = connect(); PreparedStatement query = connection.prepareStatement( + "SELECT upload_id, object_key, created_at FROM cluster_uploads WHERE bucket=?")) { + query.setString(1, bucket); + try (ResultSet result = query.executeQuery()) { + while (result.next()) { + String key = result.getString(2); + if (key.startsWith(prefix)) + uploads.add(new UploadInfo(result.getObject(1).toString(), key, result.getLong(3))); + } + } + } catch (SQLException error) { throw databaseError(error); } + uploads.sort(Comparator.comparing(UploadInfo::key).thenComparing(UploadInfo::id)); + return uploads; + } + + @Override public int activeUploads() { + try (Connection connection = connect(); PreparedStatement query = connection.prepareStatement( + "SELECT count(*) FROM cluster_uploads WHERE bucket=?")) { + query.setString(1, configuredBucket); + try (ResultSet result = query.executeQuery()) { + result.next(); + return result.getInt(1); + } + } catch (SQLException error) { throw new IllegalStateException("Could not count multipart uploads", error); } + } + + @Override public long stagedBytes() { + try (Connection connection = connect()) { return stagedBytes(connection, configuredBucket); } + catch (SQLException error) { throw new IllegalStateException("Could not count staged bytes", error); } + } + @Override public OpenObject open(String bucket, String key) throws IOException { try (Connection connection = connect()) { connection.setAutoCommit(false); @@ -353,6 +628,106 @@ final class ClusterStore implements ObjectStorage { } } } + + private static UUID uploadId(String id) { + try { + if (id == null || !id.matches("[0-9a-f-]{36}")) throw new IllegalArgumentException(); + return UUID.fromString(id); + } catch (IllegalArgumentException error) { + throw new StoreException(404, "NoSuchUpload", "Upload not found"); + } + } + + private static Upload upload(Connection connection, UUID id, String bucket, String key, boolean lock) + throws SQLException { + String sql = "SELECT bucket, object_key, content_type, created_at FROM cluster_uploads WHERE upload_id=?" + + (lock ? " FOR UPDATE" : ""); + try (PreparedStatement query = connection.prepareStatement(sql)) { + query.setObject(1, id); + try (ResultSet result = query.executeQuery()) { + if (!result.next() || !result.getString(1).equals(bucket) || !result.getString(2).equals(key)) + throw new StoreException(404, "NoSuchUpload", "Upload not found"); + return new Upload(result.getString(3)); + } + } + } + + private static long stagedBytes(Connection connection, String bucket) throws SQLException { + try (PreparedStatement query = connection.prepareStatement( + "SELECT COALESCE(sum(p.length), 0) FROM cluster_upload_parts p JOIN cluster_uploads u USING (upload_id) WHERE u.bucket=?")) { + query.setString(1, bucket); + try (ResultSet result = query.executeQuery()) { + result.next(); + return result.getLong(1); + } + } + } + + private static long uploadLength(Connection connection, UUID id) throws SQLException { + try (PreparedStatement query = connection.prepareStatement( + "SELECT COALESCE(sum(length), 0) FROM cluster_upload_parts WHERE upload_id=?")) { + query.setObject(1, id); + try (ResultSet result = query.executeQuery()) { + result.next(); + return result.getLong(1); + } + } + } + + private static long partLength(Connection connection, UUID id, int number) throws SQLException { + try (PreparedStatement query = connection.prepareStatement( + "SELECT length FROM cluster_upload_parts WHERE upload_id=? AND part_number=?")) { + query.setObject(1, id); + query.setInt(2, number); + try (ResultSet result = query.executeQuery()) { + return result.next() ? result.getLong(1) : 0; + } + } + } + + private static StoredPart storedPart(Connection connection, UUID id, int number) throws SQLException, IOException { + long length; + String etag; + try (PreparedStatement query = connection.prepareStatement( + "SELECT length, etag FROM cluster_upload_parts WHERE upload_id=? AND part_number=?")) { + query.setObject(1, id); + query.setInt(2, number); + try (ResultSet result = query.executeQuery()) { + if (!result.next()) return null; + length = result.getLong(1); + etag = result.getString(2); + } + } + List segments = new ArrayList<>(); + try (PreparedStatement query = connection.prepareStatement( + "SELECT ordinal, segment_id, length, sha256, replica_ids FROM cluster_upload_segments WHERE upload_id=? AND part_number=? ORDER BY ordinal")) { + query.setObject(1, id); + query.setInt(2, number); + try (ResultSet result = query.executeQuery()) { + long total = 0; + while (result.next()) { + if (result.getInt(1) != segments.size()) throw new IOException("Incomplete multipart manifest"); + Segment segment = new Segment((UUID) result.getObject(2), result.getInt(3), + result.getBytes(4), replicaIds(result, 5)); + total = Math.addExact(total, segment.length()); + segments.add(segment); + } + if (total != length) throw new IOException("Incomplete multipart manifest"); + } + } + return new StoredPart(length, etag, segments); + } + + private byte[] readableSegment(Segment segment) { + for (UUID id : segment.replicas()) { + int index = nodes.index(id); + if (index < 0) continue; + byte[] bytes = readableReplica(index, segment); + if (bytes != null) return bytes; + } + return null; + } + private long currentLength(Connection connection, String bucket, String key) throws SQLException { try (PreparedStatement query = connection.prepareStatement("SELECT length FROM cluster_objects WHERE bucket=? AND object_key=?")) { query.setString(1, bucket); @@ -419,14 +794,18 @@ final class ClusterStore implements ObjectStorage { try (Connection reader = connect()) { reader.setAutoCommit(false); try (PreparedStatement query = reader.prepareStatement( - "SELECT s.generation, s.ordinal, s.segment_id, s.length, s.sha256, s.replica_ids, s.placement_version FROM cluster_segments s JOIN cluster_objects o ON o.generation=s.generation ORDER BY s.generation, s.ordinal")) { + "SELECT s.generation, 0, s.ordinal, s.segment_id, s.length, s.sha256, s.replica_ids, s.placement_version " + + "FROM cluster_segments s JOIN cluster_objects o ON o.generation=s.generation " + + "UNION ALL SELECT s.upload_id, s.part_number, s.ordinal, s.segment_id, s.length, s.sha256, " + + "s.replica_ids, s.placement_version FROM cluster_upload_segments s " + + "ORDER BY 1, 2, 3")) { query.setFetchSize(128); try (ResultSet result = query.executeQuery()) { while (result.next()) { scanned++; RepairTarget target = new RepairTarget((UUID) result.getObject(1), result.getInt(2), - result.getLong(7), new Segment((UUID) result.getObject(3), result.getInt(4), - result.getBytes(5), replicaIds(result, 6))); + result.getInt(3), result.getLong(8), new Segment((UUID) result.getObject(4), + result.getInt(5), result.getBytes(6), replicaIds(result, 7))); Segment segment = target.segment(); byte[] copy = null; Set healthy = new HashSet<>(); @@ -458,12 +837,18 @@ final class ClusterStore implements ObjectStorage { Set listed = new java.util.LinkedHashSet<>(segment.replicas()); listed.addAll(healthy); if (listed.size() != segment.replicas().size()) { + String table = target.part() == 0 ? "cluster_segments" : "cluster_upload_segments"; + String identity = target.part() == 0 ? "generation=? AND ordinal=?" : + "upload_id=? AND part_number=? AND ordinal=?"; try (Connection writer = connect(); PreparedStatement update = writer.prepareStatement( - "UPDATE cluster_segments SET replica_ids=?, placement_version=placement_version+1 WHERE generation=? AND ordinal=? AND placement_version=?")) { + "UPDATE " + table + " SET replica_ids=?, placement_version=placement_version+1 WHERE " + + identity + " AND placement_version=?")) { update.setArray(1, writer.createArrayOf("uuid", listed.toArray())); - update.setObject(2, target.generation()); - update.setInt(3, target.ordinal()); - update.setLong(4, target.version()); + update.setObject(2, target.id()); + int next = 3; + if (target.part() != 0) update.setInt(next++, target.part()); + update.setInt(next++, target.ordinal()); + update.setLong(next, target.version()); update.executeUpdate(); } } diff --git a/src/cloud/lunarsky/store/Main.java b/src/cloud/lunarsky/store/Main.java index afac05a..731bf6b 100644 --- a/src/cloud/lunarsky/store/Main.java +++ b/src/cloud/lunarsky/store/Main.java @@ -86,6 +86,16 @@ public final class Main { } private void handleBucket(HttpExchange exchange, Map query, String hash) throws IOException { + if (exchange.getRequestMethod().equals("GET") && query.containsKey("uploads")) { + if (!query.get("uploads").isEmpty() || + !query.keySet().stream().allMatch(java.util.Set.of("uploads", "prefix", "key-marker", + "upload-id-marker", "max-uploads", "x-id")::contains) || + (query.containsKey("x-id") && !"ListMultipartUploads".equals(query.get("x-id")))) + unsupported("Bucket operation"); + requireEmptyBody(exchange, hash); + listUploads(exchange, query); + return; + } if (!exchange.getRequestMethod().equals("GET") || !"2".equals(query.get("list-type")) || !query.keySet().stream().allMatch(java.util.Set.of("list-type", "prefix", "delimiter", "max-keys", "continuation-token", "start-after", "encoding-type", "x-id")::contains) || @@ -183,12 +193,16 @@ public final class Main { query.keySet().stream().allMatch(java.util.Set.of("uploads", "x-id")::contains) && (!query.containsKey("x-id") || query.get("x-id").equals("CreateMultipartUpload")); if (!query.containsKey("uploadId") || - !query.keySet().stream().allMatch(java.util.Set.of("uploadId", "partNumber", "x-id")::contains)) + !query.keySet().stream().allMatch(java.util.Set.of("uploadId", "partNumber", + "part-number-marker", "max-parts", "x-id")::contains)) return false; String xId = query.get("x-id"); if (method.equals("PUT")) return query.containsKey("partNumber") && + !query.containsKey("part-number-marker") && !query.containsKey("max-parts") && (xId == null || xId.equals("UploadPart")); if (query.containsKey("partNumber")) return false; + if (method.equals("GET")) return xId == null || xId.equals("ListParts"); + if (query.containsKey("part-number-marker") || query.containsKey("max-parts")) return false; return (method.equals("POST") && (xId == null || xId.equals("CompleteMultipartUpload"))) || (method.equals("DELETE") && (xId == null || xId.equals("AbortMultipartUpload"))); } @@ -231,10 +245,81 @@ public final class Main { multipart.abort(id, bucket, key); exchange.sendResponseHeaders(204, -1); } + case "GET" -> { + requireEmptyBody(exchange, hash); + listParts(exchange, id, key, query); + } default -> unsupported("Multipart operation"); } } + private void listParts(HttpExchange exchange, String id, String key, Map query) throws IOException { + int marker = boundedNumber(query.get("part-number-marker"), 0, 10000, 0); + int maxParts = boundedNumber(query.get("max-parts"), 1, 1000, 1000); + MultipartStorage.PartPage page = multipart.listParts(id, bucket, key, marker, maxParts); + StringBuilder body = new StringBuilder("").append(xml(bucket)) + .append("").append(xml(key)).append("").append(xml(id)) + .append("").append(marker) + .append("").append(page.nextMarker()) + .append("").append(maxParts) + .append("").append(page.truncated()).append(""); + for (MultipartStorage.PartInfo part : page.parts()) { + body.append("").append(part.number()).append("") + .append(Instant.ofEpochMilli(part.modified())).append(""") + .append(part.etag()).append(""").append(part.length()).append(""); + } + sendXml(exchange, 200, body.append("").toString()); + } + + private void listUploads(HttpExchange exchange, Map query) throws IOException { + String prefix = query.getOrDefault("prefix", ""); + String marker = query.getOrDefault("key-marker", ""); + String uploadMarker = query.getOrDefault("upload-id-marker", ""); + if (!uploadMarker.isEmpty() && marker.isEmpty()) + throw new StoreException(400, "InvalidArgument", "Upload ID marker requires a key marker"); + int maximum = boundedNumber(query.get("max-uploads"), 1, 1000, 1000); + List uploads = multipart.listUploads(bucket, prefix); + List page = new ArrayList<>(); + boolean truncated = false; + for (MultipartStorage.UploadInfo upload : uploads) { + if (upload.key().compareTo(marker) < 0 || + (upload.key().equals(marker) && upload.id().compareTo(uploadMarker) <= 0)) continue; + if (page.size() == maximum) { + truncated = true; + break; + } + page.add(upload); + } + StringBuilder body = new StringBuilder("").append(xml(bucket)) + .append("").append(xml(marker)).append("") + .append(xml(uploadMarker)).append("").append(maximum) + .append("").append(truncated).append(""); + if (truncated) { + MultipartStorage.UploadInfo last = page.getLast(); + body.append("").append(xml(last.key())).append("") + .append(xml(last.id())).append(""); + } + for (MultipartStorage.UploadInfo upload : page) { + body.append("").append(xml(upload.key())).append("") + .append(xml(upload.id())).append("") + .append(Instant.ofEpochMilli(upload.created())).append(""); + } + sendXml(exchange, 200, body.append("").toString()); + } + + private static int boundedNumber(String text, int minimum, int maximum, int defaultValue) { + if (text == null) return defaultValue; + int value; + try { + value = Integer.parseInt(text); + } catch (NumberFormatException error) { + throw new StoreException(400, "InvalidArgument", "Invalid listing limit or marker"); + } + if (value < minimum || value > maximum) + throw new StoreException(400, "InvalidArgument", "Invalid listing limit or marker"); + return value; + } + private static long contentLength(com.sun.net.httpserver.Headers headers) { String text = SigV4.single(headers, "content-length"); if (text == null) return -1; @@ -489,7 +574,7 @@ public final class Main { required(env, "POSTGRES_PASSWORD"), bucket, java.util.Arrays.stream(urls).map(URI::create).toList(), required(env, "CLUSTER_TOKEN"), null, maxObject, maxTotal, "true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS"))); - multipart = new UnavailableMultipart(); + multipart = (ClusterStore) store; } else if (mode.equals("disk")) { DiskStore disk = new DiskStore(Path.of(env.getOrDefault("DATA_DIR", "/data")), maxObject, maxTotal); store = disk; diff --git a/src/cloud/lunarsky/store/MultipartStorage.java b/src/cloud/lunarsky/store/MultipartStorage.java index 7caf0b2..fdf0d59 100644 --- a/src/cloud/lunarsky/store/MultipartStorage.java +++ b/src/cloud/lunarsky/store/MultipartStorage.java @@ -6,12 +6,17 @@ import java.util.List; interface MultipartStorage { record Part(int number, String etag) {} + record PartInfo(int number, long length, String etag, long modified) {} + record PartPage(List parts, int nextMarker, boolean truncated) {} + record UploadInfo(String id, String key, long created) {} String create(String bucket, String key, String contentType) throws IOException; String putPart(String id, String bucket, String key, int number, InputStream input, long length, String expectedHash, String checksum) throws IOException; ObjectStorage.Metadata complete(String id, String bucket, String key, List parts) throws IOException; void abort(String id, String bucket, String key) throws IOException; + PartPage listParts(String id, String bucket, String key, int marker, int maxParts) throws IOException; + List listUploads(String bucket, String prefix) throws IOException; int activeUploads(); long stagedBytes(); } diff --git a/src/cloud/lunarsky/store/MultipartStore.java b/src/cloud/lunarsky/store/MultipartStore.java index a0068aa..e9684af 100644 --- a/src/cloud/lunarsky/store/MultipartStore.java +++ b/src/cloud/lunarsky/store/MultipartStore.java @@ -164,6 +164,48 @@ final class MultipartStore implements MultipartStorage { remove(upload(id, bucket, key)); } + @Override public synchronized PartPage listParts(String id, String bucket, String key, + int marker, int maxParts) throws IOException { + Path dir = upload(id, bucket, key); + List parts = new ArrayList<>(); + boolean truncated = false; + try (var files = Files.list(dir)) { + for (Path file : files.filter(path -> path.getFileName().toString().matches("part-[0-9]{5}")) + .sorted().toList()) { + int number = Integer.parseInt(file.getFileName().toString().substring(5)); + if (number <= marker) continue; + if (parts.size() == maxParts) { + truncated = true; + break; + } + MessageDigest md5 = digest("MD5"); + try (InputStream input = Files.newInputStream(file)) { + byte[] buffer = new byte[65536]; + int count; + while ((count = input.read(buffer)) != -1) md5.update(buffer, 0, count); + } + parts.add(new PartInfo(number, Files.size(file), SigV4.hex(md5.digest()), + Files.getLastModifiedTime(file).toMillis())); + } + } + int next = parts.isEmpty() ? marker : parts.getLast().number(); + return new PartPage(parts, next, truncated); + } + + @Override public synchronized List listUploads(String bucket, String prefix) throws IOException { + List uploads = new ArrayList<>(); + try (var dirs = Files.list(root)) { + for (Path dir : dirs.filter(Files::isDirectory).toList()) { + Upload upload = readUpload(dir); + if (upload.bucket().equals(bucket) && upload.key().startsWith(prefix)) + uploads.add(new UploadInfo(dir.getFileName().toString(), upload.key(), + Files.getLastModifiedTime(dir.resolve("manifest")).toMillis())); + } + } + uploads.sort(Comparator.comparing(UploadInfo::key).thenComparing(UploadInfo::id)); + return uploads; + } + private Path upload(String id, String bucket, String key) throws IOException { if (!id.matches("[0-9a-f-]{36}")) throw new StoreException(404, "NoSuchUpload", "Upload not found"); Path dir = root.resolve(id); diff --git a/src/cloud/lunarsky/store/SchemaMigrator.java b/src/cloud/lunarsky/store/SchemaMigrator.java index f52398a..ffd5648 100644 --- a/src/cloud/lunarsky/store/SchemaMigrator.java +++ b/src/cloud/lunarsky/store/SchemaMigrator.java @@ -21,7 +21,7 @@ final class SchemaMigrator { result.next(); version = result.getInt(1); } - if (version > 2) throw new IOException("Metadata schema is newer than this ObjectStore build"); + if (version > 3) throw new IOException("Metadata schema is newer than this ObjectStore build"); if (version < 1) { statement.execute("CREATE TABLE IF NOT EXISTS cluster_usage (bucket text PRIMARY KEY, used_bytes bigint NOT NULL CHECK (used_bytes >= 0))"); statement.execute("CREATE TABLE IF NOT EXISTS cluster_objects (bucket text NOT NULL, object_key text COLLATE \"C\" NOT NULL, generation uuid NOT NULL, length bigint NOT NULL, modified bigint NOT NULL, etag text NOT NULL, sha256 bytea NOT NULL, content_type text NOT NULL, PRIMARY KEY (bucket, object_key))"); @@ -37,6 +37,12 @@ final class SchemaMigrator { statement.execute("CREATE TABLE IF NOT EXISTS cluster_format (singleton integer PRIMARY KEY CHECK (singleton=1), version integer NOT NULL)"); statement.execute("INSERT INTO cluster_schema_migrations VALUES (2)"); } + if (version < 3) { + statement.execute("CREATE TABLE cluster_uploads (upload_id uuid PRIMARY KEY, bucket text NOT NULL, object_key text COLLATE \"C\" NOT NULL, content_type text NOT NULL, created_at bigint NOT NULL)"); + statement.execute("CREATE TABLE cluster_upload_parts (upload_id uuid NOT NULL REFERENCES cluster_uploads(upload_id) ON DELETE CASCADE, part_number integer NOT NULL CHECK (part_number BETWEEN 1 AND 10000), length bigint NOT NULL CHECK (length >= 0), etag text NOT NULL, modified bigint NOT NULL, PRIMARY KEY (upload_id, part_number))"); + statement.execute("CREATE TABLE cluster_upload_segments (upload_id uuid NOT NULL, part_number integer NOT NULL, ordinal integer NOT NULL, segment_id uuid NOT NULL, length integer NOT NULL, sha256 bytea NOT NULL, replica_ids uuid[] NOT NULL, placement_version bigint NOT NULL DEFAULT 0, PRIMARY KEY (upload_id, part_number, ordinal), FOREIGN KEY (upload_id, part_number) REFERENCES cluster_upload_parts(upload_id, part_number) ON DELETE CASCADE)"); + statement.execute("INSERT INTO cluster_schema_migrations VALUES (3)"); + } statement.execute("INSERT INTO cluster_format SELECT 1, CASE WHEN EXISTS (SELECT 1 FROM cluster_segments WHERE replica_ids IS NULL) THEN 1 ELSE 2 END WHERE NOT EXISTS (SELECT 1 FROM cluster_format)"); } try (PreparedStatement insert = connection.prepareStatement("INSERT INTO cluster_usage VALUES (?, 0) ON CONFLICT DO NOTHING")) { diff --git a/src/cloud/lunarsky/store/UnavailableMultipart.java b/src/cloud/lunarsky/store/UnavailableMultipart.java deleted file mode 100644 index 258e1e0..0000000 --- a/src/cloud/lunarsky/store/UnavailableMultipart.java +++ /dev/null @@ -1,17 +0,0 @@ -package cloud.lunarsky.store; - -import java.io.InputStream; -import java.util.List; - -final class UnavailableMultipart implements MultipartStorage { - private StoreException unavailable() { - return new StoreException(501, "NotImplemented", "Multipart uploads are unavailable in the local cluster prototype"); - } - @Override public String create(String bucket, String key, String contentType) { throw unavailable(); } - @Override public String putPart(String id, String bucket, String key, int number, InputStream input, - long length, String expectedHash, String checksum) { throw unavailable(); } - @Override public ObjectStorage.Metadata complete(String id, String bucket, String key, List parts) { throw unavailable(); } - @Override public void abort(String id, String bucket, String key) { throw unavailable(); } - @Override public int activeUploads() { return 0; } - @Override public long stagedBytes() { return 0; } -} diff --git a/test/cloud/lunarsky/store/ClusterIntegrationTest.java b/test/cloud/lunarsky/store/ClusterIntegrationTest.java index f93c8ab..05aedfa 100644 --- a/test/cloud/lunarsky/store/ClusterIntegrationTest.java +++ b/test/cloud/lunarsky/store/ClusterIntegrationTest.java @@ -3,7 +3,10 @@ package cloud.lunarsky.store; import java.io.ByteArrayInputStream; import java.net.URI; import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; import java.util.Arrays; +import java.util.HexFormat; +import java.util.List; import java.util.Map; public final class ClusterIntegrationTest { @@ -62,6 +65,60 @@ public final class ClusterIntegrationTest { } System.out.println("Cluster degraded test passed"); } + case "multipart-stage" -> { + String key = "cluster-test/multipart"; + String id = store.create(bucket, key, "text/plain"); + byte[] first = "hello ".getBytes(StandardCharsets.UTF_8); + byte[] second = "world".getBytes(StandardCharsets.UTF_8); + putPart(store, id, bucket, key, 1, "old".getBytes(StandardCharsets.UTF_8)); + putPart(store, id, bucket, key, 1, first); + putPart(store, id, bucket, key, 2, second); + require(store.activeUploads() == 1, "Upload was not retained"); + require(store.stagedBytes() == first.length + second.length, "Replaced part was counted twice"); + require(store.listUploads(bucket, key).size() == 1, "Upload listing missed the staged upload"); + var firstPage = store.listParts(id, bucket, key, 0, 1); + require(firstPage.truncated() && firstPage.parts().size() == 1 && firstPage.nextMarker() == 1, + "Part listing did not paginate"); + require(store.listParts(id, bucket, key, 1, 1).parts().getFirst().number() == 2, + "Part marker skipped the second part"); + try { + store.open(bucket, key); + throw new AssertionError("Incomplete upload became visible"); + } catch (StoreException error) { require(error.status == 404, "Wrong incomplete-upload status"); } + System.out.println("Cluster multipart parts staged and listed"); + } + case "multipart-complete" -> { + String key = "cluster-test/multipart"; + var uploads = store.listUploads(bucket, key); + require(uploads.size() == 1, "Upload did not survive gateway restart"); + String id = uploads.getFirst().id(); + var listed = store.listParts(id, bucket, key, 0, 1000).parts(); + require(listed.size() == 2, "Staged parts were lost"); + try { + store.complete(id, bucket, key, List.of(new MultipartStorage.Part(1, "0".repeat(32)))); + throw new AssertionError("Wrong part ETag was accepted"); + } catch (StoreException error) { require(error.status == 400, "Wrong ETag rejection status"); } + var completed = store.complete(id, bucket, key, List.of( + new MultipartStorage.Part(1, listed.get(0).etag()), + new MultipartStorage.Part(2, listed.get(1).etag()))); + byte[] expected = "hello world".getBytes(StandardCharsets.UTF_8); + require(completed.length() == expected.length && completed.contentType().equals("text/plain"), + "Completed object metadata is wrong"); + MessageDigest digest = MessageDigest.getInstance("MD5"); + digest.update(HexFormat.of().parseHex(listed.get(0).etag())); + digest.update(HexFormat.of().parseHex(listed.get(1).etag())); + require(completed.etag().equals(HexFormat.of().formatHex(digest.digest()) + "-2"), + "Multipart ETag is wrong"); + try (var opened = store.open(bucket, key)) { + require(Arrays.equals(opened.stream().readAllBytes(), expected), "Completed multipart body is wrong"); + } + require(store.activeUploads() == 0 && store.stagedBytes() == 0, "Completed parts still count as staged"); + String aborted = store.create(bucket, "cluster-test/aborted", "text/plain"); + putPart(store, aborted, bucket, "cluster-test/aborted", 1, expected); + store.abort(aborted, bucket, "cluster-test/aborted"); + require(store.activeUploads() == 0 && store.stagedBytes() == 0, "Aborted parts still count as staged"); + System.out.println("Cluster multipart completion survived restart and node loss"); + } case "quorum-lost" -> { require(!store.ready(), "One available node must not be ready"); try { @@ -151,6 +208,11 @@ public final class ClusterIntegrationTest { store.put(bucket, key, new ByteArrayInputStream(data), data.length, SigV4.hex(SigV4.hash(data)), null, createOnly, "application/octet-stream"); } + private static void putPart(ClusterStore store, String id, String bucket, String key, int number, byte[] data) + throws Exception { + store.putPart(id, bucket, key, number, new ByteArrayInputStream(data), data.length, + SigV4.hex(SigV4.hash(data)), null); + } private static void require(boolean condition, String message) { if (!condition) throw new AssertionError(message); } diff --git a/test/cloud/lunarsky/store/HttpTest.java b/test/cloud/lunarsky/store/HttpTest.java index 96bd107..5e98494 100644 --- a/test/cloud/lunarsky/store/HttpTest.java +++ b/test/cloud/lunarsky/store/HttpTest.java @@ -144,6 +144,16 @@ public final class HttpTest { "?partNumber=2&uploadId=" + upload), "PUT", second, Map.of()), HttpResponse.BodyHandlers.ofByteArray()); status(200, partOne); status(200, partTwo); + var parts = client.send(signedUri(URI.create(base + "/objects/" + movie + + "?uploadId=" + upload + "&max-parts=1"), "GET", new byte[0], Map.of()), + HttpResponse.BodyHandlers.ofString()); + if (parts.statusCode() != 200 || !parts.body().contains("true") || + !parts.body().contains("1")) + throw new AssertionError("Multipart part listing failed: " + parts.body()); + var uploads = client.send(signedUri(URI.create(base + "/objects?uploads&prefix=folder%2F"), + "GET", new byte[0], Map.of()), HttpResponse.BodyHandlers.ofString()); + if (uploads.statusCode() != 200 || !uploads.body().contains("" + upload + "")) + throw new AssertionError("Multipart upload listing failed: " + uploads.body()); String completion = "1" + partOne.headers().firstValue("etag").orElseThrow() + "2" + @@ -156,6 +166,8 @@ public final class HttpTest { if (!"hello world".equals(new String(assembled.body(), StandardCharsets.UTF_8)) || !"video/mp4".equals(assembled.headers().firstValue("content-type").orElse(""))) throw new AssertionError("Completed multipart object mismatch"); + status(404, client.send(signedUri(URI.create(base + "/objects/" + movie + "?uploadId=" + upload), + "GET", new byte[0], Map.of()), HttpResponse.BodyHandlers.ofByteArray())); var abandoned = client.send(signedUri(URI.create(base + "/objects/abandoned?uploads="), "POST", new byte[0], Map.of()), HttpResponse.BodyHandlers.ofString()); String abandonedId = abandoned.body().split("")[1].split("")[0];