4 Commits
Author SHA1 Message Date
admin 01c82a8223 Release ObjectStore 0.0.3 2026-10-10 08:21:04 +02:00
admin 350611d562 Handle unavailable cluster nodes explicitly 2026-10-10 08:19:32 +02:00
admin 2276fee60c Refactor complex ObjectStore methods 2026-10-10 07:51:38 +02:00
admin b6b8e0adfb Restrict cluster HTTP test to local gateway 2026-10-09 22:01:39 +02:00
12 changed files with 377 additions and 171 deletions

No files matched your search

+1 -1
View File
@@ -8,7 +8,7 @@ Source: [GitHub](https://github.com/LunarSkyOSS/ObjectStore) · [Gitea mirror](h
**Development software. Do not use it for production data.** It is provided as is, without warranty, under the [MIT License](LICENSE). **Development software. Do not use it for production data.** It is provided as is, without warranty, under the [MIT License](LICENSE).
[![Stage: development](assets/badge-stage.svg)](#limits-and-safety) [![License: MIT](assets/badge-license.svg)](LICENSE) [![Runtime: Java 21](assets/badge-java.svg)](Dockerfile) [![S3 API: partial support](assets/badge-api.svg)](#s3-api-support-checklist) [![Stage: development](assets/badge-stage.svg)](#limits-and-safety) [![License: MIT](assets/badge-license.svg)](LICENSE) [![Runtime: Java 21](assets/badge-java.svg)](Dockerfile) [![S3 API: partial support](assets/badge-api.svg)](#s3-api-support-checklist) [![CodeFactor](https://www.codefactor.io/repository/github/lunarskyoss/objectstore/badge)](https://www.codefactor.io/repository/github/lunarskyoss/objectstore)
## Contents ## Contents
+9 -9
View File
@@ -2,11 +2,10 @@
import datetime import datetime
import hashlib import hashlib
import hmac import hmac
import http.client
import pathlib import pathlib
import sys import sys
import urllib.error
import urllib.parse import urllib.parse
import urllib.request
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()
@@ -14,7 +13,9 @@ values = dict(line.strip().split("=", 1) for line in pathlib.Path(sys.argv[1]).r
access = values["S3_ACCESS_KEY"] access = values["S3_ACCESS_KEY"]
secret = values["S3_SECRET_KEY"] secret = values["S3_SECRET_KEY"]
bucket = values.get("S3_BUCKET", "objects") bucket = values.get("S3_BUCKET", "objects")
port = values.get("CLUSTER_HOST_PORT", "9001") port = int(values.get("CLUSTER_HOST_PORT", "9001"))
if not 1 <= port <= 65535:
raise ValueError("CLUSTER_HOST_PORT must be between 1 and 65535")
host = f"127.0.0.1:{port}" host = f"127.0.0.1:{port}"
@@ -40,14 +41,13 @@ def request(method, path, body=b"", extra=None):
signature = hmac.new(key, to_sign.encode(), hashlib.sha256).hexdigest() signature = hmac.new(key, to_sign.encode(), hashlib.sha256).hexdigest()
headers["authorization"] = (f"AWS4-HMAC-SHA256 Credential={access}/{scope}," headers["authorization"] = (f"AWS4-HMAC-SHA256 Credential={access}/{scope},"
f"SignedHeaders={signed_names},Signature={signature}") f"SignedHeaders={signed_names},Signature={signature}")
url = f"http://{host}{path}" connection = http.client.HTTPConnection("127.0.0.1", port, timeout=30)
outgoing = urllib.request.Request(url, data=body if method == "PUT" else None,
method=method, headers=headers)
try: try:
with urllib.request.urlopen(outgoing, timeout=30) as response: connection.request(method, path, body=body if method == "PUT" else None, headers=headers)
response = connection.getresponse()
return response.status, response.read(), response.headers return response.status, response.read(), response.headers
except urllib.error.HTTPError as error: finally:
return error.code, error.read(), error.headers connection.close()
if len(sys.argv) > 2 and sys.argv[2] == "survivor": if len(sys.argv) > 2 and sys.argv[2] == "survivor":
+46 -30
View File
@@ -78,40 +78,14 @@ public final class ClusterNode implements AutoCloseable {
respond(exchange, 200, "ok"); respond(exchange, 200, "ok");
return; return;
} }
byte[] supplied = exchange.getRequestHeaders().getFirst("X-Cluster-Token") == null if (!authorized(exchange)) return;
? new byte[0] : exchange.getRequestHeaders().getFirst("X-Cluster-Token")
.getBytes(java.nio.charset.StandardCharsets.UTF_8);
if (!MessageDigest.isEqual(token, supplied)) {
respond(exchange, 403, "Forbidden");
return;
}
if (path.equals("/identity") && exchange.getRequestMethod().equals("GET")) { if (path.equals("/identity") && exchange.getRequestMethod().equals("GET")) {
respond(exchange, 200, identity.nodeId() + " " + identity.hostId()); respond(exchange, 200, identity.nodeId() + " " + identity.hostId());
return; return;
} }
if (!identity.nodeId().toString().equals(exchange.getRequestHeaders().getFirst("X-Cluster-Expected-Node"))) { if (!expectedNodeAndRepairAuthorized(exchange)) return;
respond(exchange, 409, "Wrong storage node"); String id = segmentId(exchange, path);
return; if (id == null) return;
}
if (exchange.getRequestMethod().equals("PUT") &&
"true".equals(exchange.getRequestHeaders().getFirst("X-Cluster-Repair"))) {
String suppliedRepair = exchange.getRequestHeaders().getFirst("X-Cluster-Repair-Token");
byte[] suppliedBytes = suppliedRepair == null ? new byte[0]
: suppliedRepair.getBytes(java.nio.charset.StandardCharsets.UTF_8);
if (!MessageDigest.isEqual(repairToken, suppliedBytes)) {
respond(exchange, 403, "Repair authority required");
return;
}
}
if (!path.matches("/segments/[0-9a-f-]{36}")) {
respond(exchange, 404, "Not found");
return;
}
String id = path.substring("/segments/".length());
if (!UUID.fromString(id).toString().equals(id)) {
respond(exchange, 400, "Invalid segment ID");
return;
}
switch (exchange.getRequestMethod()) { switch (exchange.getRequestMethod()) {
case "PUT" -> put(exchange, segmentPath(id, true)); case "PUT" -> put(exchange, segmentPath(id, true));
case "GET" -> get(exchange, segmentPath(id, false)); case "GET" -> get(exchange, segmentPath(id, false));
@@ -127,6 +101,48 @@ public final class ClusterNode implements AutoCloseable {
} }
} }
private boolean authorized(HttpExchange exchange) throws IOException {
byte[] supplied = exchange.getRequestHeaders().getFirst("X-Cluster-Token") == null
? new byte[0] : exchange.getRequestHeaders().getFirst("X-Cluster-Token")
.getBytes(java.nio.charset.StandardCharsets.UTF_8);
if (!MessageDigest.isEqual(token, supplied)) {
respond(exchange, 403, "Forbidden");
return false;
}
return true;
}
private boolean expectedNodeAndRepairAuthorized(HttpExchange exchange) throws IOException {
if (!identity.nodeId().toString().equals(exchange.getRequestHeaders().getFirst("X-Cluster-Expected-Node"))) {
respond(exchange, 409, "Wrong storage node");
return false;
}
if (exchange.getRequestMethod().equals("PUT") &&
"true".equals(exchange.getRequestHeaders().getFirst("X-Cluster-Repair"))) {
String suppliedRepair = exchange.getRequestHeaders().getFirst("X-Cluster-Repair-Token");
byte[] suppliedBytes = suppliedRepair == null ? new byte[0]
: suppliedRepair.getBytes(java.nio.charset.StandardCharsets.UTF_8);
if (!MessageDigest.isEqual(repairToken, suppliedBytes)) {
respond(exchange, 403, "Repair authority required");
return false;
}
}
return true;
}
private static String segmentId(HttpExchange exchange, String path) throws IOException {
if (!path.matches("/segments/[0-9a-f-]{36}")) {
respond(exchange, 404, "Not found");
return null;
}
String id = path.substring("/segments/".length());
if (!UUID.fromString(id).toString().equals(id)) {
respond(exchange, 400, "Invalid segment ID");
return null;
}
return id;
}
private synchronized Path segmentPath(String id, boolean createShard) throws IOException { private synchronized Path segmentPath(String id, boolean createShard) throws IOException {
Path shard = segments.resolve(id.substring(0, 2)); Path shard = segments.resolve(id.substring(0, 2));
if (createShard && !Files.isDirectory(shard)) { if (createShard && !Files.isDirectory(shard)) {
+77 -30
View File
@@ -52,6 +52,23 @@ final class ClusterStore implements ObjectStorage {
@Override public Metadata put(String bucket, String key, InputStream input, long length, String expectedHash, @Override public Metadata put(String bucket, String key, InputStream input, long length, String expectedHash,
String checksum, boolean createOnly, String contentType) throws IOException { String checksum, boolean createOnly, String contentType) throws IOException {
validatePut(bucket, length, contentType);
MessageDigest md5 = digest("MD5");
Path staged = Files.createTempFile("objectstore-cluster-", ".pending");
List<Segment> segments;
byte[] fullHash;
try {
fullHash = stageInput(staged, input, length, expectedHash, checksum, md5);
checkCapacity(bucket, key, length, createOnly);
segments = uploadSegments(staged, length);
} finally { Files.deleteIfExists(staged); }
Metadata metadata = new Metadata(length, Instant.now().toEpochMilli(),
HexFormat.of().formatHex(md5.digest()), fullHash, bucket, key, contentType);
persistObject(metadata, segments, createOnly);
return metadata;
}
private void validatePut(String bucket, long length, String contentType) {
if (!configuredBucket.equals(bucket)) throw new StoreException(404, "NoSuchBucket", "Bucket not found"); if (!configuredBucket.equals(bucket)) throw new StoreException(404, "NoSuchBucket", "Bucket not found");
if (length < 0) throw new StoreException(411, "MissingContentLength", "Content-Length is required"); if (length < 0) throw new StoreException(411, "MissingContentLength", "Content-Length is required");
if (length > maxObject) throw new StoreException(413, "EntityTooLarge", "Object exceeds the configured size limit"); if (length > maxObject) throw new StoreException(413, "EntityTooLarge", "Object exceeds the configured size limit");
@@ -59,11 +76,11 @@ final class ClusterStore implements ObjectStorage {
throw new StoreException(503, "SlowDown", "Fewer than two storage hosts are available"); throw new StoreException(503, "SlowDown", "Fewer than two storage hosts are available");
if (contentType.getBytes(java.nio.charset.StandardCharsets.UTF_8).length > 255) if (contentType.getBytes(java.nio.charset.StandardCharsets.UTF_8).length > 255)
throw new StoreException(400, "InvalidArgument", "Content-Type is too long"); throw new StoreException(400, "InvalidArgument", "Content-Type is too long");
MessageDigest sha = digest("SHA-256"), md5 = digest("MD5"); }
List<Segment> segments = new ArrayList<>();
byte[] fullHash; private byte[] stageInput(Path staged, InputStream input, long length, String expectedHash,
Path staged = Files.createTempFile("objectstore-cluster-", ".pending"); String checksum, MessageDigest md5) throws IOException {
try { MessageDigest sha = digest("SHA-256");
try (OutputStream output = Files.newOutputStream(staged)) { try (OutputStream output = Files.newOutputStream(staged)) {
byte[] buffer = new byte[65536]; byte[] buffer = new byte[65536];
long remaining = length; long remaining = length;
@@ -77,11 +94,15 @@ final class ClusterStore implements ObjectStorage {
} }
} }
if (input.read() != -1) throw new StoreException(413, "EntityTooLarge", "Payload exceeds declared size"); if (input.read() != -1) throw new StoreException(413, "EntityTooLarge", "Payload exceeds declared size");
fullHash = sha.digest(); byte[] fullHash = sha.digest();
if (!HexFormat.of().formatHex(fullHash).equals(expectedHash)) if (!HexFormat.of().formatHex(fullHash).equals(expectedHash))
throw new StoreException(400, "XAmzContentSHA256Mismatch", "Payload hash mismatch"); throw new StoreException(400, "XAmzContentSHA256Mismatch", "Payload hash mismatch");
if (checksum != null && !Base64.getEncoder().encodeToString(fullHash).equals(checksum)) if (checksum != null && !Base64.getEncoder().encodeToString(fullHash).equals(checksum))
throw new StoreException(400, "BadDigest", "SHA-256 checksum mismatch"); throw new StoreException(400, "BadDigest", "SHA-256 checksum mismatch");
return fullHash;
}
private void checkCapacity(String bucket, String key, long length, boolean createOnly) throws IOException {
try (Connection connection = connect()) { try (Connection connection = connect()) {
long previous = currentLength(connection, bucket, key); long previous = currentLength(connection, bucket, key);
if (createOnly && previous >= 0) if (createOnly && previous >= 0)
@@ -95,6 +116,10 @@ final class ClusterStore implements ObjectStorage {
} }
} }
} catch (SQLException error) { throw databaseError(error); } } catch (SQLException error) { throw databaseError(error); }
}
private List<Segment> uploadSegments(Path staged, long length) throws IOException {
List<Segment> segments = new ArrayList<>();
try (InputStream stagedInput = Files.newInputStream(staged)) { try (InputStream stagedInput = Files.newInputStream(staged)) {
long remaining = length; long remaining = length;
while (remaining > 0) { while (remaining > 0) {
@@ -124,9 +149,12 @@ final class ClusterStore implements ObjectStorage {
remaining -= wanted; remaining -= wanted;
} }
} }
} finally { Files.deleteIfExists(staged); } return segments;
Metadata metadata = new Metadata(length, Instant.now().toEpochMilli(), }
HexFormat.of().formatHex(md5.digest()), fullHash, bucket, key, contentType);
private void persistObject(Metadata metadata, List<Segment> segments, boolean createOnly) throws IOException {
String bucket = metadata.bucket(), key = metadata.key();
long length = metadata.length();
UUID generation = UUID.randomUUID(); UUID generation = UUID.randomUUID();
try (Connection connection = connect()) { try (Connection connection = connect()) {
connection.setAutoCommit(false); connection.setAutoCommit(false);
@@ -159,7 +187,6 @@ final class ClusterStore implements ObjectStorage {
delete.setString(1, bucket); delete.setString(2, key); delete.executeUpdate(); delete.setString(1, bucket); delete.setString(2, key); delete.executeUpdate();
} }
connection.commit(); connection.commit();
return metadata;
} catch (SQLException | RuntimeException error) { } catch (SQLException | RuntimeException error) {
connection.rollback(); connection.rollback();
if (error instanceof SQLException sql) throw databaseError(sql); if (error instanceof SQLException sql) throw databaseError(sql);
@@ -241,28 +268,35 @@ final class ClusterStore implements ObjectStorage {
} }
@Override public ListPage list(String bucket, String prefix, String delimiter, int maxKeys, String after) throws IOException { @Override public ListPage list(String bucket, String prefix, String delimiter, int maxKeys, String after) throws IOException {
List<ListedObject> entries = new ArrayList<>(); if (maxKeys == 0) return new ListPage(new ArrayList<>(), new ArrayList<>(), null, false);
List<String> prefixes = new ArrayList<>();
if (maxKeys == 0) return new ListPage(entries, prefixes, null, false);
String lastKey = null, activePrefix = null;
boolean truncated = false;
try (Connection connection = connect()) { try (Connection connection = connect()) {
connection.setAutoCommit(false); connection.setAutoCommit(false);
ListPage page;
try (PreparedStatement query = connection.prepareStatement( try (PreparedStatement query = connection.prepareStatement(
"SELECT object_key, length, modified, etag, sha256, content_type FROM cluster_objects WHERE bucket=? AND object_key>=? ORDER BY object_key")) { "SELECT object_key, length, modified, etag, sha256, content_type FROM cluster_objects WHERE bucket=? AND object_key>=? ORDER BY object_key")) {
query.setString(1, bucket); query.setString(1, bucket);
query.setString(2, after != null && after.compareTo(prefix) > 0 ? after : prefix); query.setString(2, after != null && after.compareTo(prefix) > 0 ? after : prefix);
query.setFetchSize(128); query.setFetchSize(128);
try (ResultSet result = query.executeQuery()) { try (ResultSet result = query.executeQuery()) {
page = readListPage(result, bucket, prefix, delimiter, maxKeys, after);
}
}
connection.commit();
return page;
} catch (SQLException error) { throw databaseError(error); }
}
private static ListPage readListPage(ResultSet result, String bucket, String prefix, String delimiter,
int maxKeys, String after) throws SQLException {
List<ListedObject> entries = new ArrayList<>();
List<String> prefixes = new ArrayList<>();
String lastKey = null, activePrefix = null;
boolean truncated = false;
while (result.next()) { while (result.next()) {
String key = result.getString(1); String key = result.getString(1);
if (!key.startsWith(prefix)) break; if (!key.startsWith(prefix)) break;
if (after != null && key.compareTo(after) <= 0) continue; if (after != null && key.compareTo(after) <= 0) continue;
String group = null; String group = commonPrefix(key, prefix, delimiter);
if (!delimiter.isEmpty()) {
int at = key.indexOf(delimiter, prefix.length());
if (at >= 0) group = key.substring(0, at + delimiter.length());
}
if (group != null && group.equals(activePrefix)) { lastKey = key; continue; } if (group != null && group.equals(activePrefix)) { lastKey = key; continue; }
if (entries.size() + prefixes.size() >= maxKeys) { truncated = true; break; } if (entries.size() + prefixes.size() >= maxKeys) { truncated = true; break; }
if (group != null) { prefixes.add(group); activePrefix = group; } if (group != null) { prefixes.add(group); activePrefix = group; }
@@ -273,13 +307,15 @@ final class ClusterStore implements ObjectStorage {
} }
lastKey = key; lastKey = key;
} }
}
}
connection.commit();
} catch (SQLException error) { throw databaseError(error); }
return new ListPage(entries, prefixes, truncated ? lastKey : null, truncated); return new ListPage(entries, prefixes, truncated ? lastKey : null, truncated);
} }
private static String commonPrefix(String key, String prefix, String delimiter) {
if (delimiter.isEmpty()) return null;
int at = key.indexOf(delimiter, prefix.length());
return at < 0 ? null : key.substring(0, at + delimiter.length());
}
private long lockUsage(Connection connection, String bucket) throws SQLException { private long lockUsage(Connection connection, String bucket) throws SQLException {
try (PreparedStatement query = connection.prepareStatement("SELECT used_bytes FROM cluster_usage WHERE bucket=? FOR UPDATE")) { try (PreparedStatement query = connection.prepareStatement("SELECT used_bytes FROM cluster_usage WHERE bucket=? FOR UPDATE")) {
query.setString(1, bucket); query.setString(1, bucket);
@@ -364,23 +400,21 @@ final class ClusterStore implements ObjectStorage {
for (UUID id : segment.replicas()) { for (UUID id : segment.replicas()) {
int node = nodes.index(id); int node = nodes.index(id);
if (node < 0) continue; if (node < 0) continue;
try { byte[] candidate = readableReplica(node, segment);
byte[] candidate = nodes.get(node, segment.id(), segment.length(), segment.hash()); if (candidate == null) continue;
if (copy == null) copy = candidate; if (copy == null) copy = candidate;
healthy.add(id); healthy.add(id);
healthyHosts.add(nodes.faultDomain(node, testNodeDomains)); healthyHosts.add(nodes.faultDomain(node, testNodeDomains));
} catch (IOException error) { }
} }
if (copy == null) { unrecoverable++; continue; } if (copy == null) { unrecoverable++; continue; }
for (int node : PlacementPolicy.candidates(segment.id(), nodes, testNodeDomains)) { for (int node : PlacementPolicy.candidates(segment.id(), nodes, testNodeDomains)) {
UUID host = nodes.faultDomain(node, testNodeDomains); UUID host = nodes.faultDomain(node, testNodeDomains);
if (healthyHosts.contains(host)) continue; if (healthyHosts.contains(host)) continue;
try { if (repairReplica(node, segment, copy)) {
nodes.repair(node, segment.id(), copy, segment.hash());
healthy.add(nodes.node(node).id()); healthy.add(nodes.node(node).id());
healthyHosts.add(host); healthyHosts.add(host);
restored++; restored++;
} catch (IOException error) { } }
if (healthyHosts.size() == 3) break; if (healthyHosts.size() == 3) break;
} }
if (healthyHosts.size() < 3) underReplicated++; if (healthyHosts.size() < 3) underReplicated++;
@@ -403,6 +437,19 @@ final class ClusterStore implements ObjectStorage {
} catch (SQLException error) { throw databaseError(error); } } catch (SQLException error) { throw databaseError(error); }
return new RepairReport(scanned, restored, underReplicated, unrecoverable); return new RepairReport(scanned, restored, underReplicated, unrecoverable);
} }
private byte[] readableReplica(int node, Segment segment) {
try { return nodes.get(node, segment.id(), segment.length(), segment.hash()); }
catch (IOException unavailable) { return null; }
}
private boolean repairReplica(int node, Segment segment, byte[] copy) {
try {
nodes.repair(node, segment.id(), copy, segment.hash());
return true;
} catch (IOException unavailable) { return false; }
}
@Override public void close() {} @Override public void close() {}
private final class SegmentStream extends InputStream { private final class SegmentStream extends InputStream {
+24 -9
View File
@@ -117,11 +117,25 @@ final class DiskStore implements ObjectStorage {
byte[] bucketBytes = bucket.getBytes(StandardCharsets.UTF_8); byte[] bucketBytes = bucket.getBytes(StandardCharsets.UTF_8);
byte[] keyBytes = key.getBytes(StandardCharsets.UTF_8); byte[] keyBytes = key.getBytes(StandardCharsets.UTF_8);
byte[] typeBytes = contentType.getBytes(StandardCharsets.UTF_8); byte[] typeBytes = contentType.getBytes(StandardCharsets.UTF_8);
if (bucketBytes.length > 63 || keyBytes.length > 1024 || typeBytes.length > 255) validateMetadataLengths(bucketBytes, keyBytes, typeBytes);
throw new StoreException(400, "InvalidArgument", "Object metadata is too long");
int headerLength = HEADER_V2 + bucketBytes.length + keyBytes.length + typeBytes.length;
Path destination = object(bucket, key), pending = Files.createTempFile(temporary, "upload-", ".part"); Path destination = object(bucket, key), pending = Files.createTempFile(temporary, "upload-", ".part");
try { try {
Metadata metadata = stagePut(pending, input, length, expectedHash, checksum,
bucket, key, contentType, bucketBytes, keyBytes, typeBytes);
installPending(destination, pending, metadata, createOnly);
return metadata;
} finally { Files.deleteIfExists(pending); }
}
private static void validateMetadataLengths(byte[] bucketBytes, byte[] keyBytes, byte[] typeBytes) {
if (bucketBytes.length > 63 || keyBytes.length > 1024 || typeBytes.length > 255)
throw new StoreException(400, "InvalidArgument", "Object metadata is too long");
}
private Metadata stagePut(Path pending, InputStream input, long length, String expectedHash, String checksum,
String bucket, String key, String contentType,
byte[] bucketBytes, byte[] keyBytes, byte[] typeBytes) throws IOException {
int headerLength = HEADER_V2 + bucketBytes.length + keyBytes.length + typeBytes.length;
MessageDigest sha = digest("SHA-256"), md5 = digest("MD5"); MessageDigest sha = digest("SHA-256"), md5 = digest("MD5");
long count = 0; long count = 0;
try (OutputStream out = Files.newOutputStream(pending)) { try (OutputStream out = Files.newOutputStream(pending)) {
@@ -150,7 +164,10 @@ final class DiskStore implements ObjectStorage {
while (header.hasRemaining()) file.write(header, header.position()); while (header.hasRemaining()) file.write(header, header.position());
file.force(true); file.force(true);
} }
Metadata metadata = new Metadata(count, modified, SigV4.hex(etag), hash, bucket, key, contentType); return new Metadata(count, modified, SigV4.hex(etag), hash, bucket, key, contentType);
}
private void installPending(Path destination, Path pending, Metadata metadata, boolean createOnly) throws IOException {
synchronized (lock(destination)) { synchronized (lock(destination)) {
long previous = 0; long previous = 0;
boolean existed = Files.exists(destination); boolean existed = Files.exists(destination);
@@ -164,18 +181,16 @@ final class DiskStore implements ObjectStorage {
} }
} }
synchronized (this) { synchronized (this) {
if (used - previous + count > maxTotal) if (used - previous + metadata.length() > maxTotal)
throw new StoreException(507, "InsufficientStorage", "Store capacity limit reached"); throw new StoreException(507, "InsufficientStorage", "Store capacity limit reached");
Files.move(pending, destination, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING); Files.move(pending, destination, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING);
used = used - previous + count; used = used - previous + metadata.length();
if (!existed) objectCount++; if (!existed) objectCount++;
if (legacy) legacyCount--; if (legacy) legacyCount--;
index.put(indexKey(bucket, key), metadata); index.put(indexKey(metadata.bucket(), metadata.key()), metadata);
syncDirectory(destination.getParent()); syncDirectory(destination.getParent());
} }
} }
return metadata;
} finally { Files.deleteIfExists(pending); }
} }
public OpenObject open(String bucket, String key) throws IOException { public OpenObject open(String bucket, String key) throws IOException {
+77 -40
View File
@@ -45,25 +45,42 @@ public final class Main {
exchange.getResponseHeaders().set("X-Content-Type-Options", "nosniff"); exchange.getResponseHeaders().set("X-Content-Type-Options", "nosniff");
try { try {
if (!admitted) throw new StoreException(503, "SlowDown", "Too many concurrent requests"); if (!admitted) throw new StoreException(503, "SlowDown", "Too many concurrent requests");
if (exchange.getRequestURI().getRawPath().equals("/health") && exchange.getRequestMethod().equals("GET")) { if (handleStatus(exchange)) return;
String hash = authentication.verify(exchange.getRequestMethod(), exchange.getRequestURI(), exchange.getRequestHeaders());
String path = SigV4.decode(exchange.getRequestURI().getRawPath());
Map<String, String> query = query(exchange.getRequestURI().getRawQuery());
if (path.equals("/" + bucket) || path.equals("/" + bucket + "/")) {
handleBucket(exchange, query, hash);
} else {
handleObject(exchange, path, query, hash);
}
} catch (StoreException error) { sendError(exchange, error.status, error.code, error.getMessage(), requestId); }
catch (Exception error) {
System.err.println("ObjectStore request failed: " + requestId + " " + error.getClass().getSimpleName());
sendError(exchange, 500, "InternalError", "Storage operation failed", requestId);
} finally { if (admitted) slots.release(); exchange.close(); }
}
private boolean handleStatus(HttpExchange exchange) throws IOException {
if (!exchange.getRequestMethod().equals("GET")) return false;
String path = exchange.getRequestURI().getRawPath();
if (path.equals("/health")) {
byte[] body = "{\"status\":\"ok\",\"service\":\"lunarsky-objectstore\"}".getBytes(StandardCharsets.UTF_8); byte[] body = "{\"status\":\"ok\",\"service\":\"lunarsky-objectstore\"}".getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type", "application/json"); exchange.getResponseHeaders().set("Content-Type", "application/json");
exchange.sendResponseHeaders(200, body.length); exchange.sendResponseHeaders(200, body.length);
exchange.getResponseBody().write(body); exchange.getResponseBody().write(body);
return; return true;
} }
if (exchange.getRequestURI().getRawPath().equals("/ready") && exchange.getRequestMethod().equals("GET")) { if (!path.equals("/ready")) return false;
boolean ready = store.ready(); boolean ready = store.ready();
byte[] body = (ready ? "ready" : "unavailable").getBytes(StandardCharsets.UTF_8); byte[] body = (ready ? "ready" : "unavailable").getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type", "text/plain; charset=utf-8"); exchange.getResponseHeaders().set("Content-Type", "text/plain; charset=utf-8");
exchange.sendResponseHeaders(ready ? 200 : 503, body.length); exchange.sendResponseHeaders(ready ? 200 : 503, body.length);
exchange.getResponseBody().write(body); exchange.getResponseBody().write(body);
return; return true;
} }
String hash = authentication.verify(exchange.getRequestMethod(), exchange.getRequestURI(), exchange.getRequestHeaders());
String path = SigV4.decode(exchange.getRequestURI().getRawPath()); private void handleBucket(HttpExchange exchange, Map<String, String> query, String hash) throws IOException {
Map<String, String> query = query(exchange.getRequestURI().getRawQuery());
if (path.equals("/" + bucket) || path.equals("/" + bucket + "/")) {
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) ||
@@ -71,8 +88,9 @@ public final class Main {
unsupported("Bucket operation"); unsupported("Bucket operation");
requireEmptyBody(exchange, hash); requireEmptyBody(exchange, hash);
listObjects(exchange, query); listObjects(exchange, query);
return;
} }
private void handleObject(HttpExchange exchange, String path, Map<String, String> query, String hash) throws IOException {
String prefix = "/" + bucket + "/"; String prefix = "/" + bucket + "/";
if (!path.startsWith(prefix)) throw new StoreException(404, "NoSuchBucket", "Bucket not found"); if (!path.startsWith(prefix)) throw new StoreException(404, "NoSuchBucket", "Bucket not found");
String key = path.substring(prefix.length()); String key = path.substring(prefix.length());
@@ -84,7 +102,21 @@ public final class Main {
("PutObject".equals(query.get("x-id")) || "GetObject".equals(query.get("x-id")) || ("PutObject".equals(query.get("x-id")) || "GetObject".equals(query.get("x-id")) ||
"HeadObject".equals(query.get("x-id")) || "DeleteObject".equals(query.get("x-id"))))) "HeadObject".equals(query.get("x-id")) || "DeleteObject".equals(query.get("x-id")))))
unsupported("Query operation"); unsupported("Query operation");
var headers = exchange.getRequestHeaders(); validateObjectHeaders(exchange.getRequestHeaders());
if (multipartRequest) {
handleMultipart(exchange, method, query, key, hash);
return;
}
if (!method.equals("PUT")) requireEmptyBody(exchange, hash);
switch (method) {
case "PUT" -> putObject(exchange, key, hash);
case "GET", "HEAD" -> readObject(exchange, key);
case "DELETE" -> deleteObject(exchange, key);
default -> unsupported("HTTP method");
}
}
private static void validateObjectHeaders(com.sun.net.httpserver.Headers headers) {
for (String name : headers.keySet()) { for (String name : headers.keySet()) {
String lower = name.toLowerCase(java.util.Locale.ROOT); String lower = name.toLowerCase(java.util.Locale.ROOT);
if (lower.startsWith("x-amz-") && !java.util.Set.of("x-amz-date", "x-amz-content-sha256", if (lower.startsWith("x-amz-") && !java.util.Set.of("x-amz-date", "x-amz-content-sha256",
@@ -99,40 +131,29 @@ public final class Main {
} }
String algorithm = SigV4.single(headers, "x-amz-sdk-checksum-algorithm"); String algorithm = SigV4.single(headers, "x-amz-sdk-checksum-algorithm");
if (algorithm != null && !algorithm.equals("SHA256")) unsupported("Checksum algorithm"); if (algorithm != null && !algorithm.equals("SHA256")) unsupported("Checksum algorithm");
if (multipartRequest) {
handleMultipart(exchange, method, query, key, hash);
return;
} }
if (!method.equals("PUT")) requireEmptyBody(exchange, hash);
switch (method) { private void putObject(HttpExchange exchange, String key, String hash) throws IOException {
case "PUT" -> { var headers = exchange.getRequestHeaders();
String length = SigV4.single(headers, "content-length"), condition = SigV4.single(headers, "if-none-match"); String length = SigV4.single(headers, "content-length"), condition = SigV4.single(headers, "if-none-match");
if (condition != null && !condition.equals("*")) unsupported("Write condition"); if (condition != null && !condition.equals("*")) unsupported("Write condition");
long bytes; long bytes;
try { bytes = length == null ? -1 : Long.parseLong(length); } try { bytes = length == null ? -1 : Long.parseLong(length); }
catch (NumberFormatException e) { throw new StoreException(400, "InvalidArgument", "Invalid Content-Length"); } catch (NumberFormatException e) { throw new StoreException(400, "InvalidArgument", "Invalid Content-Length"); }
if (headers.containsKey("content-encoding")) unsupported("Encoded payload"); if (headers.containsKey("content-encoding")) unsupported("Encoded payload");
String contentType = contentType(headers); String type = contentType(headers);
ObjectStorage.Metadata data = store.put(bucket, key, exchange.getRequestBody(), bytes, hash, ObjectStorage.Metadata data = store.put(bucket, key, exchange.getRequestBody(), bytes, hash,
SigV4.single(headers, "x-amz-checksum-sha256"), condition != null, contentType); SigV4.single(headers, "x-amz-checksum-sha256"), condition != null, type);
exchange.getResponseHeaders().set("ETag", "\"" + data.etag() + "\""); exchange.getResponseHeaders().set("ETag", "\"" + data.etag() + "\"");
exchange.getResponseHeaders().set("x-amz-checksum-sha256", Base64.getEncoder().encodeToString(data.sha256())); exchange.getResponseHeaders().set("x-amz-checksum-sha256", Base64.getEncoder().encodeToString(data.sha256()));
exchange.sendResponseHeaders(200, -1); exchange.sendResponseHeaders(200, -1);
} }
case "GET", "HEAD" -> readObject(exchange, key);
case "DELETE" -> { private void deleteObject(HttpExchange exchange, String key) throws IOException {
if (headers.containsKey("if-none-match")) unsupported("Conditional delete"); if (exchange.getRequestHeaders().containsKey("if-none-match")) unsupported("Conditional delete");
store.delete(bucket, key); store.delete(bucket, key);
exchange.sendResponseHeaders(204, -1); exchange.sendResponseHeaders(204, -1);
} }
default -> unsupported("HTTP method");
}
} catch (StoreException error) { sendError(exchange, error.status, error.code, error.getMessage(), requestId); }
catch (Exception error) {
System.err.println("ObjectStore request failed: " + requestId + " " + error.getClass().getSimpleName());
sendError(exchange, 500, "InternalError", "Storage operation failed", requestId);
} finally { if (admitted) slots.release(); exchange.close(); }
}
private static String contentType(com.sun.net.httpserver.Headers headers) { private static String contentType(com.sun.net.httpserver.Headers headers) {
String value = SigV4.single(headers, "content-type"); String value = SigV4.single(headers, "content-type");
@@ -337,6 +358,24 @@ public final class Main {
} }
private void listObjects(HttpExchange exchange, Map<String, String> query) throws IOException { private void listObjects(HttpExchange exchange, Map<String, String> query) throws IOException {
ListRequest request = listRequest(query);
var page = store.list(bucket, request.prefix(), request.delimiter(), request.maxKeys(), request.after());
StringBuilder xml = new StringBuilder("<?xml version=\"1.0\" encoding=\"UTF-8\"?><ListBucketResult xmlns=\"http://s3.amazonaws.com/doc/2006-03-01/\">");
appendListHeader(xml, query, request, page);
appendListEntries(xml, page, request.encoding());
if (page.truncated()) xml.append("<NextContinuationToken>")
.append(Base64.getUrlEncoder().withoutPadding().encodeToString(page.nextKey().getBytes(StandardCharsets.UTF_8)))
.append("</NextContinuationToken>");
xml.append("</ListBucketResult>");
byte[] body = xml.toString().getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type", "application/xml");
exchange.sendResponseHeaders(200, body.length);
exchange.getResponseBody().write(body);
}
private record ListRequest(String prefix, String delimiter, String encoding, int maxKeys, String after) { }
private static ListRequest listRequest(Map<String, String> query) {
String prefix = query.getOrDefault("prefix", ""), delimiter = query.getOrDefault("delimiter", ""); String prefix = query.getOrDefault("prefix", ""), delimiter = query.getOrDefault("delimiter", "");
String encoding = query.get("encoding-type"); String encoding = query.get("encoding-type");
if (encoding != null && !encoding.equals("url")) unsupported("Encoding type"); if (encoding != null && !encoding.equals("url")) unsupported("Encoding type");
@@ -356,8 +395,11 @@ public final class Main {
throw new StoreException(400, "InvalidArgument", "Invalid continuation token"); throw new StoreException(400, "InvalidArgument", "Invalid continuation token");
} }
} }
var page = store.list(bucket, prefix, delimiter, maxKeys, after); return new ListRequest(prefix, delimiter, encoding, maxKeys, after);
StringBuilder xml = new StringBuilder("<?xml version=\"1.0\" encoding=\"UTF-8\"?><ListBucketResult xmlns=\"http://s3.amazonaws.com/doc/2006-03-01/\">"); }
private void appendListHeader(StringBuilder xml, Map<String, String> query, ListRequest request, ObjectStorage.ListPage page) {
String prefix = request.prefix(), delimiter = request.delimiter(), encoding = request.encoding();
xml.append("<Name>").append(xml(bucket)).append("</Name><Prefix>").append(xml(listKey(prefix, encoding))).append("</Prefix>"); xml.append("<Name>").append(xml(bucket)).append("</Name><Prefix>").append(xml(listKey(prefix, encoding))).append("</Prefix>");
if (!delimiter.isEmpty()) xml.append("<Delimiter>").append(xml(listKey(delimiter, encoding))).append("</Delimiter>"); if (!delimiter.isEmpty()) xml.append("<Delimiter>").append(xml(listKey(delimiter, encoding))).append("</Delimiter>");
if (encoding != null) xml.append("<EncodingType>url</EncodingType>"); if (encoding != null) xml.append("<EncodingType>url</EncodingType>");
@@ -365,8 +407,11 @@ public final class Main {
.append(xml(query.get("continuation-token"))).append("</ContinuationToken>"); .append(xml(query.get("continuation-token"))).append("</ContinuationToken>");
if (query.containsKey("start-after")) xml.append("<StartAfter>") if (query.containsKey("start-after")) xml.append("<StartAfter>")
.append(xml(listKey(query.get("start-after"), encoding))).append("</StartAfter>"); .append(xml(listKey(query.get("start-after"), encoding))).append("</StartAfter>");
xml.append("<KeyCount>").append(page.keyCount()).append("</KeyCount><MaxKeys>").append(maxKeys) xml.append("<KeyCount>").append(page.keyCount()).append("</KeyCount><MaxKeys>").append(request.maxKeys())
.append("</MaxKeys><IsTruncated>").append(page.truncated()).append("</IsTruncated>"); .append("</MaxKeys><IsTruncated>").append(page.truncated()).append("</IsTruncated>");
}
private static void appendListEntries(StringBuilder xml, ObjectStorage.ListPage page, String encoding) {
int objectAt = 0, prefixAt = 0; int objectAt = 0, prefixAt = 0;
while (objectAt < page.objects().size() || prefixAt < page.prefixes().size()) { while (objectAt < page.objects().size() || prefixAt < page.prefixes().size()) {
if (objectAt < page.objects().size() && if (objectAt < page.objects().size() &&
@@ -384,14 +429,6 @@ public final class Main {
.append("</Prefix></CommonPrefixes>"); .append("</Prefix></CommonPrefixes>");
} }
} }
if (page.truncated()) xml.append("<NextContinuationToken>")
.append(Base64.getUrlEncoder().withoutPadding().encodeToString(page.nextKey().getBytes(StandardCharsets.UTF_8)))
.append("</NextContinuationToken>");
xml.append("</ListBucketResult>");
byte[] body = xml.toString().getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type", "application/xml");
exchange.sendResponseHeaders(200, body.length);
exchange.getResponseBody().write(body);
} }
private static String listKey(String key, String encoding) { private static String listKey(String key, String encoding) {
+7 -3
View File
@@ -74,6 +74,11 @@ final class NodeClient {
} }
} }
static NodeIdentity probeIfAvailable(URI url, String token) {
try { return probe(url, token); }
catch (IOException offline) { return null; }
}
static void validateUrl(URI url) { static void validateUrl(URI url) {
if (url == null || !"http".equals(url.getScheme()) || url.getHost() == null || if (url == null || !"http".equals(url.getScheme()) || url.getHost() == null ||
url.getPort() < 1 || url.getRawUserInfo() != null || url.getPort() < 1 || url.getRawUserInfo() != null ||
@@ -86,12 +91,11 @@ final class NodeClient {
Set<UUID> healthy = new HashSet<>(); Set<UUID> healthy = new HashSet<>();
for (int i = 0; i < nodes.size(); i++) { for (int i = 0; i < nodes.size(); i++) {
Node node = nodes.get(i); Node node = nodes.get(i);
try { NodeIdentity actual = probeIfAvailable(node.url(), token);
NodeIdentity actual = probe(node.url(), token); if (actual == null) continue;
if (actual.nodeId().equals(node.id()) && actual.hostId().equals(node.hostId())) if (actual.nodeId().equals(node.id()) && actual.hostId().equals(node.hostId()))
healthy.add(faultDomain(i, testNodeDomains)); healthy.add(faultDomain(i, testNodeDomains));
if (healthy.size() >= required) return true; if (healthy.size() >= required) return true;
} catch (IOException error) { }
} }
return false; return false;
} }
+35 -19
View File
@@ -69,6 +69,25 @@ final class NodeRegistry {
try (Statement statement = connection.createStatement()) { try (Statement statement = connection.createStatement()) {
statement.execute("SELECT pg_advisory_xact_lock(6834071092781)"); statement.execute("SELECT pg_advisory_xact_lock(6834071092781)");
} }
Map<String, NodeClient.Node> stored = registeredNodes(connection);
if (stored.isEmpty()) registerInitialNodes(connection, urls, token, stored);
List<NodeClient.Node> configured = configuredNodes(urls, token, stored);
ensureLiveReplicasConfigured(connection, configured);
NodeClient nodes = new NodeClient(configured, token, repairToken);
connection.commit();
return nodes;
} catch (SQLException | IOException | RuntimeException error) {
try { connection.rollback(); } catch (SQLException rollback) { error.addSuppressed(rollback); }
if (error instanceof IOException io) throw io;
if (error instanceof SQLException sql) throw new IOException("Node registry check failed", sql);
throw (RuntimeException) error;
} finally {
try { connection.setAutoCommit(true); }
catch (SQLException error) { throw new IOException("Could not restore metadata connection", error); }
}
}
private static Map<String, NodeClient.Node> registeredNodes(Connection connection) throws SQLException {
Map<String, NodeClient.Node> stored = new HashMap<>(); Map<String, NodeClient.Node> stored = new HashMap<>();
try (Statement statement = connection.createStatement(); try (Statement statement = connection.createStatement();
ResultSet result = statement.executeQuery("SELECT node_id, host_id, endpoint FROM cluster_nodes WHERE state <> 'retired'")) { ResultSet result = statement.executeQuery("SELECT node_id, host_id, endpoint FROM cluster_nodes WHERE state <> 'retired'")) {
@@ -78,7 +97,11 @@ final class NodeRegistry {
(UUID) result.getObject(2), url)); (UUID) result.getObject(2), url));
} }
} }
if (stored.isEmpty()) { return stored;
}
private static void registerInitialNodes(Connection connection, List<URI> urls, String token,
Map<String, NodeClient.Node> stored) throws SQLException, IOException {
try (Statement statement = connection.createStatement(); try (Statement statement = connection.createStatement();
ResultSet result = statement.executeQuery("SELECT EXISTS (SELECT 1 FROM cluster_segments)")) { ResultSet result = statement.executeQuery("SELECT EXISTS (SELECT 1 FROM cluster_segments)")) {
result.next(); result.next();
@@ -96,20 +119,25 @@ final class NodeRegistry {
stored.put(url.toString(), new NodeClient.Node(identity.nodeId(), identity.hostId(), url)); stored.put(url.toString(), new NodeClient.Node(identity.nodeId(), identity.hostId(), url));
} }
} }
private static List<NodeClient.Node> configuredNodes(List<URI> urls, String token,
Map<String, NodeClient.Node> stored) throws IOException {
List<NodeClient.Node> configured = new ArrayList<>(); List<NodeClient.Node> configured = new ArrayList<>();
Set<UUID> configuredIds = new HashSet<>();
for (URI url : urls) { for (URI url : urls) {
NodeClient.Node node = stored.get(url.toString()); NodeClient.Node node = stored.get(url.toString());
if (node == null) throw new IOException("Unregistered storage node URL: " + url); if (node == null) throw new IOException("Unregistered storage node URL: " + url);
NodeIdentity actual = null; NodeIdentity actual = NodeClient.probeIfAvailable(url, token);
try {
actual = NodeClient.probe(url, token);
} catch (IOException offline) { }
if (actual != null && (!actual.nodeId().equals(node.id()) || !actual.hostId().equals(node.hostId()))) if (actual != null && (!actual.nodeId().equals(node.id()) || !actual.hostId().equals(node.hostId())))
throw new IOException("Storage node identity changed at " + url); throw new IOException("Storage node identity changed at " + url);
configured.add(node); configured.add(node);
configuredIds.add(node.id());
} }
return configured;
}
private static void ensureLiveReplicasConfigured(Connection connection, List<NodeClient.Node> configured)
throws SQLException, IOException {
Set<UUID> configuredIds = new HashSet<>();
for (NodeClient.Node node : configured) configuredIds.add(node.id());
try (Statement statement = connection.createStatement(); try (Statement statement = connection.createStatement();
ResultSet result = statement.executeQuery( ResultSet result = statement.executeQuery(
"SELECT DISTINCT unnest(s.replica_ids) FROM cluster_segments s JOIN cluster_objects o ON o.generation=s.generation")) { "SELECT DISTINCT unnest(s.replica_ids) FROM cluster_segments s JOIN cluster_objects o ON o.generation=s.generation")) {
@@ -119,17 +147,5 @@ final class NodeRegistry {
throw new IOException("A live segment refers to a node missing from CLUSTER_NODES: " + id); throw new IOException("A live segment refers to a node missing from CLUSTER_NODES: " + id);
} }
} }
NodeClient nodes = new NodeClient(configured, token, repairToken);
connection.commit();
return nodes;
} catch (SQLException | IOException | RuntimeException error) {
try { connection.rollback(); } catch (SQLException rollback) { error.addSuppressed(rollback); }
if (error instanceof IOException io) throw io;
if (error instanceof SQLException sql) throw new IOException("Node registry check failed", sql);
throw (RuntimeException) error;
} finally {
try { connection.setAutoCommit(true); }
catch (SQLException error) { throw new IOException("Could not restore metadata connection", error); }
}
} }
} }
+38 -12
View File
@@ -28,6 +28,23 @@ final class SigV4 {
} }
String verify(String method, URI uri, Headers headers) { String verify(String method, URI uri, Headers headers) {
Map<String, String> fields = authorizationFields(headers);
String[] credential = credentialScope(fields.get("Credential"));
String date = signingDate(headers, credential[1]);
String payload = payloadHash(headers);
String signedHeaders = fields.get("SignedHeaders");
String canonicalHeaders = canonicalHeaders(headers, signedHeaders);
String canonical = method + "\n" + encode(decode(uri.getRawPath()), true) + "\n"
+ canonicalQuery(uri.getRawQuery()) + "\n" + canonicalHeaders + "\n" + signedHeaders + "\n" + payload;
String scope = String.join("/", Arrays.copyOfRange(credential, 1, 5));
String toSign = "AWS4-HMAC-SHA256\n" + date + "\n" + scope + "\n" + hex(hash(canonical.getBytes(StandardCharsets.UTF_8)));
byte[] signingKey = signingKey(secretKey, credential[1], region);
String signature = fields.get("Signature");
if (!HEX.matcher(signature).matches() || !MessageDigest.isEqual(hmac(signingKey, toSign), HexFormat.of().parseHex(signature))) denied("Signature mismatch");
return payload;
}
private static Map<String, String> authorizationFields(Headers headers) {
String authorization = single(headers, "authorization"); String authorization = single(headers, "authorization");
if (authorization == null || !authorization.startsWith("AWS4-HMAC-SHA256 ")) denied("Signed requests are required"); if (authorization == null || !authorization.startsWith("AWS4-HMAC-SHA256 ")) denied("Signed requests are required");
Map<String, String> fields = new TreeMap<>(); Map<String, String> fields = new TreeMap<>();
@@ -36,20 +53,36 @@ final class SigV4 {
if (pair.length != 2 || fields.put(pair[0], pair[1]) != null) denied("Invalid authorization header"); if (pair.length != 2 || fields.put(pair[0], pair[1]) != null) denied("Invalid authorization header");
} }
if (!fields.keySet().equals(java.util.Set.of("Credential", "SignedHeaders", "Signature"))) denied("Invalid authorization fields"); if (!fields.keySet().equals(java.util.Set.of("Credential", "SignedHeaders", "Signature"))) denied("Invalid authorization fields");
String[] credential = fields.get("Credential").split("/", -1); return fields;
}
private String[] credentialScope(String value) {
String[] credential = value.split("/", -1);
if (credential.length != 5 || !credential[0].equals(accessKey) || !credential[2].equals(region) if (credential.length != 5 || !credential[0].equals(accessKey) || !credential[2].equals(region)
|| !credential[3].equals("s3") || !credential[4].equals("aws4_request")) denied("Invalid credential scope"); || !credential[3].equals("s3") || !credential[4].equals("aws4_request")) denied("Invalid credential scope");
String date = single(headers, "x-amz-date"), payload = single(headers, "x-amz-content-sha256"); return credential;
if (date == null || !credential[1].matches("[0-9]{8}") || !date.matches("[0-9]{8}T[0-9]{6}Z") || !date.startsWith(credential[1])) denied("Invalid signing date"); }
private String signingDate(Headers headers, String credentialDate) {
String date = single(headers, "x-amz-date");
if (date == null || !credentialDate.matches("[0-9]{8}") || !date.matches("[0-9]{8}T[0-9]{6}Z") || !date.startsWith(credentialDate)) denied("Invalid signing date");
try { try {
Instant signed = Instant.from(DATE.parse(date)); Instant signed = Instant.from(DATE.parse(date));
if (Duration.between(signed, clock.instant()).abs().compareTo(Duration.ofMinutes(5)) > 0) if (Duration.between(signed, clock.instant()).abs().compareTo(Duration.ofMinutes(5)) > 0)
throw new StoreException(403, "RequestTimeTooSkewed", "Request timestamp is outside the permitted window"); throw new StoreException(403, "RequestTimeTooSkewed", "Request timestamp is outside the permitted window");
} catch (java.time.DateTimeException e) { denied("Invalid signing date"); } } catch (java.time.DateTimeException e) { denied("Invalid signing date"); }
return date;
}
private static String payloadHash(Headers headers) {
String payload = single(headers, "x-amz-content-sha256");
if (payload == null || !HEX.matcher(payload).matches()) if (payload == null || !HEX.matcher(payload).matches())
throw new StoreException(400, "NotImplemented", "A hexadecimal SHA-256 payload hash is required; unsigned and chunk-signed payloads are unsupported"); throw new StoreException(400, "NotImplemented", "A hexadecimal SHA-256 payload hash is required; unsigned and chunk-signed payloads are unsupported");
if (headers.containsKey("x-amz-security-token")) denied("Temporary credentials are unsupported"); if (headers.containsKey("x-amz-security-token")) denied("Temporary credentials are unsupported");
String signedHeaders = fields.get("SignedHeaders"); return payload;
}
private static String canonicalHeaders(Headers headers, String signedHeaders) {
String[] names = signedHeaders.split(";", -1); String[] names = signedHeaders.split(";", -1);
if (names.length > 32 || !signedHeaders.equals(String.join(";", Arrays.stream(names).distinct().sorted().toList()))) denied("Signed headers must be unique and sorted"); if (names.length > 32 || !signedHeaders.equals(String.join(";", Arrays.stream(names).distinct().sorted().toList()))) denied("Signed headers must be unique and sorted");
var namesSet = java.util.Set.copyOf(Arrays.asList(names)); var namesSet = java.util.Set.copyOf(Arrays.asList(names));
@@ -66,14 +99,7 @@ final class SigV4 {
if (value == null) denied("Missing signed header"); if (value == null) denied("Missing signed header");
canonicalHeaders.append(name).append(':').append(value.trim().replaceAll("[\\t ]+", " ")).append('\n'); canonicalHeaders.append(name).append(':').append(value.trim().replaceAll("[\\t ]+", " ")).append('\n');
} }
String canonical = method + "\n" + encode(decode(uri.getRawPath()), true) + "\n" return canonicalHeaders.toString();
+ canonicalQuery(uri.getRawQuery()) + "\n" + canonicalHeaders + "\n" + signedHeaders + "\n" + payload;
String scope = String.join("/", Arrays.copyOfRange(credential, 1, 5));
String toSign = "AWS4-HMAC-SHA256\n" + date + "\n" + scope + "\n" + hex(hash(canonical.getBytes(StandardCharsets.UTF_8)));
byte[] signingKey = signingKey(secretKey, credential[1], region);
String signature = fields.get("Signature");
if (!HEX.matcher(signature).matches() || !MessageDigest.isEqual(hmac(signingKey, toSign), HexFormat.of().parseHex(signature))) denied("Signature mismatch");
return payload;
} }
static String single(Headers headers, String name) { static String single(Headers headers, String name) {
+1 -1
View File
@@ -1,7 +1,7 @@
package cloud.lunarsky.store; package cloud.lunarsky.store;
final class Version { final class Version {
static final String VALUE = "0.0.2"; static final String VALUE = "0.0.3";
private Version() {} private Version() {}
} }
+30 -13
View File
@@ -66,19 +66,7 @@ public final class HttpTest {
} }
} }
public static void main(String[] args) throws Exception { private static void testObjects(HttpClient client, String base) throws Exception {
Path root = Files.createTempDirectory("store-http-test-");
var executor = Executors.newVirtualThreadPerTaskExecutor();
HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 16);
DiskStore store = new DiskStore(root, 1024, 4096);
try {
var app = new Main(store,
new SigV4(ACCESS, SECRET, REGION, Clock.systemUTC()), "objects");
server.setExecutor(executor);
server.createContext("/", app::handle);
server.start();
String base = "http://127.0.0.1:" + server.getAddress().getPort();
HttpClient client = HttpClient.newHttpClient();
status(200, client.send(HttpRequest.newBuilder(URI.create(base + "/health")).GET().build(), status(200, client.send(HttpRequest.newBuilder(URI.create(base + "/health")).GET().build(),
HttpResponse.BodyHandlers.ofByteArray())); HttpResponse.BodyHandlers.ofByteArray()));
String key = "folder/moon-☾.txt"; String key = "folder/moon-☾.txt";
@@ -109,6 +97,9 @@ public final class HttpTest {
throw new AssertionError("Range response mismatch"); throw new AssertionError("Range response mismatch");
status(416, client.send(signedUri(URI.create(base + "/objects/" + other), "GET", status(416, client.send(signedUri(URI.create(base + "/objects/" + other), "GET",
new byte[0], Map.of("range", "bytes=20-30")), HttpResponse.BodyHandlers.ofByteArray())); new byte[0], Map.of("range", "bytes=20-30")), HttpResponse.BodyHandlers.ofByteArray()));
}
private static void testListing(HttpClient client, String base) throws Exception {
var listed = client.send(signedUri(URI.create(base + "/objects?list-type=2&prefix=folder%2F"), var listed = client.send(signedUri(URI.create(base + "/objects?list-type=2&prefix=folder%2F"),
"GET", new byte[0], Map.of()), HttpResponse.BodyHandlers.ofString()); "GET", new byte[0], Map.of()), HttpResponse.BodyHandlers.ofString());
if (listed.statusCode() != 200 || !listed.body().contains("<Key>folder/stars.txt</Key>") || if (listed.statusCode() != 200 || !listed.body().contains("<Key>folder/stars.txt</Key>") ||
@@ -136,6 +127,9 @@ public final class HttpTest {
"GET", new byte[0], Map.of()), HttpResponse.BodyHandlers.ofString()); "GET", new byte[0], Map.of()), HttpResponse.BodyHandlers.ofString());
if (emptyPage.statusCode() != 200 || !emptyPage.body().contains("<KeyCount>0</KeyCount>")) if (emptyPage.statusCode() != 200 || !emptyPage.body().contains("<KeyCount>0</KeyCount>"))
throw new AssertionError("Empty list page failed: " + emptyPage.body()); throw new AssertionError("Empty list page failed: " + emptyPage.body());
}
private static void testMultipart(HttpClient client, String base) throws Exception {
String movie = "folder/video.mp4"; String movie = "folder/video.mp4";
URI initiate = URI.create(base + "/objects/" + movie + "?uploads="); URI initiate = URI.create(base + "/objects/" + movie + "?uploads=");
var created = client.send(signedUri(initiate, "POST", new byte[0], var created = client.send(signedUri(initiate, "POST", new byte[0],
@@ -166,6 +160,10 @@ public final class HttpTest {
String abandonedId = abandoned.body().split("<UploadId>")[1].split("</UploadId>")[0]; String abandonedId = abandoned.body().split("<UploadId>")[1].split("</UploadId>")[0];
status(204, client.send(signedUri(URI.create(base + "/objects/abandoned?uploadId=" + abandonedId), status(204, client.send(signedUri(URI.create(base + "/objects/abandoned?uploadId=" + abandonedId),
"DELETE", new byte[0], Map.of()), HttpResponse.BodyHandlers.ofByteArray())); "DELETE", new byte[0], Map.of()), HttpResponse.BodyHandlers.ofByteArray()));
}
private static void testDelete(HttpClient client, String base) throws Exception {
String key = "folder/moon-☾.txt";
var head = client.send(signed(base, "HEAD", key, new byte[0]), var head = client.send(signed(base, "HEAD", key, new byte[0]),
HttpResponse.BodyHandlers.ofByteArray()); HttpResponse.BodyHandlers.ofByteArray());
status(200, head); status(200, head);
@@ -174,6 +172,25 @@ public final class HttpTest {
HttpResponse.BodyHandlers.ofByteArray())); HttpResponse.BodyHandlers.ofByteArray()));
status(404, client.send(signed(base, "GET", key, new byte[0]), status(404, client.send(signed(base, "GET", key, new byte[0]),
HttpResponse.BodyHandlers.ofByteArray())); HttpResponse.BodyHandlers.ofByteArray()));
}
public static void main(String[] args) throws Exception {
Path root = Files.createTempDirectory("store-http-test-");
var executor = Executors.newVirtualThreadPerTaskExecutor();
HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 16);
DiskStore store = new DiskStore(root, 1024, 4096);
try {
var app = new Main(store,
new SigV4(ACCESS, SECRET, REGION, Clock.systemUTC()), "objects");
server.setExecutor(executor);
server.createContext("/", app::handle);
server.start();
String base = "http://127.0.0.1:" + server.getAddress().getPort();
HttpClient client = HttpClient.newHttpClient();
testObjects(client, base);
testListing(client, base);
testMultipart(client, base);
testDelete(client, base);
System.out.println("HTTP tests passed: health, authentication, PUT, GET, HEAD, DELETE, MIME, ranges, listing, multipart"); System.out.println("HTTP tests passed: health, authentication, PUT, GET, HEAD, DELETE, MIME, ranges, listing, multipart");
} finally { } finally {
server.stop(0); server.stop(0);
+32 -4
View File
@@ -14,7 +14,7 @@ public final class StoreTest {
static ObjectStorage.Metadata put(DiskStore store,String key,byte[] body,boolean only)throws Exception{ static ObjectStorage.Metadata put(DiskStore store,String key,byte[] body,boolean only)throws Exception{
return store.put("test",key,new ByteArrayInputStream(body),body.length,SigV4.hex(SigV4.hash(body)),null,only,"application/octet-stream"); return store.put("test",key,new ByteArrayInputStream(body),body.length,SigV4.hex(SigV4.hash(body)),null,only,"application/octet-stream");
} }
public static void main(String[] args)throws Exception{ private static void testSignature() throws Exception {
var headers=new com.sun.net.httpserver.Headers(); var headers=new com.sun.net.httpserver.Headers();
headers.set("host","examplebucket.s3.amazonaws.com");headers.set("range","bytes=0-9"); headers.set("host","examplebucket.s3.amazonaws.com");headers.set("range","bytes=0-9");
headers.set("x-amz-date","20130524T000000Z"); headers.set("x-amz-date","20130524T000000Z");
@@ -29,8 +29,9 @@ public final class StoreTest {
fails(403,()->new SigV4("AKIAIOSFODNN7EXAMPLE","wrong","us-east-1",java.time.Clock.systemUTC()).verify("GET",uri,headers)); fails(403,()->new SigV4("AKIAIOSFODNN7EXAMPLE","wrong","us-east-1",java.time.Clock.systemUTC()).verify("GET",uri,headers));
headers.add("host","duplicate");fails(403,()->auth.verify("GET",uri,headers)); headers.add("host","duplicate");fails(403,()->auth.verify("GET",uri,headers));
System.out.println("SigV4 official vector and tampering tests passed"); System.out.println("SigV4 official vector and tampering tests passed");
Path root=Files.createTempDirectory("store-test-"); }
try{
private static void testInitialStore(Path root) throws Exception {
try(var store=new DiskStore(root,8,10)){ try(var store=new DiskStore(root,8,10)){
byte[] body={1,2,3,4,5,6}; byte[] body={1,2,3,4,5,6};
put(store,"../nested/☾",body,true); put(store,"../nested/☾",body,true);
@@ -52,6 +53,9 @@ public final class StoreTest {
if(!expected.getMessage().contains("already in use"))throw expected; if(!expected.getMessage().contains("already in use"))throw expected;
} }
} }
}
private static void testRestart(Path root) throws Exception {
try(var restarted=new DiskStore(root,8,10)){ try(var restarted=new DiskStore(root,8,10)){
try(var obj=restarted.open("test","../nested/☾")){if(obj.stream().read()!=9)throw new AssertionError("Persistence");} try(var obj=restarted.open("test","../nested/☾")){if(obj.stream().read()!=9)throw new AssertionError("Persistence");}
if(restarted.list("test","","",100,null).objects().size()!=2)throw new AssertionError("Index persistence"); if(restarted.list("test","","",100,null).objects().size()!=2)throw new AssertionError("Index persistence");
@@ -64,6 +68,9 @@ public final class StoreTest {
SigV4.hex(SigV4.hash(part)),null); SigV4.hex(SigV4.hash(part)),null);
Files.writeString(root.resolve("pending-upload-id"),upload); Files.writeString(root.resolve("pending-upload-id"),upload);
} }
}
private static void testMultipartRecovery(Path root) throws Exception {
try(var resumed=new DiskStore(root,8,10)){ try(var resumed=new DiskStore(root,8,10)){
Path unfinished=root.resolve("multipart/.creating-00000000-0000-0000-0000-000000000000"); Path unfinished=root.resolve("multipart/.creating-00000000-0000-0000-0000-000000000000");
Files.createDirectory(unfinished); Files.createDirectory(unfinished);
@@ -80,6 +87,9 @@ public final class StoreTest {
} }
resumed.delete("test","from-parts"); resumed.delete("test","from-parts");
} }
}
private static void testLegacyRecord(Path root) throws Exception {
byte[] old={4,5,6}; byte[] old={4,5,6};
String oldId=SigV4.hex(SigV4.hash("test/legacy".getBytes(java.nio.charset.StandardCharsets.UTF_8))); String oldId=SigV4.hex(SigV4.hash("test/legacy".getBytes(java.nio.charset.StandardCharsets.UTF_8)));
Path oldPath=root.resolve("objects").resolve(oldId.substring(0,2)).resolve(oldId); Path oldPath=root.resolve("objects").resolve(oldId.substring(0,2)).resolve(oldId);
@@ -100,6 +110,9 @@ public final class StoreTest {
if(migrated.list("test","","",100,null).objects().stream().noneMatch(entry->entry.key().equals("legacy"))) if(migrated.list("test","","",100,null).objects().stream().noneMatch(entry->entry.key().equals("legacy")))
throw new AssertionError("Legacy overwrite was not indexed"); throw new AssertionError("Legacy overwrite was not indexed");
} }
}
private static void testCorruption(Path root) throws Exception {
Files.delete(root.resolve("pending-upload-id")); Files.delete(root.resolve("pending-upload-id"));
String id=SigV4.hex(SigV4.hash("test/empty".getBytes(java.nio.charset.StandardCharsets.UTF_8))); String id=SigV4.hex(SigV4.hash("test/empty".getBytes(java.nio.charset.StandardCharsets.UTF_8)));
Files.write(root.resolve("objects").resolve(id.substring(0,2)).resolve(id),new byte[]{1},StandardOpenOption.APPEND); Files.write(root.resolve("objects").resolve(id.substring(0,2)).resolve(id),new byte[]{1},StandardOpenOption.APPEND);
@@ -110,7 +123,22 @@ public final class StoreTest {
if(!expected.getMessage().contains("object record"))throw expected; if(!expected.getMessage().contains("object record"))throw expected;
} }
try(var pending=Files.list(root.resolve("pending"))){if(pending.count()!=0)throw new AssertionError("Pending cleanup");} try(var pending=Files.list(root.resolve("pending"))){if(pending.count()!=0)throw new AssertionError("Pending cleanup");}
}
public static void main(String[] args) throws Exception {
testSignature();
Path root = Files.createTempDirectory("store-test-");
try {
testInitialStore(root);
testRestart(root);
testMultipartRecovery(root);
testLegacyRecord(root);
testCorruption(root);
System.out.println("Java storage tests passed: roundtrip, quota, indexing, persistence, multipart recovery, legacy reads, locking, corruption, delete"); System.out.println("Java storage tests passed: roundtrip, quota, indexing, persistence, multipart recovery, legacy reads, locking, corruption, delete");
}finally{try(var paths=Files.walk(root)){for(var p:paths.sorted(java.util.Comparator.reverseOrder()).toList())Files.delete(p);}} } finally {
try (var paths = Files.walk(root)) {
for (var path : paths.sorted(java.util.Comparator.reverseOrder()).toList()) Files.delete(path);
}
}
} }
} }