Add cluster multipart uploads and listings

This commit is contained in:
admin committed 2026-10-10 08:48:15 +02:00
1 parent 11bfe71124
commit 00d2af2fb0
12 files changed
+660 -38

No files matched your search

+5 -3
View File
@@ -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 - ✅ `PutObject`, `GetObject`, `HeadObject`, and `DeleteObject` in both modes
- ✅ Single-range GET and `ListObjectsV2` in both modes - ✅ Single-range GET and `ListObjectsV2` in both modes
- ✅ SHA-256 payload verification and `x-amz-checksum-sha256` in both modes - ✅ SHA-256 payload verification and `x-amz-checksum-sha256` in both modes
- ✅ `CreateMultipartUpload`, `UploadPart`, `CompleteMultipartUpload`, and `AbortMultipartUpload` in single-node mode - ✅ `CreateMultipartUpload`, `UploadPart`, `CompleteMultipartUpload`, and `AbortMultipartUpload` in both modes
- ⬜ Multipart uploads in cluster mode - ✅ `ListParts` and `ListMultipartUploads` in both modes
- ⬜ Presigned URLs and streaming Signature V4 uploads - ⬜ Presigned URLs and streaming Signature V4 uploads
- ⬜ `CopyObject`, `ListParts`, and `ListMultipartUploads` - ⬜ `CopyObject`
- ⬜ `Content-MD5` and checksum algorithms other than SHA-256 - ⬜ `Content-MD5` and checksum algorithms other than SHA-256
- ⬜ Bucket creation and listing, object versioning, ACLs, tags, and user metadata - ⬜ 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. 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. 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 ## Migrating a local cluster
+2 -2
View File
@@ -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. | | `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. | | `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. | | `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. | | `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. 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). `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).
+28 -2
View File
@@ -6,6 +6,7 @@ import http.client
import pathlib import pathlib
import sys import sys
import urllib.parse 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() 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}") f"SignedHeaders={signed_names},Signature={signature}")
connection = http.client.HTTPConnection("127.0.0.1", port, timeout=30) connection = http.client.HTTPConnection("127.0.0.1", port, timeout=30)
try: 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() response = connection.getresponse()
return response.status, response.read(), response.headers return response.status, response.read(), response.headers
finally: finally:
@@ -81,4 +82,29 @@ status, _, _ = request("DELETE", key)
assert status == 204, status assert status == 204, status
status, _, _ = request("GET", key) status, _, _ = request("GET", key)
assert status == 404, status 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"<IsTruncated>true</IsTruncated>" 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 = "<CompleteMultipartUpload>" + "".join(
f"<Part><PartNumber>{number}</PartNumber><ETag>{etag}</ETag></Part>"
for number, etag in enumerate(part_etags, start=1)) + "</CompleteMultipartUpload>"
status, content, _ = request("POST", multipart_key + f"?uploadId={upload_id}", completion.encode(),
{"content-type": "application/xml"})
assert status == 200 and b"<CompleteMultipartUploadResult>" 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")
+14
View File
@@ -20,10 +20,24 @@ wait_ready() {
done done
} }
run_phase basic 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 same-host
run_phase concurrent run_phase concurrent
compose stop node-a compose stop node-a
run_phase degraded run_phase degraded
run_phase multipart-complete
compose stop node-b compose stop node-b
run_phase quorum-lost run_phase quorum-lost
compose start node-a node-b compose start node-a node-b
+396 -11
View File
@@ -18,15 +18,18 @@ import java.sql.SQLException;
import java.time.Instant; import java.time.Instant;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Base64; import java.util.Base64;
import java.util.Comparator;
import java.util.HexFormat; import java.util.HexFormat;
import java.util.HashSet; import java.util.HashSet;
import java.util.List; import java.util.List;
import java.util.Set; import java.util.Set;
import java.util.UUID; 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<UUID> replicas) {} private record Segment(UUID id, int length, byte[] hash, List<UUID> 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<Segment> segments) {}
record RepairReport(int scanned, int restored, int underReplicated, int unrecoverable) {} record RepairReport(int scanned, int restored, int underReplicated, int unrecoverable) {}
private final String jdbcUrl, user, password, configuredBucket; private final String jdbcUrl, user, password, configuredBucket;
private final NodeClient nodes; private final NodeClient nodes;
@@ -116,7 +119,8 @@ final class ClusterStore implements ObjectStorage {
query.setString(1, bucket); query.setString(1, bucket);
try (ResultSet result = query.executeQuery()) { try (ResultSet result = query.executeQuery()) {
if (!result.next()) throw new SQLException("Bucket quota row is missing"); 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"); throw new StoreException(507, "InsufficientStorage", "Store capacity limit reached");
} }
} }
@@ -168,7 +172,8 @@ final class ClusterStore implements ObjectStorage {
long previous = currentLength(connection, bucket, key); long previous = currentLength(connection, bucket, key);
if (createOnly && previous >= 0) if (createOnly && previous >= 0)
throw new StoreException(412, "PreconditionFailed", "Object already exists"); 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"); throw new StoreException(507, "InsufficientStorage", "Store capacity limit reached");
try (PreparedStatement insert = connection.prepareStatement( try (PreparedStatement insert = connection.prepareStatement(
"INSERT INTO cluster_segments (generation, ordinal, segment_id, length, sha256, replicas, replica_ids) VALUES (?, ?, ?, ?, ?, 'v2', ?)")) { "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); } } 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<Segment> 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<MultipartStorage.Part> 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<Segment> 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<PartInfo> 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<UploadInfo> listUploads(String bucket, String prefix) throws IOException {
if (!configuredBucket.equals(bucket)) throw new StoreException(404, "NoSuchBucket", "Bucket not found");
List<UploadInfo> 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 { @Override public OpenObject open(String bucket, String key) throws IOException {
try (Connection connection = connect()) { try (Connection connection = connect()) {
connection.setAutoCommit(false); 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<Segment> 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 { 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=?")) { try (PreparedStatement query = connection.prepareStatement("SELECT length FROM cluster_objects WHERE bucket=? AND object_key=?")) {
query.setString(1, bucket); query.setString(1, bucket);
@@ -419,14 +794,18 @@ final class ClusterStore implements ObjectStorage {
try (Connection reader = connect()) { try (Connection reader = connect()) {
reader.setAutoCommit(false); reader.setAutoCommit(false);
try (PreparedStatement query = reader.prepareStatement( 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); query.setFetchSize(128);
try (ResultSet result = query.executeQuery()) { try (ResultSet result = query.executeQuery()) {
while (result.next()) { while (result.next()) {
scanned++; scanned++;
RepairTarget target = new RepairTarget((UUID) result.getObject(1), result.getInt(2), RepairTarget target = new RepairTarget((UUID) result.getObject(1), result.getInt(2),
result.getLong(7), new Segment((UUID) result.getObject(3), result.getInt(4), result.getInt(3), result.getLong(8), new Segment((UUID) result.getObject(4),
result.getBytes(5), replicaIds(result, 6))); result.getInt(5), result.getBytes(6), replicaIds(result, 7)));
Segment segment = target.segment(); Segment segment = target.segment();
byte[] copy = null; byte[] copy = null;
Set<UUID> healthy = new HashSet<>(); Set<UUID> healthy = new HashSet<>();
@@ -458,12 +837,18 @@ final class ClusterStore implements ObjectStorage {
Set<UUID> listed = new java.util.LinkedHashSet<>(segment.replicas()); Set<UUID> listed = new java.util.LinkedHashSet<>(segment.replicas());
listed.addAll(healthy); listed.addAll(healthy);
if (listed.size() != segment.replicas().size()) { 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( 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.setArray(1, writer.createArrayOf("uuid", listed.toArray()));
update.setObject(2, target.generation()); update.setObject(2, target.id());
update.setInt(3, target.ordinal()); int next = 3;
update.setLong(4, target.version()); if (target.part() != 0) update.setInt(next++, target.part());
update.setInt(next++, target.ordinal());
update.setLong(next, target.version());
update.executeUpdate(); update.executeUpdate();
} }
} }
+87 -2
View File
@@ -86,6 +86,16 @@ public final class Main {
} }
private void handleBucket(HttpExchange exchange, Map<String, String> query, String hash) throws IOException { private void handleBucket(HttpExchange exchange, Map<String, String> 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")) || 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", !query.keySet().stream().allMatch(java.util.Set.of("list-type", "prefix", "delimiter", "max-keys",
"continuation-token", "start-after", "encoding-type", "x-id")::contains) || "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.keySet().stream().allMatch(java.util.Set.of("uploads", "x-id")::contains) &&
(!query.containsKey("x-id") || query.get("x-id").equals("CreateMultipartUpload")); (!query.containsKey("x-id") || query.get("x-id").equals("CreateMultipartUpload"));
if (!query.containsKey("uploadId") || 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; return false;
String xId = query.get("x-id"); String xId = query.get("x-id");
if (method.equals("PUT")) return query.containsKey("partNumber") && if (method.equals("PUT")) return query.containsKey("partNumber") &&
!query.containsKey("part-number-marker") && !query.containsKey("max-parts") &&
(xId == null || xId.equals("UploadPart")); (xId == null || xId.equals("UploadPart"));
if (query.containsKey("partNumber")) return false; 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"))) || return (method.equals("POST") && (xId == null || xId.equals("CompleteMultipartUpload"))) ||
(method.equals("DELETE") && (xId == null || xId.equals("AbortMultipartUpload"))); (method.equals("DELETE") && (xId == null || xId.equals("AbortMultipartUpload")));
} }
@@ -231,10 +245,81 @@ public final class Main {
multipart.abort(id, bucket, key); multipart.abort(id, bucket, key);
exchange.sendResponseHeaders(204, -1); exchange.sendResponseHeaders(204, -1);
} }
case "GET" -> {
requireEmptyBody(exchange, hash);
listParts(exchange, id, key, query);
}
default -> unsupported("Multipart operation"); default -> unsupported("Multipart operation");
} }
} }
private void listParts(HttpExchange exchange, String id, String key, Map<String, String> 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("<ListPartsResult><Bucket>").append(xml(bucket))
.append("</Bucket><Key>").append(xml(key)).append("</Key><UploadId>").append(xml(id))
.append("</UploadId><PartNumberMarker>").append(marker)
.append("</PartNumberMarker><NextPartNumberMarker>").append(page.nextMarker())
.append("</NextPartNumberMarker><MaxParts>").append(maxParts)
.append("</MaxParts><IsTruncated>").append(page.truncated()).append("</IsTruncated>");
for (MultipartStorage.PartInfo part : page.parts()) {
body.append("<Part><PartNumber>").append(part.number()).append("</PartNumber><LastModified>")
.append(Instant.ofEpochMilli(part.modified())).append("</LastModified><ETag>&quot;")
.append(part.etag()).append("&quot;</ETag><Size>").append(part.length()).append("</Size></Part>");
}
sendXml(exchange, 200, body.append("</ListPartsResult>").toString());
}
private void listUploads(HttpExchange exchange, Map<String, String> 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<MultipartStorage.UploadInfo> uploads = multipart.listUploads(bucket, prefix);
List<MultipartStorage.UploadInfo> 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("<ListMultipartUploadsResult><Bucket>").append(xml(bucket))
.append("</Bucket><KeyMarker>").append(xml(marker)).append("</KeyMarker><UploadIdMarker>")
.append(xml(uploadMarker)).append("</UploadIdMarker><MaxUploads>").append(maximum)
.append("</MaxUploads><IsTruncated>").append(truncated).append("</IsTruncated>");
if (truncated) {
MultipartStorage.UploadInfo last = page.getLast();
body.append("<NextKeyMarker>").append(xml(last.key())).append("</NextKeyMarker><NextUploadIdMarker>")
.append(xml(last.id())).append("</NextUploadIdMarker>");
}
for (MultipartStorage.UploadInfo upload : page) {
body.append("<Upload><Key>").append(xml(upload.key())).append("</Key><UploadId>")
.append(xml(upload.id())).append("</UploadId><Initiated>")
.append(Instant.ofEpochMilli(upload.created())).append("</Initiated></Upload>");
}
sendXml(exchange, 200, body.append("</ListMultipartUploadsResult>").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) { private static long contentLength(com.sun.net.httpserver.Headers headers) {
String text = SigV4.single(headers, "content-length"); String text = SigV4.single(headers, "content-length");
if (text == null) return -1; if (text == null) return -1;
@@ -489,7 +574,7 @@ public final class Main {
required(env, "POSTGRES_PASSWORD"), bucket, required(env, "POSTGRES_PASSWORD"), bucket,
java.util.Arrays.stream(urls).map(URI::create).toList(), required(env, "CLUSTER_TOKEN"), null, java.util.Arrays.stream(urls).map(URI::create).toList(), required(env, "CLUSTER_TOKEN"), null,
maxObject, maxTotal, "true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS"))); maxObject, maxTotal, "true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS")));
multipart = new UnavailableMultipart(); multipart = (ClusterStore) store;
} else if (mode.equals("disk")) { } else if (mode.equals("disk")) {
DiskStore disk = new DiskStore(Path.of(env.getOrDefault("DATA_DIR", "/data")), maxObject, maxTotal); DiskStore disk = new DiskStore(Path.of(env.getOrDefault("DATA_DIR", "/data")), maxObject, maxTotal);
store = disk; store = disk;
@@ -6,12 +6,17 @@ import java.util.List;
interface MultipartStorage { interface MultipartStorage {
record Part(int number, String etag) {} record Part(int number, String etag) {}
record PartInfo(int number, long length, String etag, long modified) {}
record PartPage(List<PartInfo> parts, int nextMarker, boolean truncated) {}
record UploadInfo(String id, String key, long created) {}
String create(String bucket, String key, String contentType) throws IOException; String create(String bucket, String key, String contentType) throws IOException;
String putPart(String id, String bucket, String key, int number, InputStream input, String putPart(String id, String bucket, String key, int number, InputStream input,
long length, String expectedHash, String checksum) throws IOException; long length, String expectedHash, String checksum) throws IOException;
ObjectStorage.Metadata complete(String id, String bucket, String key, List<Part> parts) throws IOException; ObjectStorage.Metadata complete(String id, String bucket, String key, List<Part> parts) throws IOException;
void abort(String id, String bucket, String key) 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<UploadInfo> listUploads(String bucket, String prefix) throws IOException;
int activeUploads(); int activeUploads();
long stagedBytes(); long stagedBytes();
} }
@@ -164,6 +164,48 @@ final class MultipartStore implements MultipartStorage {
remove(upload(id, bucket, key)); 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<PartInfo> 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<UploadInfo> listUploads(String bucket, String prefix) throws IOException {
List<UploadInfo> 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 { 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"); if (!id.matches("[0-9a-f-]{36}")) throw new StoreException(404, "NoSuchUpload", "Upload not found");
Path dir = root.resolve(id); Path dir = root.resolve(id);
+7 -1
View File
@@ -21,7 +21,7 @@ final class SchemaMigrator {
result.next(); result.next();
version = result.getInt(1); 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) { 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_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))"); 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("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)"); 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)"); 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")) { try (PreparedStatement insert = connection.prepareStatement("INSERT INTO cluster_usage VALUES (?, 0) ON CONFLICT DO NOTHING")) {
@@ -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<Part> 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; }
}
@@ -3,7 +3,10 @@ package cloud.lunarsky.store;
import java.io.ByteArrayInputStream; import java.io.ByteArrayInputStream;
import java.net.URI; import java.net.URI;
import java.nio.charset.StandardCharsets; import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.util.Arrays; import java.util.Arrays;
import java.util.HexFormat;
import java.util.List;
import java.util.Map; import java.util.Map;
public final class ClusterIntegrationTest { public final class ClusterIntegrationTest {
@@ -62,6 +65,60 @@ public final class ClusterIntegrationTest {
} }
System.out.println("Cluster degraded test passed"); 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" -> { case "quorum-lost" -> {
require(!store.ready(), "One available node must not be ready"); require(!store.ready(), "One available node must not be ready");
try { try {
@@ -151,6 +208,11 @@ public final class ClusterIntegrationTest {
store.put(bucket, key, new ByteArrayInputStream(data), data.length, SigV4.hex(SigV4.hash(data)), store.put(bucket, key, new ByteArrayInputStream(data), data.length, SigV4.hex(SigV4.hash(data)),
null, createOnly, "application/octet-stream"); 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) { private static void require(boolean condition, String message) {
if (!condition) throw new AssertionError(message); if (!condition) throw new AssertionError(message);
} }
+12
View File
@@ -144,6 +144,16 @@ public final class HttpTest {
"?partNumber=2&uploadId=" + upload), "PUT", second, Map.of()), HttpResponse.BodyHandlers.ofByteArray()); "?partNumber=2&uploadId=" + upload), "PUT", second, Map.of()), HttpResponse.BodyHandlers.ofByteArray());
status(200, partOne); status(200, partOne);
status(200, partTwo); 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("<IsTruncated>true</IsTruncated>") ||
!parts.body().contains("<PartNumber>1</PartNumber>"))
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("<UploadId>" + upload + "</UploadId>"))
throw new AssertionError("Multipart upload listing failed: " + uploads.body());
String completion = "<CompleteMultipartUpload><Part><PartNumber>1</PartNumber><ETag>" + String completion = "<CompleteMultipartUpload><Part><PartNumber>1</PartNumber><ETag>" +
partOne.headers().firstValue("etag").orElseThrow() + partOne.headers().firstValue("etag").orElseThrow() +
"</ETag></Part><Part><PartNumber>2</PartNumber><ETag>" + "</ETag></Part><Part><PartNumber>2</PartNumber><ETag>" +
@@ -156,6 +166,8 @@ public final class HttpTest {
if (!"hello world".equals(new String(assembled.body(), StandardCharsets.UTF_8)) || if (!"hello world".equals(new String(assembled.body(), StandardCharsets.UTF_8)) ||
!"video/mp4".equals(assembled.headers().firstValue("content-type").orElse(""))) !"video/mp4".equals(assembled.headers().firstValue("content-type").orElse("")))
throw new AssertionError("Completed multipart object mismatch"); 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="), var abandoned = client.send(signedUri(URI.create(base + "/objects/abandoned?uploads="),
"POST", new byte[0], Map.of()), HttpResponse.BodyHandlers.ofString()); "POST", new byte[0], Map.of()), HttpResponse.BodyHandlers.ofString());
String abandonedId = abandoned.body().split("<UploadId>")[1].split("</UploadId>")[0]; String abandonedId = abandoned.body().split("<UploadId>")[1].split("</UploadId>")[0];