Release ObjectStore 0.0.4

This commit is contained in:
admin committed 2026-10-10 10:01:43 +02:00
1 parent 00d2af2fb0
commit ed7712a8af
19 files changed
+1025 -114

No files matched your search

+43
View File
@@ -0,0 +1,43 @@
package cloud.lunarsky.store;
import java.net.URI;
import java.util.Arrays;
import java.util.Map;
public final class ClusterGc {
public static void main(String[] args) throws Exception {
if (args.length > 1 || (args.length == 1 && !args[0].equals("--apply")))
throw new IllegalArgumentException("Usage: objectstore cluster-gc [--apply]");
run(System.getenv(), args.length == 1);
}
static void run(Map<String, String> env, boolean apply) throws Exception {
if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV")))
throw new IllegalArgumentException("Cluster garbage collection is only enabled in local cluster mode");
boolean testDomains = "true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS"));
boolean disposableTest = testDomains && "true".equals(env.get("CLUSTER_GC_TEST_MODE"));
long age = Long.parseLong(env.getOrDefault("CLUSTER_GC_MIN_AGE_SECONDS", "1209600"));
long backupRetention = Long.parseLong(env.getOrDefault("CLUSTER_BACKUP_RETENTION_SECONDS", "0"));
if (age < 0 || age > 315360000 || (age == 0 && !disposableTest))
throw new IllegalArgumentException("Invalid garbage collection age");
if (apply && !disposableTest && (backupRetention < 86400 || age <= backupRetention))
throw new IllegalArgumentException("Set a garbage collection age longer than the backup retention");
try (ClusterStore store = new ClusterStore(env.get("POSTGRES_JDBC_URL"), env.get("POSTGRES_USER"),
env.get("POSTGRES_PASSWORD"), env.get("S3_BUCKET"),
Arrays.stream(env.get("CLUSTER_NODES").split(",")).map(URI::create).toList(),
env.get("CLUSTER_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"), 134217728, 2147483648L,
testDomains)) {
if (apply) {
var repair = store.repairOnce();
if (repair.underReplicated() > 0 || repair.unrecoverable() > 0)
throw new IllegalStateException("Refusing cleanup while live segments need repair");
}
var report = store.collectGarbage(age * 1000, apply);
System.out.println("segments_scanned=" + report.scanned());
System.out.println("orphan_candidates=" + report.eligible());
System.out.println("segments_deleted=" + report.deleted());
System.out.println("unavailable_nodes=" + report.unavailableNodes());
if (report.unavailableNodes() > 0) throw new IllegalStateException("Cleanup did not scan every node");
}
}
}
+80
View File
@@ -15,7 +15,9 @@ import java.nio.file.StandardCopyOption;
import java.nio.file.StandardOpenOption;
import java.security.MessageDigest;
import java.util.Arrays;
import java.util.Comparator;
import java.util.HexFormat;
import java.util.List;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.Executors;
@@ -84,11 +86,18 @@ public final class ClusterNode implements AutoCloseable {
return;
}
if (!expectedNodeAndRepairAuthorized(exchange)) return;
if (path.equals("/segments") && exchange.getRequestMethod().equals("GET")) {
if (maintenanceAuthorized(exchange)) inventory(exchange);
return;
}
String id = segmentId(exchange, path);
if (id == null) return;
switch (exchange.getRequestMethod()) {
case "PUT" -> put(exchange, segmentPath(id, true));
case "GET" -> get(exchange, segmentPath(id, false));
case "DELETE" -> {
if (maintenanceAuthorized(exchange)) delete(exchange, segmentPath(id, false));
}
default -> respond(exchange, 405, "Method not allowed");
}
} catch (IllegalArgumentException error) {
@@ -130,6 +139,77 @@ public final class ClusterNode implements AutoCloseable {
return true;
}
private boolean maintenanceAuthorized(HttpExchange exchange) throws IOException {
String supplied = exchange.getRequestHeaders().getFirst("X-Cluster-Repair-Token");
byte[] value = supplied == null ? new byte[0] : supplied.getBytes(java.nio.charset.StandardCharsets.UTF_8);
if (MessageDigest.isEqual(repairToken, value)) return true;
respond(exchange, 403, "Repair authority required");
return false;
}
private void inventory(HttpExchange exchange) throws IOException {
String query = exchange.getRequestURI().getRawQuery();
if (query == null || !query.matches("shard=[0-9a-f]{2}(&after=[0-9a-f-]{36})?")) {
respond(exchange, 400, "Invalid inventory request");
return;
}
String shard = query.substring(6, 8);
String after = query.length() > 8 ? query.substring(15) : "";
if (!after.isEmpty() && (!after.startsWith(shard) || !UUID.fromString(after).toString().equals(after))) {
respond(exchange, 400, "Invalid inventory cursor");
return;
}
Path directory = segments.resolve(shard);
if (!Files.isDirectory(directory)) {
respond(exchange, 200, "");
return;
}
List<Path> files;
try (var entries = Files.list(directory)) {
files = entries.filter(Files::isRegularFile).sorted(Comparator.comparing(path ->
path.getFileName().toString())).toList();
}
StringBuilder body = new StringBuilder();
int count = 0;
for (Path file : files) {
String id = file.getFileName().toString();
if (id.compareTo(after) <= 0 || !id.matches("[0-9a-f-]{36}")) continue;
if (!UUID.fromString(id).toString().equals(id)) continue;
body.append(id).append(' ').append(Files.getLastModifiedTime(file).toMillis()).append('\n');
if (++count == 1000) break;
}
respond(exchange, 200, body.toString());
}
private synchronized void delete(HttpExchange exchange, Path target) throws IOException {
String expected = exchange.getRequestHeaders().getFirst("X-Cluster-Expected-Mtime");
String age = exchange.getRequestHeaders().getFirst("X-Cluster-Gc-Min-Age-Millis");
long expectedTime, minimumAge;
try {
expectedTime = Long.parseLong(expected);
minimumAge = Long.parseLong(age);
} catch (NumberFormatException error) {
respond(exchange, 400, "Invalid deletion guard");
return;
}
if (minimumAge < 0 || minimumAge > System.currentTimeMillis()) {
respond(exchange, 400, "Invalid deletion age");
return;
}
if (!Files.isRegularFile(target)) {
exchange.sendResponseHeaders(404, -1);
return;
}
long modified = Files.getLastModifiedTime(target).toMillis();
if (modified != expectedTime || modified > System.currentTimeMillis() - minimumAge) {
exchange.sendResponseHeaders(409, -1);
return;
}
Files.delete(target);
DiskStore.syncDirectory(target.getParent());
exchange.sendResponseHeaders(204, -1);
}
private static String segmentId(HttpExchange exchange, String path) throws IOException {
if (!path.matches("/segments/[0-9a-f-]{36}")) {
respond(exchange, 404, "Not found");
+33 -13
View File
@@ -6,21 +6,41 @@ import java.util.Map;
public final class ClusterRepair {
public static void main(String[] args) throws Exception {
if (args.length != 0) throw new IllegalArgumentException("Usage: objectstore cluster-repair");
if (args.length > 1 || (args.length == 1 && !args[0].equals("--loop")))
throw new IllegalArgumentException("Usage: objectstore cluster-repair [--loop]");
Map<String, String> env = System.getenv();
if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV")))
throw new IllegalArgumentException("Cluster repair is only enabled in local cluster mode");
try (ClusterStore store = new ClusterStore(env.get("POSTGRES_JDBC_URL"), env.get("POSTGRES_USER"),
env.get("POSTGRES_PASSWORD"), env.get("S3_BUCKET"),
Arrays.stream(env.get("CLUSTER_NODES").split(",")).map(URI::create).toList(),
env.get("CLUSTER_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"), 134217728, 2147483648L,
"true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS")))) {
var report = store.repairOnce();
System.out.println("segments_scanned=" + report.scanned());
System.out.println("replicas_restored=" + report.restored());
System.out.println("segments_under_replicated=" + report.underReplicated());
System.out.println("segments_unrecoverable=" + report.unrecoverable());
if (report.unrecoverable() > 0) System.exit(1);
}
boolean loop = args.length == 1;
long seconds = Long.parseLong(env.getOrDefault("CLUSTER_MAINTENANCE_INTERVAL_SECONDS", "60"));
if (seconds < 1 || seconds > 3600) throw new IllegalArgumentException("Invalid maintenance interval");
boolean gcEnabled = loop && "true".equals(env.get("CLUSTER_GC_ENABLED"));
long gcInterval = Long.parseLong(env.getOrDefault("CLUSTER_GC_INTERVAL_SECONDS", "86400"));
if (gcEnabled && (gcInterval < 1 || gcInterval > 604800))
throw new IllegalArgumentException("Invalid garbage collection interval");
long nextGc = 0;
do {
try (ClusterStore store = new ClusterStore(env.get("POSTGRES_JDBC_URL"), env.get("POSTGRES_USER"),
env.get("POSTGRES_PASSWORD"), env.get("S3_BUCKET"),
Arrays.stream(env.get("CLUSTER_NODES").split(",")).map(URI::create).toList(),
env.get("CLUSTER_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"), 134217728, 2147483648L,
"true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS")))) {
var report = store.repairOnce();
System.out.println("segments_scanned=" + report.scanned());
System.out.println("replicas_restored=" + report.restored());
System.out.println("segments_rebalanced=" + report.rebalanced());
System.out.println("segments_under_replicated=" + report.underReplicated());
System.out.println("segments_unrecoverable=" + report.unrecoverable());
if (!loop && report.unrecoverable() > 0) System.exit(1);
if (gcEnabled && System.currentTimeMillis() >= nextGc) {
ClusterGc.run(env, true);
nextGc = System.currentTimeMillis() + gcInterval * 1000;
}
} catch (Exception error) {
if (!loop) throw error;
System.err.println("Cluster maintenance failed: " + error.getMessage());
}
if (loop) Thread.sleep(seconds * 1000);
} while (loop);
}
}
+215 -67
View File
@@ -21,6 +21,7 @@ import java.util.Base64;
import java.util.Comparator;
import java.util.HexFormat;
import java.util.HashSet;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Set;
import java.util.UUID;
@@ -30,7 +31,8 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
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 rebalanced, int underReplicated, int unrecoverable) {}
record GcReport(int scanned, int eligible, int deleted, int unavailableNodes) {}
private final String jdbcUrl, user, password, configuredBucket;
private final NodeClient nodes;
private final long maxObject, maxTotal;
@@ -57,22 +59,29 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
private Connection connect() throws SQLException { return DriverManager.getConnection(jdbcUrl, user, password); }
private static void lockGc(Connection connection, boolean shared) throws SQLException {
try (var statement = connection.createStatement()) {
statement.execute("SELECT pg_advisory_lock" + (shared ? "_shared" : "") + "(6834071092783)");
}
}
@Override public Metadata put(String bucket, String key, InputStream input, long length, String expectedHash,
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);
byte[] fullHash = stageInput(staged, input, length, expectedHash, checksum, md5);
checkCapacity(bucket, key, length, createOnly);
segments = uploadSegments(staged, length);
try (Connection connection = connect()) {
lockGc(connection, true);
List<Segment> segments = uploadSegments(staged, length);
Metadata metadata = new Metadata(length, Instant.now().toEpochMilli(),
HexFormat.of().formatHex(md5.digest()), fullHash, bucket, key, contentType);
persistObject(connection, metadata, segments, createOnly);
return metadata;
} catch (SQLException error) { throw databaseError(error); }
} 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) {
@@ -161,11 +170,12 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
return segments;
}
private void persistObject(Metadata metadata, List<Segment> segments, boolean createOnly) throws IOException {
private void persistObject(Connection connection, Metadata metadata, List<Segment> segments,
boolean createOnly) throws IOException {
String bucket = metadata.bucket(), key = metadata.key();
long length = metadata.length();
UUID generation = UUID.randomUUID();
try (Connection connection = connect()) {
try {
connection.setAutoCommit(false);
try {
long used = lockUsage(connection, bucket);
@@ -258,58 +268,58 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
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);
try (Connection connection = connect()) {
lockGc(connection, true);
List<Segment> segments = uploadSegments(staged, length);
String etag = HexFormat.of().formatHex(md5.digest());
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.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.setLong(3, length);
insert.setString(4, etag);
insert.setLong(5, Instant.now().toEpochMilli());
insert.executeUpdate();
}
insert.executeBatch();
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;
}
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;
return etag;
} catch (SQLException error) { throw databaseError(error); }
} finally { Files.deleteIfExists(staged); }
}
@Override public Metadata complete(String id, String bucket, String key, List<MultipartStorage.Part> parts) throws IOException {
@@ -487,6 +497,7 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
try (Connection connection = connect()) {
connection.setAutoCommit(false);
connection.setTransactionIsolation(Connection.TRANSACTION_REPEATABLE_READ);
lockGc(connection, true);
try {
Metadata metadata;
UUID generation;
@@ -790,9 +801,12 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
} catch (SQLException error) { return false; }
}
RepairReport repairOnce() throws IOException {
int scanned = 0, restored = 0, underReplicated = 0, unrecoverable = 0;
int scanned = 0, restored = 0, rebalanced = 0, underReplicated = 0, unrecoverable = 0;
try (Connection reader = connect()) {
reader.setAutoCommit(false);
try (var lock = reader.createStatement()) {
lock.execute("SELECT pg_advisory_xact_lock(6834071092782)");
}
try (PreparedStatement query = reader.prepareStatement(
"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 " +
@@ -808,7 +822,7 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
result.getInt(5), result.getBytes(6), replicaIds(result, 7)));
Segment segment = target.segment();
byte[] copy = null;
Set<UUID> healthy = new HashSet<>();
Set<UUID> healthy = new LinkedHashSet<>();
Set<UUID> healthyHosts = new HashSet<>();
for (UUID id : segment.replicas()) {
int node = nodes.index(id);
@@ -823,20 +837,45 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
unrecoverable++;
continue;
}
List<UUID> preferred = new ArrayList<>();
Set<UUID> preferredHosts = new HashSet<>();
for (int node : PlacementPolicy.candidates(segment.id(), nodes, testNodeDomains)) {
UUID host = nodes.faultDomain(node, testNodeDomains);
if (healthyHosts.contains(host)) continue;
if (!preferredHosts.add(host)) continue;
preferred.add(nodes.node(node).id());
if (preferred.size() == 3) break;
}
for (UUID id : preferred) {
int node = nodes.index(id);
UUID host = nodes.faultDomain(node, testNodeDomains);
if (healthy.contains(id)) continue;
if (repairReplica(node, segment, copy)) {
healthy.add(nodes.node(node).id());
healthy.add(id);
healthyHosts.add(host);
restored++;
}
if (healthyHosts.size() == 3) break;
}
if (healthyHosts.size() < 3) underReplicated++;
Set<UUID> listed = new java.util.LinkedHashSet<>(segment.replicas());
listed.addAll(healthy);
if (listed.size() != segment.replicas().size()) {
List<UUID> listed = new ArrayList<>();
Set<UUID> listedHosts = new HashSet<>();
for (UUID id : preferred) {
if (healthy.contains(id)) {
listed.add(id);
listedHosts.add(nodes.faultDomain(nodes.index(id), testNodeDomains));
}
}
for (UUID id : segment.replicas()) {
if (listed.size() == 3) break;
int node = nodes.index(id);
if (healthy.contains(id) && node >= 0 &&
listedHosts.add(nodes.faultDomain(node, testNodeDomains))) listed.add(id);
}
if (listed.size() < 3) {
for (UUID id : segment.replicas()) {
if (!listed.contains(id)) listed.add(id);
}
}
if (!listed.equals(segment.replicas())) {
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=?";
@@ -849,7 +888,8 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
if (target.part() != 0) update.setInt(next++, target.part());
update.setInt(next++, target.ordinal());
update.setLong(next, target.version());
update.executeUpdate();
if (update.executeUpdate() == 1 && preferred.stream().anyMatch(id ->
!segment.replicas().contains(id) && listed.contains(id))) rebalanced++;
}
}
}
@@ -857,7 +897,7 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
}
reader.commit();
} catch (SQLException error) { throw databaseError(error); }
return new RepairReport(scanned, restored, underReplicated, unrecoverable);
return new RepairReport(scanned, restored, rebalanced, underReplicated, unrecoverable);
}
private byte[] readableReplica(int node, Segment segment) {
@@ -872,6 +912,114 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
} catch (IOException unavailable) { return false; }
}
GcReport collectGarbage(long minimumAgeMillis, boolean apply) throws IOException {
if (minimumAgeMillis < 0 || (minimumAgeMillis == 0 && !testNodeDomains))
throw new IllegalArgumentException("Invalid garbage collection age");
if (nodes.count() < 2 || nodes.repairTokenUnavailable())
throw new IllegalStateException("Garbage collection requires repair authority");
int scanned = 0, eligible = 0, deleted = 0, unavailable = 0;
try (Connection connection = connect()) {
try (var lock = connection.createStatement()) {
lock.execute("SELECT pg_advisory_lock(6834071092782)");
}
lockGc(connection, false);
try (PreparedStatement referenced = connection.prepareStatement(
"SELECT EXISTS (SELECT 1 FROM cluster_segments s JOIN cluster_objects o " +
"ON o.generation=s.generation WHERE s.segment_id=? AND ?=ANY(s.replica_ids) " +
"UNION ALL SELECT 1 FROM cluster_upload_segments s " +
"WHERE s.segment_id=? AND ?=ANY(s.replica_ids))");
PreparedStatement candidate = connection.prepareStatement(
"SELECT observed_mtime, first_seen FROM cluster_gc_candidates WHERE node_id=? AND segment_id=?");
PreparedStatement mark = connection.prepareStatement(
"INSERT INTO cluster_gc_candidates VALUES (?, ?, ?, ?) ON CONFLICT DO NOTHING");
PreparedStatement reset = connection.prepareStatement(
"UPDATE cluster_gc_candidates SET observed_mtime=?, first_seen=? WHERE node_id=? AND segment_id=?");
PreparedStatement clear = connection.prepareStatement(
"DELETE FROM cluster_gc_candidates WHERE node_id=? AND segment_id=?")) {
for (int node = 0; node < nodes.count(); node++) {
boolean reachable = true;
for (int shard = 0; shard < 256 && reachable; shard++) {
String prefix = "%02x".formatted(shard);
UUID after = null;
while (true) {
List<NodeClient.StoredSegment> page;
try { page = nodes.inventory(node, prefix, after); }
catch (IOException error) {
System.err.println("Cluster inventory failed for node " + nodes.node(node).id() +
": " + error.getMessage());
unavailable++;
reachable = false;
break;
}
for (NodeClient.StoredSegment segment : page) {
scanned++;
UUID nodeId = nodes.node(node).id();
referenced.setObject(1, segment.id());
referenced.setObject(2, nodeId);
referenced.setObject(3, segment.id());
referenced.setObject(4, nodeId);
try (ResultSet result = referenced.executeQuery()) {
result.next();
if (result.getBoolean(1)) {
if (apply) clearCandidate(clear, nodeId, segment.id());
continue;
}
}
eligible++;
if (apply) {
candidate.setObject(1, nodeId);
candidate.setObject(2, segment.id());
long now = System.currentTimeMillis();
boolean firstObservation = false;
long firstSeen = now;
try (ResultSet result = candidate.executeQuery()) {
if (!result.next()) firstObservation = true;
else if (result.getLong(1) != segment.modified()) firstObservation = true;
else firstSeen = result.getLong(2);
}
if (firstObservation) {
reset.setLong(1, segment.modified());
reset.setLong(2, now);
reset.setObject(3, nodeId);
reset.setObject(4, segment.id());
if (reset.executeUpdate() == 0) {
mark.setObject(1, nodeId);
mark.setObject(2, segment.id());
mark.setLong(3, segment.modified());
mark.setLong(4, now);
mark.executeUpdate();
}
continue;
}
if (now - firstSeen < minimumAgeMillis) continue;
try {
if (nodes.deleteOrphan(node, segment, minimumAgeMillis)) deleted++;
clearCandidate(clear, nodeId, segment.id());
} catch (IOException error) {
System.err.println("Cluster deletion failed for segment " + segment.id() +
": " + error.getMessage());
unavailable++;
reachable = false;
break;
}
}
}
if (!reachable || page.size() < 1000) break;
after = page.getLast().id();
}
}
}
}
} catch (SQLException error) { throw databaseError(error); }
return new GcReport(scanned, eligible, deleted, unavailable);
}
private static void clearCandidate(PreparedStatement clear, UUID node, UUID segment) throws SQLException {
clear.setObject(1, node);
clear.setObject(2, segment);
clear.executeUpdate();
}
@Override public void close() {}
private final class SegmentStream extends InputStream {
+67 -14
View File
@@ -114,7 +114,8 @@ public final class Main {
String method = exchange.getRequestMethod();
boolean multipartRequest = multipartRequest(method, query);
if (!multipartRequest && !query.isEmpty() && !(query.size() == 1 &&
("PutObject".equals(query.get("x-id")) || "GetObject".equals(query.get("x-id")) ||
("PutObject".equals(query.get("x-id")) || "CopyObject".equals(query.get("x-id")) ||
"GetObject".equals(query.get("x-id")) ||
"HeadObject".equals(query.get("x-id")) || "DeleteObject".equals(query.get("x-id")))))
unsupported("Query operation");
validateObjectHeaders(exchange.getRequestHeaders());
@@ -122,9 +123,19 @@ public final class Main {
handleMultipart(exchange, method, query, key, hash);
return;
}
if (!method.equals("PUT")) requireEmptyBody(exchange, hash);
boolean copy = exchange.getRequestHeaders().containsKey("x-amz-copy-source");
if (!method.equals("PUT") && (copy || exchange.getRequestHeaders().containsKey("content-md5") ||
exchange.getRequestHeaders().containsKey("x-amz-metadata-directive") ||
exchange.getRequestHeaders().containsKey("x-amz-sdk-checksum-algorithm") ||
exchange.getRequestHeaders().keySet().stream().anyMatch(name ->
name.toLowerCase(Locale.ROOT).startsWith("x-amz-checksum-"))))
unsupported("Object upload header");
if (!method.equals("PUT") || copy) requireEmptyBody(exchange, hash);
switch (method) {
case "PUT" -> putObject(exchange, key, hash);
case "PUT" -> {
if (copy) copyObject(exchange, key);
else putObject(exchange, key, hash);
}
case "GET", "HEAD" -> readObject(exchange, key);
case "DELETE" -> deleteObject(exchange, key);
default -> unsupported("HTTP method");
@@ -135,17 +146,14 @@ public final class Main {
for (String name : headers.keySet()) {
String lower = name.toLowerCase(java.util.Locale.ROOT);
if (lower.startsWith("x-amz-") && !java.util.Set.of("x-amz-date", "x-amz-content-sha256",
"x-amz-checksum-sha256", "x-amz-sdk-checksum-algorithm", "x-amz-user-agent").contains(lower))
"x-amz-sdk-checksum-algorithm", "x-amz-user-agent", "x-amz-copy-source",
"x-amz-metadata-directive").contains(lower) && !lower.startsWith("x-amz-checksum-"))
unsupported("Amazon header");
if (lower.startsWith("x-amz-meta-") || lower.startsWith("x-amz-server-side-") ||
lower.startsWith("x-amz-copy-") || lower.startsWith("x-amz-acl") ||
lower.startsWith("x-amz-acl") ||
lower.startsWith("x-amz-grant") || lower.startsWith("x-amz-tagging") ||
lower.equals("content-md5")) unsupported("Object metadata, encryption, ACL, copy, tagging or MD5 header");
if (lower.startsWith("x-amz-checksum-") && !lower.equals("x-amz-checksum-sha256"))
unsupported("Checksum algorithm");
lower.startsWith("x-amz-copy-source-")) unsupported("Object metadata, encryption, ACL or tagging header");
}
String algorithm = SigV4.single(headers, "x-amz-sdk-checksum-algorithm");
if (algorithm != null && !algorithm.equals("SHA256")) unsupported("Checksum algorithm");
}
private void putObject(HttpExchange exchange, String key, String hash) throws IOException {
@@ -156,14 +164,51 @@ public final class Main {
try { bytes = length == null ? -1 : Long.parseLong(length); }
catch (NumberFormatException e) { throw new StoreException(400, "InvalidArgument", "Invalid Content-Length"); }
if (headers.containsKey("content-encoding")) unsupported("Encoded payload");
if (headers.containsKey("x-amz-metadata-directive")) unsupported("Copy metadata directive");
UploadChecksums checksums = UploadChecksums.from(headers);
String type = contentType(headers);
ObjectStorage.Metadata data = store.put(bucket, key, exchange.getRequestBody(), bytes, hash,
SigV4.single(headers, "x-amz-checksum-sha256"), condition != null, type);
ObjectStorage.Metadata data = store.put(bucket, key, checksums.verifying(exchange.getRequestBody()),
bytes, hash, checksums.sha256(), condition != null, type);
exchange.getResponseHeaders().set("ETag", "\"" + data.etag() + "\"");
exchange.getResponseHeaders().set("x-amz-checksum-sha256", Base64.getEncoder().encodeToString(data.sha256()));
checksums.response(exchange.getResponseHeaders());
exchange.sendResponseHeaders(200, -1);
}
private void copyObject(HttpExchange exchange, String key) throws IOException {
var headers = exchange.getRequestHeaders();
if (headers.containsKey("content-encoding") || headers.containsKey("content-md5") ||
headers.containsKey("if-none-match") || headers.containsKey("x-amz-sdk-checksum-algorithm") ||
headers.keySet().stream().anyMatch(name -> name.toLowerCase(Locale.ROOT).startsWith("x-amz-checksum-")))
unsupported("Copy request header");
String source = SigV4.single(headers, "x-amz-copy-source");
if (source == null) throw new StoreException(400, "InvalidArgument", "Missing copy source");
if (source.startsWith("/")) source = source.substring(1);
int separator = source.indexOf('/');
if (separator <= 0 || separator == source.length() - 1 || source.indexOf('?') >= 0)
throw new StoreException(400, "InvalidArgument", "Invalid copy source");
String sourceBucket = SigV4.decode(source.substring(0, separator));
String sourceKey = SigV4.decode(source.substring(separator + 1));
if (!sourceBucket.equals(bucket)) throw new StoreException(404, "NoSuchBucket", "Bucket not found");
if (sourceKey.isEmpty() || sourceKey.getBytes(StandardCharsets.UTF_8).length > 1024 ||
sourceKey.indexOf('\0') >= 0)
throw new StoreException(400, "InvalidArgument", "Invalid copy source key");
String directive = SigV4.single(headers, "x-amz-metadata-directive");
if (directive != null && !directive.equals("COPY") && !directive.equals("REPLACE"))
throw new StoreException(400, "InvalidArgument", "Invalid metadata directive");
if (!"REPLACE".equals(directive) && headers.containsKey("content-type"))
unsupported("Content-Type requires REPLACE metadata directive");
try (var object = store.open(bucket, sourceKey)) {
var sourceMetadata = object.metadata();
String type = "REPLACE".equals(directive) ? contentType(headers) : sourceMetadata.contentType();
var copied = store.put(bucket, key, object.stream(), sourceMetadata.length(),
SigV4.hex(sourceMetadata.sha256()), null, false, type);
sendXml(exchange, 200, "<CopyObjectResult><LastModified>" +
Instant.ofEpochMilli(copied.modified()) + "</LastModified><ETag>&quot;" +
copied.etag() + "&quot;</ETag></CopyObjectResult>");
}
}
private void deleteObject(HttpExchange exchange, String key) throws IOException {
if (exchange.getRequestHeaders().containsKey("if-none-match")) unsupported("Conditional delete");
store.delete(bucket, key);
@@ -212,6 +257,12 @@ public final class Main {
var headers = exchange.getRequestHeaders();
if (headers.containsKey("content-encoding") || headers.containsKey("if-none-match"))
unsupported("Multipart request header");
if (headers.containsKey("x-amz-copy-source") || headers.containsKey("x-amz-metadata-directive"))
unsupported("Multipart copy request");
if (!method.equals("PUT") && (headers.containsKey("content-md5") ||
headers.containsKey("x-amz-sdk-checksum-algorithm") ||
headers.keySet().stream().anyMatch(name -> name.toLowerCase(Locale.ROOT).startsWith("x-amz-checksum-"))))
unsupported("Multipart checksum header");
if (query.containsKey("uploads")) {
requireEmptyBody(exchange, hash);
String id = multipart.create(bucket, key, contentType(headers));
@@ -227,9 +278,11 @@ public final class Main {
try { number = Integer.parseInt(query.get("partNumber")); }
catch (NumberFormatException e) { throw new StoreException(400, "InvalidArgument", "Invalid part number"); }
long length = contentLength(headers);
String etag = multipart.putPart(id, bucket, key, number, exchange.getRequestBody(), length,
hash, SigV4.single(headers, "x-amz-checksum-sha256"));
UploadChecksums checksums = UploadChecksums.from(headers);
String etag = multipart.putPart(id, bucket, key, number,
checksums.verifying(exchange.getRequestBody()), length, hash, checksums.sha256());
exchange.getResponseHeaders().set("ETag", "\"" + etag + "\"");
checksums.response(exchange.getResponseHeaders());
exchange.sendResponseHeaders(200, -1);
}
case "POST" -> {
+56
View File
@@ -9,6 +9,7 @@ import java.net.http.HttpResponse;
import java.security.MessageDigest;
import java.time.Duration;
import java.util.HashSet;
import java.util.ArrayList;
import java.util.HexFormat;
import java.util.List;
import java.util.Set;
@@ -16,6 +17,7 @@ import java.util.UUID;
final class NodeClient {
record Node(UUID id, UUID hostId, URI url) {}
record StoredSegment(UUID id, long modified) {}
private static final HttpClient IDENTITY_HTTP = HttpClient.newBuilder()
.connectTimeout(Duration.ofSeconds(2)).build();
@@ -42,6 +44,7 @@ final class NodeClient {
int count() { return nodes.size(); }
List<Node> nodes() { return nodes; }
Node node(int index) { return nodes.get(index); }
boolean repairTokenUnavailable() { return repairToken == null || repairToken.length() < 32; }
int index(UUID id) {
for (int i = 0; i < nodes.size(); i++) {
if (nodes.get(i).id().equals(id)) return i;
@@ -143,6 +146,59 @@ final class NodeClient {
}
}
List<StoredSegment> inventory(int index, String shard, UUID after) throws IOException {
if (repairToken == null || repairToken.length() < 32)
throw new IOException("Repair authority is not available to this process");
Node node = nodes.get(index);
String path = "/segments?shard=" + shard + (after == null ? "" : "&after=" + after);
HttpRequest request = HttpRequest.newBuilder(node.url().resolve(path))
.timeout(Duration.ofSeconds(30)).header("X-Cluster-Token", token)
.header("X-Cluster-Expected-Node", node.id().toString())
.header("X-Cluster-Repair-Token", repairToken).GET().build();
HttpResponse<InputStream> response = send(request, HttpResponse.BodyHandlers.ofInputStream());
try (InputStream body = response.body()) {
if (response.statusCode() != 200) throw new IOException("Node inventory failed: " + response.statusCode());
byte[] bytes = body.readNBytes(70001);
if (bytes.length > 70000) throw new IOException("Node inventory response is too large");
List<StoredSegment> result = new ArrayList<>();
String last = after == null ? "" : after.toString();
for (String line : new String(bytes, java.nio.charset.StandardCharsets.US_ASCII).split("\n")) {
if (line.isEmpty()) continue;
String[] fields = line.split(" ", -1);
if (fields.length != 2) throw new IOException("Invalid node inventory response");
try {
UUID id = UUID.fromString(fields[0]);
if (!id.toString().equals(fields[0]) || !fields[0].startsWith(shard) ||
fields[0].compareTo(last) <= 0)
throw new IOException("Invalid node inventory cursor");
result.add(new StoredSegment(id, Long.parseLong(fields[1])));
last = fields[0];
} catch (IllegalArgumentException error) {
throw new IOException("Invalid node inventory response", error);
}
}
if (result.size() > 1000) throw new IOException("Node inventory page is too large");
return result;
}
}
boolean deleteOrphan(int index, StoredSegment segment, long minimumAgeMillis) throws IOException {
if (repairToken == null || repairToken.length() < 32)
throw new IOException("Repair authority is not available to this process");
Node node = nodes.get(index);
HttpRequest request = HttpRequest.newBuilder(node.url().resolve("/segments/" + segment.id()))
.timeout(Duration.ofSeconds(30)).header("X-Cluster-Token", token)
.header("X-Cluster-Expected-Node", node.id().toString())
.header("X-Cluster-Repair-Token", repairToken)
.header("X-Cluster-Expected-Mtime", Long.toString(segment.modified()))
.header("X-Cluster-Gc-Min-Age-Millis", Long.toString(minimumAgeMillis))
.DELETE().build();
HttpResponse<Void> response = send(request, HttpResponse.BodyHandlers.discarding());
if (response.statusCode() == 204) return true;
if (response.statusCode() == 404 || response.statusCode() == 409) return false;
throw new IOException("Node refused orphan deletion: " + response.statusCode());
}
private <T> HttpResponse<T> send(HttpRequest request, HttpResponse.BodyHandler<T> handler) throws IOException {
try { return http.send(request, handler); }
catch (InterruptedException error) {
+5 -1
View File
@@ -21,7 +21,7 @@ final class SchemaMigrator {
result.next();
version = result.getInt(1);
}
if (version > 3) throw new IOException("Metadata schema is newer than this ObjectStore build");
if (version > 4) throw new IOException("Metadata schema is newer than this ObjectStore build");
if (version < 1) {
statement.execute("CREATE TABLE IF NOT EXISTS cluster_usage (bucket text PRIMARY KEY, used_bytes bigint NOT NULL CHECK (used_bytes >= 0))");
statement.execute("CREATE TABLE IF NOT EXISTS cluster_objects (bucket text NOT NULL, object_key text COLLATE \"C\" NOT NULL, generation uuid NOT NULL, length bigint NOT NULL, modified bigint NOT NULL, etag text NOT NULL, sha256 bytea NOT NULL, content_type text NOT NULL, PRIMARY KEY (bucket, object_key))");
@@ -43,6 +43,10 @@ final class SchemaMigrator {
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)");
}
if (version < 4) {
statement.execute("CREATE TABLE cluster_gc_candidates (node_id uuid NOT NULL, segment_id uuid NOT NULL, observed_mtime bigint NOT NULL, first_seen bigint NOT NULL, PRIMARY KEY (node_id, segment_id))");
statement.execute("INSERT INTO cluster_schema_migrations VALUES (4)");
}
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")) {
@@ -0,0 +1,143 @@
package cloud.lunarsky.store;
import com.sun.net.httpserver.Headers;
import java.io.FilterInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
import java.util.Base64;
import java.util.Locale;
import java.util.zip.CRC32;
import java.util.zip.CRC32C;
import java.util.zip.Checksum;
final class UploadChecksums {
private record Algorithm(String name, String header, int length) {}
private final byte[] contentMd5;
private final Algorithm algorithm;
private final byte[] expected;
private final String encoded;
private UploadChecksums(byte[] contentMd5, Algorithm algorithm, byte[] expected, String encoded) {
this.contentMd5 = contentMd5;
this.algorithm = algorithm;
this.expected = expected;
this.encoded = encoded;
}
static UploadChecksums from(Headers headers) {
String md5 = SigV4.single(headers, "content-md5");
byte[] contentMd5 = md5 == null ? null : decode(md5, 16);
Algorithm algorithm = null;
String encoded = null;
for (String name : headers.keySet()) {
String lower = name.toLowerCase(Locale.ROOT);
if (!lower.startsWith("x-amz-checksum-")) continue;
if (algorithm != null)
throw new StoreException(400, "InvalidRequest", "Supply one checksum algorithm");
algorithm = algorithm(lower);
encoded = SigV4.single(headers, name);
}
String selected = SigV4.single(headers, "x-amz-sdk-checksum-algorithm");
if (selected != null && (algorithm == null || !selected.equals(algorithm.name())))
throw new StoreException(400, "InvalidRequest", "Checksum algorithm and value must match");
byte[] expected = algorithm == null ? null : decode(encoded, algorithm.length());
return new UploadChecksums(contentMd5, algorithm, expected, encoded);
}
private static Algorithm algorithm(String header) {
return switch (header) {
case "x-amz-checksum-crc32" -> new Algorithm("CRC32", header, 4);
case "x-amz-checksum-crc32c" -> new Algorithm("CRC32C", header, 4);
case "x-amz-checksum-sha1" -> new Algorithm("SHA1", header, 20);
case "x-amz-checksum-sha256" -> new Algorithm("SHA256", header, 32);
case "x-amz-checksum-sha512" -> new Algorithm("SHA512", header, 64);
case "x-amz-checksum-md5" -> new Algorithm("MD5", header, 16);
default -> throw new StoreException(501, "NotImplemented", "Checksum algorithm is unsupported");
};
}
private static byte[] decode(String value, int length) {
try {
byte[] decoded = Base64.getDecoder().decode(value);
if (decoded.length == length) return decoded;
} catch (IllegalArgumentException ignored) { }
throw new StoreException(400, "InvalidDigest", "Invalid checksum encoding or length");
}
String sha256() {
return algorithm != null && algorithm.name().equals("SHA256") ? encoded : null;
}
void response(Headers headers) {
if (algorithm != null) headers.set(algorithm.header(), encoded);
}
InputStream verifying(InputStream input) {
if (contentMd5 == null && (algorithm == null || algorithm.name().equals("SHA256"))) return input;
return new VerifiedInput(input);
}
private final class VerifiedInput extends FilterInputStream {
private final MessageDigest md5 = contentMd5 != null ||
(algorithm != null && algorithm.name().equals("MD5")) ? digest("MD5") : null;
private final MessageDigest hash = algorithm == null ? null : switch (algorithm.name()) {
case "SHA1" -> digest("SHA-1");
case "SHA512" -> digest("SHA-512");
default -> null;
};
private final Checksum crc = algorithm == null ? null : switch (algorithm.name()) {
case "CRC32" -> new CRC32();
case "CRC32C" -> new CRC32C();
default -> null;
};
private boolean checked;
private VerifiedInput(InputStream input) { super(input); }
@Override public int read() throws IOException {
int value = in.read();
if (value < 0) verify();
else update(new byte[]{(byte) value}, 0, 1);
return value;
}
@Override public int read(byte[] bytes, int offset, int length) throws IOException {
int count = in.read(bytes, offset, length);
if (count < 0) verify();
else if (count > 0) update(bytes, offset, count);
return count;
}
private void update(byte[] bytes, int offset, int length) {
if (md5 != null) md5.update(bytes, offset, length);
if (hash != null) hash.update(bytes, offset, length);
if (crc != null) crc.update(bytes, offset, length);
}
private void verify() {
if (checked) return;
checked = true;
byte[] actualMd5 = md5 == null ? null : md5.digest();
if (contentMd5 != null && !MessageDigest.isEqual(contentMd5, actualMd5))
throw new StoreException(400, "BadDigest", "Content-MD5 mismatch");
if (algorithm == null || algorithm.name().equals("SHA256")) return;
byte[] actual;
if (crc != null) {
long value = crc.getValue();
actual = new byte[]{(byte) (value >>> 24), (byte) (value >>> 16),
(byte) (value >>> 8), (byte) value};
} else if (algorithm.name().equals("MD5")) actual = actualMd5;
else actual = hash.digest();
if (!MessageDigest.isEqual(expected, actual))
throw new StoreException(400, "BadDigest", algorithm.name() + " checksum mismatch");
}
}
private static MessageDigest digest(String algorithm) {
try { return MessageDigest.getInstance(algorithm); }
catch (NoSuchAlgorithmException error) { throw new IllegalStateException(error); }
}
}
+1 -1
View File
@@ -1,7 +1,7 @@
package cloud.lunarsky.store;
final class Version {
static final String VALUE = "0.0.3";
static final String VALUE = "0.0.4";
private Version() {}
}