Release ObjectStore 0.0.2

This commit is contained in:
admin committed 2026-10-09 21:37:52 +02:00
1 parent 51e206f8df
commit 10e0a8061f
43 files changed
+3668 -172

No files matched your search

+154
View File
@@ -0,0 +1,154 @@
package cloud.lunarsky.store;
import java.io.DataInputStream;
import java.io.IOException;
import java.io.PrintStream;
import java.nio.channels.Channels;
import java.nio.channels.FileChannel;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.NoSuchFileException;
import java.nio.file.Path;
import java.nio.file.StandardOpenOption;
import java.security.MessageDigest;
import java.util.ArrayList;
import java.util.List;
public final class Cli {
private Cli() {}
static final class Report {
long objects, legacyObjects, payloadBytes, recordBytes, uploads, stagedBytes, checked, changedDuringScan;
long errors;
final List<String> problems = new ArrayList<>();
void problem(String message) {
errors++;
if (problems.size() < 100) problems.add(message);
}
}
public static void main(String[] args) {
int result = run(args, Path.of(System.getenv().getOrDefault("DATA_DIR", "/data")), System.out, System.err);
if (result != 0) System.exit(result);
}
static int run(String[] args, Path root, PrintStream out, PrintStream err) {
if (args.length == 1 && (args[0].equals("version") || args[0].equals("--version"))) {
out.println("ObjectStore " + Version.VALUE);
return 0;
}
if (args.length == 1 && (args[0].equals("help") || args[0].equals("--help"))) {
out.println("Usage: objectstore status|verify|version");
out.println("status Show stored object and multipart usage");
out.println("verify Check object records, paths and payload checksums");
out.println("version Show the ObjectStore version");
return 0;
}
if (args.length != 1 || !(args[0].equals("status") || args[0].equals("verify"))) {
err.println("Usage: objectstore status|verify|version");
return 2;
}
try {
boolean verify = args[0].equals("verify");
Report report = inspect(root, verify);
out.println("version=" + Version.VALUE);
out.println("objects=" + report.objects);
out.println("legacy_objects=" + report.legacyObjects);
out.println("payload_bytes=" + report.payloadBytes);
out.println("record_bytes=" + report.recordBytes);
out.println("multipart_uploads=" + report.uploads);
out.println("multipart_staged_bytes=" + report.stagedBytes);
if (verify) out.println("verified_objects=" + report.checked);
if (report.changedDuringScan > 0) out.println("changed_during_scan=" + report.changedDuringScan);
out.println("errors=" + report.errors);
for (String problem : report.problems) err.println(problem);
if (report.errors > report.problems.size())
err.println((report.errors - report.problems.size()) + " further errors omitted");
return report.errors == 0 ? 0 : 1;
} catch (IOException error) {
err.println("ObjectStore inspection failed: " + error.getMessage());
return 1;
}
}
static Report inspect(Path root, boolean verify) throws IOException {
Path objects = root.resolve("objects");
if (!Files.isDirectory(objects)) throw new IOException("Object data directory does not exist: " + objects);
Report report = new Report();
try (var paths = Files.walk(objects)) {
for (Path path : paths.filter(Files::isRegularFile).toList()) inspectObject(objects, path, verify, report);
}
Path multipart = root.resolve("multipart");
if (Files.isDirectory(multipart)) {
try (var uploads = Files.list(multipart)) {
for (Path dir : uploads.toList()) inspectUpload(dir, report);
}
}
return report;
}
private static void inspectObject(Path objects, Path path, boolean verify, Report report) {
try (var channel = FileChannel.open(path, StandardOpenOption.READ);
var input = new DataInputStream(Channels.newInputStream(channel))) {
long size = channel.size();
var record = DiskStore.readRecord(input);
var meta = record.metadata();
report.objects++;
report.payloadBytes += meta.length();
report.recordBytes += size;
if (meta.key() == null) report.legacyObjects++;
else {
String id = SigV4.hex(SigV4.hash((meta.bucket() + "/" + meta.key()).getBytes(StandardCharsets.UTF_8)));
Path expected = objects.resolve(id.substring(0, 2)).resolve(id);
if (!path.equals(expected)) report.problem("Mismatched object path: " + path);
}
if (size - record.headerLength() != meta.length()) {
report.problem("Invalid object length: " + path);
return;
}
if (verify) {
MessageDigest sha = digest("SHA-256"), md5 = digest("MD5");
byte[] buffer = new byte[65536]; long count = 0; int n;
while ((n = input.read(buffer)) != -1) {
count += n;
sha.update(buffer, 0, n);
md5.update(buffer, 0, n);
}
report.checked++;
if (count != meta.length() || !MessageDigest.isEqual(sha.digest(), meta.sha256()) ||
!SigV4.hex(md5.digest()).equals(meta.etag()))
report.problem("Object checksum mismatch: " + path);
}
} catch (NoSuchFileException error) {
report.changedDuringScan++;
} catch (IOException | RuntimeException error) {
report.problem("Unreadable object record: " + path + " (" + error.getClass().getSimpleName() + ")");
}
}
private static void inspectUpload(Path dir, Report report) throws IOException {
if (!Files.isDirectory(dir)) {
report.problem("Unexpected multipart entry: " + dir);
return;
}
report.uploads++;
if (!Files.isRegularFile(dir.resolve("manifest"))) report.problem("Missing multipart manifest: " + dir);
try (var files = Files.list(dir)) {
for (Path path : files.toList()) {
String name = path.getFileName().toString();
if (name.matches("part-[0-9]{5}")) {
try { report.stagedBytes += Files.size(path); }
catch (NoSuchFileException error) { report.changedDuringScan++; }
} else if (!name.equals("manifest")) report.problem("Unexpected multipart file: " + path);
}
} catch (NoSuchFileException error) {
report.changedDuringScan++;
}
}
private static MessageDigest digest(String name) {
try { return MessageDigest.getInstance(name); }
catch (java.security.NoSuchAlgorithmException error) { throw new IllegalStateException(error); }
}
}
+29
View File
@@ -0,0 +1,29 @@
package cloud.lunarsky.store;
import java.net.URI;
import java.sql.DriverManager;
import java.util.Map;
import java.util.UUID;
public final class ClusterJoin {
private ClusterJoin() {}
public static void main(String[] args) throws Exception {
if (args.length != 2)
throw new IllegalArgumentException("Usage: objectstore cluster-join node-url expected-host-uuid");
Map<String, String> env = System.getenv();
if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV")))
throw new IllegalArgumentException("Node registration is only enabled in local cluster mode");
URI url = URI.create(args[0]);
UUID host = UUID.fromString(args[1]);
try (var connection = DriverManager.getConnection(env.get("POSTGRES_JDBC_URL"),
env.get("POSTGRES_USER"), env.get("POSTGRES_PASSWORD"))) {
if (SchemaMigrator.prepare(connection, env.get("S3_BUCKET")) != 2)
throw new IllegalStateException("Migrate legacy replicas before joining nodes");
NodeClient.Node node = NodeRegistry.join(connection, url, host, env.get("CLUSTER_TOKEN"));
System.out.println("node_id=" + node.id());
System.out.println("host_id=" + node.hostId());
System.out.println("endpoint=" + node.url());
}
}
}
@@ -0,0 +1,157 @@
package cloud.lunarsky.store;
import java.io.IOException;
import java.net.URI;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
public final class ClusterMigrate {
private ClusterMigrate() {}
public static void main(String[] args) throws Exception {
if ((args.length != 1 || !args[0].equals("--check")) &&
(args.length != 2 || !args[0].equals("--apply")))
throw new IllegalArgumentException("Usage: objectstore cluster-migrate --check | --apply node-id-0,node-id-1,node-id-2");
Map<String, String> env = System.getenv();
if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV")))
throw new IllegalArgumentException("Migration is enabled only in local cluster mode");
List<URI> urls = Arrays.stream(env.get("CLUSTER_NODES").split(",", -1)).map(URI::create).toList();
if (urls.size() != 3) throw new IllegalArgumentException("Legacy migration requires the original three URLs in their original order");
String token = env.get("CLUSTER_TOKEN");
List<NodeClient.Node> addresses = new ArrayList<>();
for (URI url : urls) {
NodeIdentity identity = NodeClient.probe(url, token);
addresses.add(new NodeClient.Node(identity.nodeId(), identity.hostId(), url));
}
NodeClient nodes = new NodeClient(addresses, token, null);
if (args[0].equals("--apply")) {
String actual = addresses.stream().map(node -> node.id().toString())
.collect(java.util.stream.Collectors.joining(","));
if (!actual.equals(args[1]))
throw new IllegalArgumentException("Confirmed legacy node mapping differs from the current ordered node identities");
}
try (Connection connection = DriverManager.getConnection(env.get("POSTGRES_JDBC_URL"),
env.get("POSTGRES_USER"), env.get("POSTGRES_PASSWORD"))) {
int format = SchemaMigrator.prepare(connection, env.get("S3_BUCKET"));
System.out.println("cluster_format=" + format);
if (format == 2) return;
for (int index = 0; index < nodes.count(); index++) {
NodeClient.Node node = nodes.node(index);
System.out.println("legacy_" + index + "=" + node.id() + " host=" + node.hostId() +
" endpoint=" + node.url());
}
long verified = verifyLiveSegments(connection, nodes);
System.out.println("live_segments_verified=" + verified);
if (args[0].equals("--check")) return;
apply(connection, nodes);
System.out.println("cluster_format=2");
}
}
private static long verifyLiveSegments(Connection connection, NodeClient nodes) throws SQLException, IOException {
long verified = 0;
try (PreparedStatement query = connection.prepareStatement(
"SELECT s.segment_id, s.length, s.sha256, s.replicas FROM cluster_segments s JOIN cluster_objects o ON o.generation=s.generation");
ResultSet result = query.executeQuery()) {
while (result.next()) {
UUID segment = (UUID) result.getObject(1);
int length = result.getInt(2);
byte[] hash = result.getBytes(3);
for (int index : legacyIndices(result.getString(4), nodes.count())) {
try {
nodes.get(index, segment, length, hash);
} catch (IOException offlineOrCorrupt) {
throw new IOException("A listed live replica is unavailable or corrupt: segment " +
segment + " legacy node " + index, offlineOrCorrupt);
}
}
verified++;
}
}
return verified;
}
private static List<Integer> legacyIndices(String text, int count) throws IOException {
if (text == null || text.isBlank()) throw new IOException("Missing legacy replica list");
List<Integer> indices = new ArrayList<>();
Set<Integer> unique = new HashSet<>();
for (String part : text.split(",", -1)) {
int index;
try { index = Integer.parseInt(part); }
catch (NumberFormatException error) { throw new IOException("Invalid legacy replica index", error); }
if (index < 0 || index >= count || !unique.add(index))
throw new IOException("Invalid or duplicate legacy replica index");
indices.add(index);
}
return indices;
}
private static void apply(Connection connection, NodeClient nodes) throws SQLException, IOException {
try {
connection.setAutoCommit(false);
try (Statement statement = connection.createStatement()) {
statement.execute("SELECT pg_advisory_xact_lock(6834071092781)");
try (ResultSet result = statement.executeQuery("SELECT version FROM cluster_format WHERE singleton=1 FOR UPDATE")) {
if (!result.next() || result.getInt(1) != 1)
throw new IOException("Cluster format changed during migration");
}
try (ResultSet result = statement.executeQuery("SELECT EXISTS (SELECT 1 FROM cluster_nodes)")) {
result.next();
if (result.getBoolean(1)) throw new IOException("Legacy migration already has node bindings");
}
}
try (PreparedStatement insert = connection.prepareStatement(
"INSERT INTO cluster_nodes (node_id, host_id, endpoint, legacy_index, state) VALUES (?, ?, ?, ?, 'active')")) {
for (int index = 0; index < nodes.count(); index++) {
NodeClient.Node node = nodes.node(index);
insert.setObject(1, node.id());
insert.setObject(2, node.hostId());
insert.setString(3, node.url().toString());
insert.setInt(4, index);
insert.addBatch();
}
insert.executeBatch();
}
try (PreparedStatement select = connection.prepareStatement(
"SELECT generation, ordinal, replicas FROM cluster_segments WHERE replica_ids IS NULL");
ResultSet result = select.executeQuery();
PreparedStatement update = connection.prepareStatement(
"UPDATE cluster_segments SET replica_ids=? WHERE generation=? AND ordinal=? AND replica_ids IS NULL")) {
while (result.next()) {
List<UUID> ids = new ArrayList<>();
for (int index : legacyIndices(result.getString(3), nodes.count()))
ids.add(nodes.node(index).id());
update.setArray(1, connection.createArrayOf("uuid", ids.toArray()));
update.setObject(2, result.getObject(1));
update.setInt(3, result.getInt(2));
if (update.executeUpdate() != 1) throw new IOException("Segment changed during migration");
}
}
try (Statement statement = connection.createStatement()) {
try (ResultSet result = statement.executeQuery("SELECT EXISTS (SELECT 1 FROM cluster_segments WHERE replica_ids IS NULL)")) {
result.next();
if (result.getBoolean(1)) throw new IOException("Unconverted legacy segments remain");
}
statement.executeUpdate("UPDATE cluster_format SET version=2 WHERE singleton=1 AND version=1");
statement.execute("ALTER TABLE cluster_segments VALIDATE CONSTRAINT cluster_replica_ids_required");
}
connection.commit();
} 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 sql;
throw (RuntimeException) error;
} finally { connection.setAutoCommit(true); }
}
}
+244
View File
@@ -0,0 +1,244 @@
package cloud.lunarsky.store;
import com.sun.net.httpserver.HttpExchange;
import com.sun.net.httpserver.HttpServer;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.net.InetSocketAddress;
import java.nio.channels.FileChannel;
import java.nio.channels.FileLock;
import java.nio.channels.OverlappingFileLockException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.nio.file.StandardOpenOption;
import java.security.MessageDigest;
import java.util.Arrays;
import java.util.HexFormat;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.Executors;
/** Serves verified object segments to authenticated cluster peers. */
public final class ClusterNode implements AutoCloseable {
static final int MAX_SEGMENT = 8 * 1024 * 1024;
private final Path root, segments, pending;
private final byte[] token;
private final byte[] repairToken;
private final NodeIdentity identity;
private final FileChannel lockChannel;
private final FileLock lock;
ClusterNode(Path root, String token, String repairToken, UUID hostId) throws IOException {
if (token == null || token.length() < 32) throw new IllegalArgumentException("Cluster token must have at least 32 characters");
if (repairToken == null || repairToken.length() < 32 || repairToken.equals(token))
throw new IllegalArgumentException("A separate repair token of at least 32 characters is required");
if (hostId == null) throw new IllegalArgumentException("Storage host ID is required");
this.root = root;
this.token = token.getBytes(java.nio.charset.StandardCharsets.UTF_8);
this.repairToken = repairToken.getBytes(java.nio.charset.StandardCharsets.UTF_8);
segments = root.resolve("segments");
pending = root.resolve("pending");
Files.createDirectories(root);
lockChannel = FileChannel.open(root.resolve(".process.lock"), StandardOpenOption.CREATE, StandardOpenOption.WRITE);
FileLock acquired;
try { acquired = lockChannel.tryLock(); }
catch (OverlappingFileLockException error) {
lockChannel.close();
throw new IOException("Node data directory is already in use", error);
}
if (acquired == null) {
lockChannel.close();
throw new IOException("Node data directory is already in use");
}
lock = acquired;
try {
identity = NodeIdentity.open(root, hostId);
Files.createDirectories(segments);
Files.createDirectories(pending);
DiskStore.syncDirectory(root);
try (var files = Files.list(pending)) {
for (Path file : files.toList()) {
if (!Files.isRegularFile(file)) throw new IOException("Invalid pending entry: " + file);
Files.delete(file);
}
}
DiskStore.syncDirectory(pending);
} catch (IOException error) {
close();
throw error;
}
}
void handle(HttpExchange exchange) throws IOException {
try {
String path = exchange.getRequestURI().getPath();
if (path.equals("/health") && exchange.getRequestMethod().equals("GET")) {
respond(exchange, 200, "ok");
return;
}
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;
}
if (path.equals("/identity") && exchange.getRequestMethod().equals("GET")) {
respond(exchange, 200, identity.nodeId() + " " + identity.hostId());
return;
}
if (!identity.nodeId().toString().equals(exchange.getRequestHeaders().getFirst("X-Cluster-Expected-Node"))) {
respond(exchange, 409, "Wrong storage node");
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()) {
case "PUT" -> put(exchange, segmentPath(id, true));
case "GET" -> get(exchange, segmentPath(id, false));
default -> respond(exchange, 405, "Method not allowed");
}
} catch (IllegalArgumentException error) {
respond(exchange, 400, "Invalid request");
} catch (IOException error) {
if (exchange.getResponseCode() == -1) respond(exchange, 500, "Storage failure");
throw error;
} finally {
exchange.close();
}
}
private synchronized Path segmentPath(String id, boolean createShard) throws IOException {
Path shard = segments.resolve(id.substring(0, 2));
if (createShard && !Files.isDirectory(shard)) {
Files.createDirectories(shard);
DiskStore.syncDirectory(segments);
}
return shard.resolve(id);
}
private void put(HttpExchange exchange, Path target) throws IOException {
String hash = exchange.getRequestHeaders().getFirst("X-Cluster-Sha256");
String lengthText = exchange.getRequestHeaders().getFirst("Content-Length");
if (hash == null || !hash.matches("[0-9a-f]{64}") || lengthText == null) {
respond(exchange, 400, "Missing checksum or length");
return;
}
long length = Long.parseLong(lengthText);
if (length < 1 || length > MAX_SEGMENT) {
respond(exchange, 413, "Segment too large");
return;
}
Path temp = Files.createTempFile(pending, "segment-", ".part");
try {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
long count = 0;
try (InputStream input = exchange.getRequestBody(); OutputStream output = Files.newOutputStream(temp)) {
byte[] buffer = new byte[65536];
int read;
while ((read = input.read(buffer)) != -1) {
count += read;
if (count > length) {
respond(exchange, 400, "Body longer than declared length");
return;
}
digest.update(buffer, 0, read);
output.write(buffer, 0, read);
}
}
if (count != length || !hash.equals(HexFormat.of().formatHex(digest.digest()))) {
respond(exchange, 400, "Segment checksum mismatch");
return;
}
try (FileChannel channel = FileChannel.open(temp, StandardOpenOption.WRITE)) {
channel.force(true);
}
synchronized (this) {
if (Files.exists(target)) {
byte[] existing = Files.readAllBytes(target);
if (!Arrays.equals(existing, Files.readAllBytes(temp))) {
if (!"true".equals(exchange.getRequestHeaders().getFirst("X-Cluster-Repair")) ||
hash.equals(HexFormat.of().formatHex(SigV4.hash(existing)))) {
respond(exchange, 409, "Segment ID conflict");
return;
}
Files.move(temp, target, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING);
DiskStore.syncDirectory(target.getParent());
}
} else {
Files.move(temp, target, StandardCopyOption.ATOMIC_MOVE);
DiskStore.syncDirectory(target.getParent());
}
}
respond(exchange, 200, "ok");
} catch (java.security.NoSuchAlgorithmException error) {
throw new IllegalStateException(error);
} finally {
Files.deleteIfExists(temp);
}
}
private void get(HttpExchange exchange, Path target) throws IOException {
if (!Files.isRegularFile(target)) {
respond(exchange, 404, "Segment not found");
return;
}
long length = Files.size(target);
if (length > MAX_SEGMENT) throw new IOException("Segment exceeds maximum length");
exchange.sendResponseHeaders(200, length);
try (InputStream input = Files.newInputStream(target)) {
input.transferTo(exchange.getResponseBody());
}
}
private static void respond(HttpExchange exchange, int status, String message) throws IOException {
byte[] body = message.getBytes(java.nio.charset.StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type", "text/plain; charset=utf-8");
exchange.sendResponseHeaders(status, body.length);
exchange.getResponseBody().write(body);
}
@Override public void close() throws IOException {
lock.release();
lockChannel.close();
}
public static void main(String[] args) throws Exception {
Map<String, String> env = System.getenv();
String token = env.get("CLUSTER_TOKEN");
var node = new ClusterNode(Path.of(env.getOrDefault("DATA_DIR", "/data")), token,
env.get("CLUSTER_REPAIR_TOKEN"),
UUID.fromString(env.get("CLUSTER_HOST_ID")));
int port = Integer.parseInt(env.getOrDefault("NODE_PORT", "9100"));
var server = HttpServer.create(new InetSocketAddress(env.getOrDefault("NODE_BIND", "127.0.0.1"), port), 64);
var executor = Executors.newVirtualThreadPerTaskExecutor();
server.setExecutor(executor);
server.createContext("/", node::handle);
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
server.stop(5);
executor.close();
try { node.close(); } catch (IOException error) { System.err.println("Node close failed: " + error); }
}));
server.start();
System.out.println("ObjectStore cluster node listening on :" + port);
}
}
@@ -0,0 +1,26 @@
package cloud.lunarsky.store;
import java.net.URI;
import java.util.Arrays;
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");
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);
}
}
}
+442
View File
@@ -0,0 +1,442 @@
package cloud.lunarsky.store;
import java.io.ByteArrayInputStream;
import java.io.FilterInputStream;
import java.io.IOException;
import java.io.InputStream;
import java.io.OutputStream;
import java.net.URI;
import java.nio.file.Files;
import java.nio.file.Path;
import java.security.MessageDigest;
import java.security.DigestInputStream;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.time.Instant;
import java.util.ArrayList;
import java.util.Base64;
import java.util.HexFormat;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.UUID;
final class ClusterStore implements ObjectStorage {
private record Segment(UUID id, int length, byte[] hash, List<UUID> replicas) {}
private record RepairTarget(UUID generation, int ordinal, long version, Segment segment) {}
record RepairReport(int scanned, int restored, int underReplicated, int unrecoverable) {}
private final String jdbcUrl, user, password, configuredBucket;
private final NodeClient nodes;
private final long maxObject, maxTotal;
private final boolean testNodeDomains;
ClusterStore(String jdbcUrl, String user, String password, String bucket,
List<URI> nodeUrls, String token, String repairToken, long maxObject, long maxTotal,
boolean testNodeDomains) throws IOException {
if (jdbcUrl == null || !jdbcUrl.startsWith("jdbc:postgresql://") || user == null || password == null)
throw new IllegalArgumentException("Invalid metadata database configuration");
this.jdbcUrl = jdbcUrl; this.user = user; this.password = password;
this.configuredBucket = bucket; this.maxObject = maxObject; this.maxTotal = maxTotal;
this.testNodeDomains = testNodeDomains;
try (Connection connection = connect()) {
int format = SchemaMigrator.prepare(connection, bucket);
if (format != 2) throw new IOException("Legacy replica positions require objectstore cluster-migrate before this gateway can start");
nodes = NodeRegistry.load(connection, nodeUrls, token, repairToken);
} catch (SQLException error) { throw databaseError(error); }
}
private Connection connect() throws SQLException { return DriverManager.getConnection(jdbcUrl, user, password); }
@Override public Metadata put(String bucket, String key, InputStream input, long length, String expectedHash,
String checksum, boolean createOnly, String contentType) throws IOException {
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 > maxObject) throw new StoreException(413, "EntityTooLarge", "Object exceeds the configured size limit");
if (!nodes.availableHostsAtLeast(2, testNodeDomains))
throw new StoreException(503, "SlowDown", "Fewer than two storage hosts are available");
if (contentType.getBytes(java.nio.charset.StandardCharsets.UTF_8).length > 255)
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;
Path staged = Files.createTempFile("objectstore-cluster-", ".pending");
try {
try (OutputStream output = Files.newOutputStream(staged)) {
byte[] buffer = new byte[65536];
long remaining = length;
while (remaining > 0) {
int count = input.read(buffer, 0, (int) Math.min(buffer.length, remaining));
if (count < 0) throw new StoreException(400, "IncompleteBody", "Payload length does not match Content-Length");
if (count == 0) continue;
sha.update(buffer, 0, count); md5.update(buffer, 0, count);
output.write(buffer, 0, count);
remaining -= count;
}
}
if (input.read() != -1) throw new StoreException(413, "EntityTooLarge", "Payload exceeds declared size");
fullHash = sha.digest();
if (!HexFormat.of().formatHex(fullHash).equals(expectedHash))
throw new StoreException(400, "XAmzContentSHA256Mismatch", "Payload hash mismatch");
if (checksum != null && !Base64.getEncoder().encodeToString(fullHash).equals(checksum))
throw new StoreException(400, "BadDigest", "SHA-256 checksum mismatch");
try (Connection connection = connect()) {
long previous = currentLength(connection, bucket, key);
if (createOnly && previous >= 0)
throw new StoreException(412, "PreconditionFailed", "Object already exists");
try (PreparedStatement query = connection.prepareStatement("SELECT used_bytes FROM cluster_usage WHERE bucket=?")) {
query.setString(1, bucket);
try (ResultSet result = query.executeQuery()) {
if (!result.next()) throw new SQLException("Bucket quota row is missing");
if (result.getLong(1) - Math.max(0, previous) > maxTotal - length)
throw new StoreException(507, "InsufficientStorage", "Store capacity limit reached");
}
}
} catch (SQLException error) { throw databaseError(error); }
try (InputStream stagedInput = Files.newInputStream(staged)) {
long remaining = length;
while (remaining > 0) {
int wanted = (int) Math.min(ClusterNode.MAX_SEGMENT, remaining);
byte[] bytes = stagedInput.readNBytes(wanted);
if (bytes.length != wanted) throw new IOException("Staged object was truncated");
byte[] segmentHash = SigV4.hash(bytes);
UUID id = UUID.randomUUID();
List<UUID> replicas = new ArrayList<>();
Set<UUID> acceptedHosts = new HashSet<>();
for (int index : PlacementPolicy.candidates(id, nodes, testNodeDomains)) {
UUID host = nodes.faultDomain(index, testNodeDomains);
if (acceptedHosts.contains(host)) continue;
try {
nodes.put(index, id, bytes, segmentHash);
replicas.add(nodes.node(index).id());
acceptedHosts.add(host);
if (acceptedHosts.size() == 3) break;
} catch (IOException error) {
System.err.println("Cluster node " + nodes.node(index).id() +
" did not accept segment " + id + ": " + error.getMessage());
}
}
if (acceptedHosts.size() < 2)
throw new StoreException(503, "SlowDown", "Fewer than two storage hosts accepted the segment");
segments.add(new Segment(id, wanted, segmentHash, List.copyOf(replicas)));
remaining -= wanted;
}
}
} finally { Files.deleteIfExists(staged); }
Metadata metadata = new Metadata(length, Instant.now().toEpochMilli(),
HexFormat.of().formatHex(md5.digest()), fullHash, bucket, key, contentType);
UUID generation = UUID.randomUUID();
try (Connection connection = connect()) {
connection.setAutoCommit(false);
try {
long used = lockUsage(connection, bucket);
long previous = currentLength(connection, bucket, key);
if (createOnly && previous >= 0)
throw new StoreException(412, "PreconditionFailed", "Object already exists");
if (used - Math.max(0, previous) > maxTotal - length)
throw new StoreException(507, "InsufficientStorage", "Store capacity limit reached");
try (PreparedStatement insert = connection.prepareStatement(
"INSERT INTO cluster_segments (generation, ordinal, segment_id, length, sha256, replicas, replica_ids) VALUES (?, ?, ?, ?, ?, 'v2', ?)")) {
for (int i = 0; i < segments.size(); i++) {
Segment segment = segments.get(i);
insert.setObject(1, generation); insert.setInt(2, i); insert.setObject(3, segment.id());
insert.setInt(4, segment.length()); insert.setBytes(5, segment.hash());
insert.setArray(6, connection.createArrayOf("uuid", segment.replicas().toArray()));
insert.addBatch();
}
insert.executeBatch();
}
try (PreparedStatement update = connection.prepareStatement(
"INSERT INTO cluster_objects VALUES (?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT (bucket, object_key) DO UPDATE SET generation=EXCLUDED.generation, length=EXCLUDED.length, modified=EXCLUDED.modified, etag=EXCLUDED.etag, sha256=EXCLUDED.sha256, content_type=EXCLUDED.content_type")) {
bindObject(update, metadata, generation); update.executeUpdate();
}
try (PreparedStatement update = connection.prepareStatement("UPDATE cluster_usage SET used_bytes=? WHERE bucket=?")) {
update.setLong(1, used - Math.max(0, previous) + length); update.setString(2, bucket); update.executeUpdate();
}
try (PreparedStatement delete = connection.prepareStatement("DELETE FROM cluster_tombstones WHERE bucket=? AND object_key=?")) {
delete.setString(1, bucket); delete.setString(2, key); delete.executeUpdate();
}
connection.commit();
return metadata;
} catch (SQLException | RuntimeException error) {
connection.rollback();
if (error instanceof SQLException sql) throw databaseError(sql);
throw error;
}
} catch (SQLException error) { throw databaseError(error); }
}
@Override public OpenObject open(String bucket, String key) throws IOException {
try (Connection connection = connect()) {
connection.setAutoCommit(false);
connection.setTransactionIsolation(Connection.TRANSACTION_REPEATABLE_READ);
try {
Metadata metadata;
UUID generation;
try (PreparedStatement query = connection.prepareStatement(
"SELECT generation, length, modified, etag, sha256, content_type FROM cluster_objects WHERE bucket=? AND object_key=?")) {
query.setString(1, bucket); query.setString(2, key);
try (ResultSet result = query.executeQuery()) {
if (!result.next()) throw new StoreException(404, "NoSuchKey", "Object not found");
generation = (UUID) result.getObject(1);
metadata = new Metadata(result.getLong(2), result.getLong(3), result.getString(4),
result.getBytes(5), bucket, key, result.getString(6));
}
}
List<Segment> parts = new ArrayList<>();
try (PreparedStatement query = connection.prepareStatement(
"SELECT segment_id, length, sha256, replica_ids FROM cluster_segments WHERE generation=? ORDER BY ordinal")) {
query.setObject(1, generation);
try (ResultSet result = query.executeQuery()) {
long total = 0;
while (result.next()) {
Segment segment = new Segment((UUID) result.getObject(1), result.getInt(2),
result.getBytes(3), replicaIds(result, 4));
total = Math.addExact(total, segment.length());
parts.add(segment);
}
if (total != metadata.length()) throw new IOException("Incomplete object manifest");
}
}
connection.commit();
return new OpenObject(metadata, verifiedObject(parts, metadata));
} catch (SQLException | RuntimeException | IOException error) {
connection.rollback();
if (error instanceof SQLException sql) throw databaseError(sql);
if (error instanceof IOException io) throw io;
throw error;
}
} catch (SQLException error) { throw databaseError(error); }
}
@Override public void delete(String bucket, String key) throws IOException {
try (Connection connection = connect()) {
connection.setAutoCommit(false);
try {
long used = lockUsage(connection, bucket);
long previous = currentLength(connection, bucket, key);
try (PreparedStatement delete = connection.prepareStatement("DELETE FROM cluster_objects WHERE bucket=? AND object_key=?")) {
delete.setString(1, bucket); delete.setString(2, key); delete.executeUpdate();
}
try (PreparedStatement update = connection.prepareStatement(
"INSERT INTO cluster_tombstones VALUES (?, ?, ?, ?) ON CONFLICT (bucket, object_key) DO UPDATE SET generation=EXCLUDED.generation, deleted_at=EXCLUDED.deleted_at")) {
update.setString(1, bucket); update.setString(2, key);
update.setObject(3, UUID.randomUUID()); update.setLong(4, Instant.now().toEpochMilli());
update.executeUpdate();
}
if (previous >= 0) {
try (PreparedStatement update = connection.prepareStatement("UPDATE cluster_usage SET used_bytes=? WHERE bucket=?")) {
update.setLong(1, used - previous); update.setString(2, bucket); update.executeUpdate();
}
}
connection.commit();
} catch (SQLException | RuntimeException error) {
connection.rollback();
if (error instanceof SQLException sql) throw databaseError(sql);
throw error;
}
} catch (SQLException error) { throw databaseError(error); }
}
@Override public ListPage list(String bucket, String prefix, String delimiter, int maxKeys, String after) throws IOException {
List<ListedObject> entries = new ArrayList<>();
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()) {
connection.setAutoCommit(false);
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")) {
query.setString(1, bucket);
query.setString(2, after != null && after.compareTo(prefix) > 0 ? after : prefix);
query.setFetchSize(128);
try (ResultSet result = query.executeQuery()) {
while (result.next()) {
String key = result.getString(1);
if (!key.startsWith(prefix)) break;
if (after != null && key.compareTo(after) <= 0) continue;
String group = null;
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 (entries.size() + prefixes.size() >= maxKeys) { truncated = true; break; }
if (group != null) { prefixes.add(group); activePrefix = group; }
else {
entries.add(new ListedObject(key, new Metadata(result.getLong(2), result.getLong(3),
result.getString(4), result.getBytes(5), bucket, key, result.getString(6))));
activePrefix = null;
}
lastKey = key;
}
}
}
connection.commit();
} catch (SQLException error) { throw databaseError(error); }
return new ListPage(entries, prefixes, truncated ? lastKey : null, truncated);
}
private long lockUsage(Connection connection, String bucket) throws SQLException {
try (PreparedStatement query = connection.prepareStatement("SELECT used_bytes FROM cluster_usage WHERE bucket=? FOR UPDATE")) {
query.setString(1, bucket);
try (ResultSet result = query.executeQuery()) {
if (!result.next()) throw new SQLException("Bucket quota row is missing");
return result.getLong(1);
}
}
}
private long currentLength(Connection connection, String bucket, String key) throws SQLException {
try (PreparedStatement query = connection.prepareStatement("SELECT length FROM cluster_objects WHERE bucket=? AND object_key=?")) {
query.setString(1, bucket); query.setString(2, key);
try (ResultSet result = query.executeQuery()) { return result.next() ? result.getLong(1) : -1; }
}
}
private static void bindObject(PreparedStatement update, Metadata data, UUID generation) throws SQLException {
update.setString(1, data.bucket()); update.setString(2, data.key()); update.setObject(3, generation);
update.setLong(4, data.length()); update.setLong(5, data.modified()); update.setString(6, data.etag());
update.setBytes(7, data.sha256()); update.setString(8, data.contentType());
}
private static List<UUID> replicaIds(ResultSet result, int column) throws SQLException, IOException {
java.sql.Array value = result.getArray(column);
if (value == null) throw new IOException("Segment has no migrated replica identities");
try {
Object[] ids = (Object[]) value.getArray();
List<UUID> replicas = new ArrayList<>(ids.length);
for (Object id : ids) replicas.add((UUID) id);
return List.copyOf(replicas);
} finally { value.free(); }
}
private static IOException databaseError(SQLException error) { return new IOException("Metadata database operation failed", error); }
private static MessageDigest digest(String algorithm) {
try { return MessageDigest.getInstance(algorithm); }
catch (java.security.NoSuchAlgorithmException error) { throw new IllegalStateException(error); }
}
private InputStream verifiedObject(List<Segment> segments, Metadata metadata) throws IOException {
Path staged = Files.createTempFile("objectstore-read-", ".pending");
boolean ready = false;
try {
MessageDigest hash = digest("SHA-256");
long count;
try (InputStream source = new DigestInputStream(new SegmentStream(segments), hash);
OutputStream output = Files.newOutputStream(staged)) {
count = source.transferTo(output);
}
if (count != metadata.length() || !MessageDigest.isEqual(hash.digest(), metadata.sha256()))
throw new IOException("Object manifest failed integrity verification");
InputStream file = Files.newInputStream(staged);
ready = true;
return new FilterInputStream(file) {
@Override public void close() throws IOException {
try { super.close(); }
finally { Files.deleteIfExists(staged); }
}
};
} finally { if (!ready) Files.deleteIfExists(staged); }
}
@Override public boolean ready() {
if (!nodes.availableHostsAtLeast(2, testNodeDomains)) return false;
try (Connection connection = connect(); var statement = connection.createStatement();
ResultSet result = statement.executeQuery("SELECT 1")) {
return result.next() && result.getInt(1) == 1;
} catch (SQLException error) { return false; }
}
RepairReport repairOnce() throws IOException {
int scanned = 0, restored = 0, underReplicated = 0, unrecoverable = 0;
try (Connection reader = connect()) {
reader.setAutoCommit(false);
try (PreparedStatement query = reader.prepareStatement(
"SELECT s.generation, s.ordinal, s.segment_id, s.length, s.sha256, s.replica_ids, s.placement_version FROM cluster_segments s JOIN cluster_objects o ON o.generation=s.generation ORDER BY s.generation, s.ordinal")) {
query.setFetchSize(128);
try (ResultSet result = query.executeQuery()) {
while (result.next()) {
scanned++;
RepairTarget target = new RepairTarget((UUID) result.getObject(1), result.getInt(2),
result.getLong(7), new Segment((UUID) result.getObject(3), result.getInt(4),
result.getBytes(5), replicaIds(result, 6)));
Segment segment = target.segment();
byte[] copy = null;
Set<UUID> healthy = new HashSet<>();
Set<UUID> healthyHosts = new HashSet<>();
for (UUID id : segment.replicas()) {
int node = nodes.index(id);
if (node < 0) continue;
try {
byte[] candidate = nodes.get(node, segment.id(), segment.length(), segment.hash());
if (copy == null) copy = candidate;
healthy.add(id);
healthyHosts.add(nodes.faultDomain(node, testNodeDomains));
} catch (IOException error) { }
}
if (copy == null) { unrecoverable++; continue; }
for (int node : PlacementPolicy.candidates(segment.id(), nodes, testNodeDomains)) {
UUID host = nodes.faultDomain(node, testNodeDomains);
if (healthyHosts.contains(host)) continue;
try {
nodes.repair(node, segment.id(), copy, segment.hash());
healthy.add(nodes.node(node).id());
healthyHosts.add(host);
restored++;
} catch (IOException error) { }
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()) {
try (Connection writer = connect(); PreparedStatement update = writer.prepareStatement(
"UPDATE cluster_segments SET replica_ids=?, placement_version=placement_version+1 WHERE generation=? AND ordinal=? AND placement_version=?")) {
update.setArray(1, writer.createArrayOf("uuid", listed.toArray()));
update.setObject(2, target.generation());
update.setInt(3, target.ordinal());
update.setLong(4, target.version());
update.executeUpdate();
}
}
}
}
}
reader.commit();
} catch (SQLException error) { throw databaseError(error); }
return new RepairReport(scanned, restored, underReplicated, unrecoverable);
}
@Override public void close() {}
private final class SegmentStream extends InputStream {
private final List<Segment> segments;
private int position;
private ByteArrayInputStream current;
private boolean closed;
SegmentStream(List<Segment> segments) { this.segments = segments; }
@Override public int read() throws IOException {
byte[] one = new byte[1];
int count = read(one, 0, 1);
return count < 0 ? -1 : one[0] & 255;
}
@Override public int read(byte[] buffer, int offset, int length) throws IOException {
if (closed) throw new IOException("Object stream is closed");
if (length == 0) return 0;
while (current == null || current.available() == 0) {
if (position == segments.size()) return -1;
Segment segment = segments.get(position++);
IOException failure = null;
for (UUID replica : segment.replicas()) {
int node = nodes.index(replica);
if (node < 0) continue;
try {
byte[] bytes = nodes.get(node, segment.id(), segment.length(), segment.hash());
current = new ByteArrayInputStream(bytes);
break;
} catch (IOException error) { failure = error; }
}
if (current == null || current.available() == 0)
throw new StoreException(503, "SlowDown", "No verified replica is currently available", failure);
}
return current.read(buffer, offset, length);
}
@Override public void close() { closed = true; current = null; }
}
}
+206 -73
View File
@@ -5,28 +5,35 @@ import java.nio.ByteBuffer;
import java.nio.channels.FileChannel;
import java.nio.channels.FileLock;
import java.nio.channels.OverlappingFileLockException;
import java.nio.charset.CodingErrorAction;
import java.nio.charset.StandardCharsets;
import java.nio.file.*;
import java.security.MessageDigest;
import java.time.Instant;
import java.util.Arrays;
import java.util.HexFormat;
import java.util.*;
import cloud.lunarsky.store.ObjectStorage.Metadata;
import cloud.lunarsky.store.ObjectStorage.OpenObject;
import cloud.lunarsky.store.ObjectStorage.ListedObject;
import cloud.lunarsky.store.ObjectStorage.ListPage;
final class DiskStore implements AutoCloseable {
private static final long MAGIC = 0x4c534f424a303031L;
private static final int HEADER = 72;
private final Path objects, temporary;
final class DiskStore implements ObjectStorage {
private static final long MAGIC_V1 = 0x4c534f424a303031L;
private static final long MAGIC_V2 = 0x4c534f424a303032L;
private static final int HEADER_V1 = 72;
private static final int HEADER_V2 = 78;
private final Path root, objects, temporary;
private final FileChannel lockChannel;
private final FileLock processLock;
private final long maxObject, maxTotal;
private final Object[] locks = new Object[128];
private final NavigableMap<String, Metadata> index = new TreeMap<>();
private long used;
record Metadata(long length, long modified, String etag, byte[] sha256) {}
record OpenObject(Metadata metadata, InputStream stream) implements AutoCloseable {
public void close() throws IOException { stream.close(); }
}
private long objectCount, legacyCount;
record Record(Metadata metadata, int headerLength) {}
DiskStore(Path root, long maxObject, long maxTotal) throws IOException {
this.root = root;
objects = root.resolve("objects"); temporary = root.resolve("pending");
this.maxObject = maxObject; this.maxTotal = maxTotal;
Arrays.setAll(locks, i -> new Object());
@@ -39,16 +46,25 @@ final class DiskStore implements AutoCloseable {
catch (OverlappingFileLockException e) { throw new IOException("Data directory is already in use", e); }
if (acquired == null) throw new IOException("Data directory is already in use");
Files.createDirectories(objects); Files.createDirectories(temporary);
try (var paths=Files.list(temporary)) {
for(Path p:paths.toList()) if(p.getFileName().toString().endsWith(".part"))Files.delete(p);
syncDirectory(root);
try (var paths = Files.list(temporary)) {
for (Path p : paths.toList()) if (p.getFileName().toString().endsWith(".part")) Files.delete(p);
}
try (var paths=Files.walk(objects)) {
for(Path p:paths.filter(Files::isRegularFile).toList()) {
try(var in=new DataInputStream(Files.newInputStream(p))) {
long length = metadata(in).length();
if(Files.size(p)-HEADER != length) throw new IOException("Truncated or oversized object record: "+p);
used = Math.addExact(used, length);
}
try (var paths = Files.walk(objects)) {
for (Path p : paths.filter(Files::isRegularFile).toList()) {
Record record;
try (var in = new DataInputStream(Files.newInputStream(p))) { record = readRecord(in); }
Metadata meta = record.metadata();
if (Files.size(p) - record.headerLength() != meta.length())
throw new IOException("Truncated or oversized object record: " + p);
if (meta.key() != null) {
if (!p.equals(objectPath(meta.bucket(), meta.key())))
throw new IOException("Mismatched object record: " + p);
if (index.put(indexKey(meta.bucket(), meta.key()), meta) != null)
throw new IOException("Duplicate object record: " + p);
} else legacyCount++;
objectCount++;
used = Math.addExact(used, meta.length());
}
}
ready = true;
@@ -66,72 +82,189 @@ final class DiskStore implements AutoCloseable {
processLock.release();
lockChannel.close();
}
Path root() { return root; }
long maxObject() { return maxObject; }
long maxTotal() { return maxTotal; }
synchronized long usedBytes() { return used; }
synchronized int indexedObjects() { return index.size(); }
synchronized long objectCount() { return objectCount; }
synchronized long legacyObjects() { return legacyCount; }
private Path object(String bucket, String key) throws IOException {
String id=SigV4.hex(SigV4.hash((bucket+"/"+key).getBytes(StandardCharsets.UTF_8)));
Path shard=objects.resolve(id.substring(0,2));Files.createDirectories(shard);
return shard.resolve(id);
private static String indexKey(String bucket, String key) { return bucket + "\0" + key; }
private Path objectPath(String bucket, String key) {
String id = SigV4.hex(SigV4.hash((bucket + "/" + key).getBytes(StandardCharsets.UTF_8)));
return objects.resolve(id.substring(0, 2)).resolve(id);
}
private Object lock(Path p){return locks[(p.hashCode()&0x7fffffff)%locks.length];}
private synchronized Path object(String bucket, String key) throws IOException {
Path path = objectPath(bucket, key);
if (!Files.isDirectory(path.getParent())) {
Files.createDirectories(path.getParent());
syncDirectory(objects);
}
return path;
}
static void syncDirectory(Path directory) throws IOException {
try (FileChannel channel = FileChannel.open(directory, StandardOpenOption.READ)) {
channel.force(true);
}
}
private Object lock(Path p) { return locks[(p.hashCode() & 0x7fffffff) % locks.length]; }
Metadata put(String bucket,String key,InputStream input,long length,String expectedHash,String checksum,boolean createOnly) throws IOException {
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");
Path destination=object(bucket,key),pending=Files.createTempFile(temporary,"upload-",".part");
public Metadata put(String bucket, String key, InputStream input, long length, String expectedHash,
String checksum, boolean createOnly, String contentType) throws IOException {
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");
byte[] bucketBytes = bucket.getBytes(StandardCharsets.UTF_8);
byte[] keyBytes = key.getBytes(StandardCharsets.UTF_8);
byte[] typeBytes = contentType.getBytes(StandardCharsets.UTF_8);
if (bucketBytes.length > 63 || keyBytes.length > 1024 || typeBytes.length > 255)
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");
try {
MessageDigest sha=digest("SHA-256"),md5=digest("MD5");
long count=0;
try(OutputStream out=Files.newOutputStream(pending)){
out.write(new byte[HEADER]);byte[] buffer=new byte[65536];int n;
while((n=input.read(buffer))!=-1){count+=n;if(count>length||count>maxObject)throw new StoreException(413,"EntityTooLarge","Payload exceeds declared size");sha.update(buffer,0,n);md5.update(buffer,0,n);out.write(buffer,0,n);}
}
if(count!=length)throw new StoreException(400,"IncompleteBody","Payload length does not match Content-Length");
byte[] hash=sha.digest(),etag=md5.digest();
if(!MessageDigest.isEqual(hash,HexFormat.of().parseHex(expectedHash)))throw new StoreException(400,"XAmzContentSHA256Mismatch","Payload hash mismatch");
if(checksum!=null&&!java.util.Base64.getEncoder().encodeToString(hash).equals(checksum))throw new StoreException(400,"BadDigest","SHA-256 checksum mismatch");
long modified=Instant.now().toEpochMilli();
try(FileChannel file=FileChannel.open(pending,StandardOpenOption.WRITE)){
ByteBuffer header=ByteBuffer.allocate(HEADER).putLong(MAGIC).putLong(count).putLong(modified).put(etag).put(hash);header.flip();
while(header.hasRemaining())file.write(header);file.force(true);
}
synchronized(lock(destination)){
long previous=0;
if(Files.exists(destination)){
if(createOnly)throw new StoreException(412,"PreconditionFailed","Object already exists");
try(var in=new DataInputStream(Files.newInputStream(destination))){previous=metadata(in).length();}
}
synchronized(this){
if(used-previous+count>maxTotal)throw new StoreException(507,"InsufficientStorage","Store capacity limit reached");
Files.move(pending,destination,StandardCopyOption.ATOMIC_MOVE,StandardCopyOption.REPLACE_EXISTING);
used=used-previous+count;
MessageDigest sha = digest("SHA-256"), md5 = digest("MD5");
long count = 0;
try (OutputStream out = Files.newOutputStream(pending)) {
out.write(new byte[headerLength]);
byte[] buffer = new byte[65536]; int n;
while ((n = input.read(buffer)) != -1) {
count += n;
if (count > length || count > maxObject)
throw new StoreException(413, "EntityTooLarge", "Payload exceeds declared size");
sha.update(buffer, 0, n); md5.update(buffer, 0, n); out.write(buffer, 0, n);
}
}
return new Metadata(count,modified,SigV4.hex(etag),hash);
} finally {Files.deleteIfExists(pending);}
if (count != length) throw new StoreException(400, "IncompleteBody", "Payload length does not match Content-Length");
byte[] hash = sha.digest(), etag = md5.digest();
if (!MessageDigest.isEqual(hash, HexFormat.of().parseHex(expectedHash)))
throw new StoreException(400, "XAmzContentSHA256Mismatch", "Payload hash mismatch");
if (checksum != null && !Base64.getEncoder().encodeToString(hash).equals(checksum))
throw new StoreException(400, "BadDigest", "SHA-256 checksum mismatch");
long modified = Instant.now().toEpochMilli();
ByteBuffer header = ByteBuffer.allocate(headerLength).putLong(MAGIC_V2).putLong(count)
.putLong(modified).put(etag).put(hash).putShort((short) bucketBytes.length)
.putShort((short) keyBytes.length).putShort((short) typeBytes.length)
.put(bucketBytes).put(keyBytes).put(typeBytes);
header.flip();
try (FileChannel file = FileChannel.open(pending, StandardOpenOption.WRITE)) {
while (header.hasRemaining()) file.write(header, header.position());
file.force(true);
}
Metadata metadata = new Metadata(count, modified, SigV4.hex(etag), hash, bucket, key, contentType);
synchronized (lock(destination)) {
long previous = 0;
boolean existed = Files.exists(destination);
boolean legacy = false;
if (existed) {
if (createOnly) throw new StoreException(412, "PreconditionFailed", "Object already exists");
try (var in = new DataInputStream(Files.newInputStream(destination))) {
Metadata old = readRecord(in).metadata();
previous = old.length();
legacy = old.key() == null;
}
}
synchronized (this) {
if (used - previous + count > maxTotal)
throw new StoreException(507, "InsufficientStorage", "Store capacity limit reached");
Files.move(pending, destination, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING);
used = used - previous + count;
if (!existed) objectCount++;
if (legacy) legacyCount--;
index.put(indexKey(bucket, key), metadata);
syncDirectory(destination.getParent());
}
}
return metadata;
} finally { Files.deleteIfExists(pending); }
}
OpenObject open(String bucket,String key) throws IOException {
Path destination=object(bucket,key);
synchronized(lock(destination)){
public OpenObject open(String bucket, String key) throws IOException {
Path destination = object(bucket, key);
synchronized (lock(destination)) {
final DataInputStream input;
try{input=new DataInputStream(Files.newInputStream(destination));}
catch(NoSuchFileException e){throw new StoreException(404,"NoSuchKey","Object not found");}
try{return new OpenObject(metadata(input),input);}catch(IOException e){input.close();throw e;}
try { input = new DataInputStream(Files.newInputStream(destination)); }
catch (NoSuchFileException e) { throw new StoreException(404, "NoSuchKey", "Object not found"); }
try { return new OpenObject(readRecord(input).metadata(), input); }
catch (IOException e) { input.close(); throw e; }
}
}
void delete(String bucket,String key) throws IOException {
Path destination=object(bucket,key);
synchronized(lock(destination)){
if(!Files.exists(destination))return;
long length;try(var input=new DataInputStream(Files.newInputStream(destination))){length=metadata(input).length();}
synchronized(this){Files.delete(destination);used-=length;}
public void delete(String bucket, String key) throws IOException {
Path destination = object(bucket, key);
synchronized (lock(destination)) {
if (!Files.exists(destination)) return;
long length;
boolean legacy;
try (var input = new DataInputStream(Files.newInputStream(destination))) {
Metadata old = readRecord(input).metadata();
length = old.length();
legacy = old.key() == null;
}
synchronized (this) {
Files.delete(destination);
used -= length;
objectCount--;
if (legacy) legacyCount--;
index.remove(indexKey(bucket, key));
syncDirectory(destination.getParent());
}
}
}
private static Metadata metadata(DataInputStream in) throws IOException {
if(in.readLong()!=MAGIC)throw new IOException("Invalid object record");
long length=in.readLong(),modified=in.readLong();byte[] md5=new byte[16],sha=new byte[32];in.readFully(md5);in.readFully(sha);
if(length<0)throw new IOException("Invalid object length");
return new Metadata(length,modified,SigV4.hex(md5),sha);
public synchronized ListPage list(String bucket, String prefix, String delimiter, int maxKeys, String after) {
List<ListedObject> entries = new ArrayList<>();
List<String> prefixes = new ArrayList<>();
if (maxKeys == 0) return new ListPage(entries, prefixes, null, false);
String lastKey = null;
boolean truncated = false;
String activePrefix = null;
for (Metadata meta : index.values()) {
if (!meta.bucket().equals(bucket) || !meta.key().startsWith(prefix)) continue;
String key = meta.key();
if (after != null && key.compareTo(after) <= 0) continue;
String group = null;
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 (entries.size() + prefixes.size() >= maxKeys) { truncated = true; break; }
if (group != null) { prefixes.add(group); activePrefix = group; }
else { entries.add(new ListedObject(key, meta)); activePrefix = null; }
lastKey = key;
}
return new ListPage(entries, prefixes, truncated ? lastKey : null, truncated);
}
static Record readRecord(DataInputStream in) throws IOException {
long magic = in.readLong();
if (magic != MAGIC_V1 && magic != MAGIC_V2) throw new IOException("Invalid object record");
long length = in.readLong(), modified = in.readLong();
byte[] md5 = new byte[16], sha = new byte[32];
in.readFully(md5); in.readFully(sha);
if (length < 0) throw new IOException("Invalid object record length");
if (magic == MAGIC_V1)
return new Record(new Metadata(length, modified, SigV4.hex(md5), sha,
null, null, "application/octet-stream"), HEADER_V1);
int bucketLength = in.readUnsignedShort(), keyLength = in.readUnsignedShort(), typeLength = in.readUnsignedShort();
if (bucketLength < 1 || bucketLength > 63 || keyLength < 1 || keyLength > 1024 || typeLength < 1 || typeLength > 255)
throw new IOException("Invalid object record metadata");
String bucket = utf8(in.readNBytes(bucketLength));
String key = utf8(in.readNBytes(keyLength));
String contentType = utf8(in.readNBytes(typeLength));
if (bucket.getBytes(StandardCharsets.UTF_8).length != bucketLength ||
key.getBytes(StandardCharsets.UTF_8).length != keyLength ||
contentType.getBytes(StandardCharsets.UTF_8).length != typeLength)
throw new IOException("Invalid object record metadata");
return new Record(new Metadata(length, modified, SigV4.hex(md5), sha,
bucket, key, contentType), HEADER_V2 + bucketLength + keyLength + typeLength);
}
private static String utf8(byte[] bytes) throws IOException {
return StandardCharsets.UTF_8.newDecoder().onMalformedInput(CodingErrorAction.REPORT)
.decode(ByteBuffer.wrap(bytes)).toString();
}
private static MessageDigest digest(String algorithm) {
try { return MessageDigest.getInstance(algorithm); }
catch (java.security.NoSuchAlgorithmException e) { throw new IllegalStateException(e); }
}
private static MessageDigest digest(String algorithm){try{return MessageDigest.getInstance(algorithm);}catch(java.security.NoSuchAlgorithmException e){throw new IllegalStateException(e);}}
}
+462 -76
View File
@@ -3,112 +3,498 @@ package cloud.lunarsky.store;
import com.sun.net.httpserver.HttpExchange;
import com.sun.net.httpserver.HttpServer;
import java.io.IOException;
import java.io.ByteArrayInputStream;
import java.net.InetSocketAddress;
import java.net.URI;
import java.nio.charset.StandardCharsets;
import java.nio.file.Path;
import java.time.Clock;
import java.time.Instant;
import java.time.ZoneOffset;
import java.time.format.DateTimeFormatter;
import java.util.Base64;
import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
import java.util.ArrayList;
import java.util.List;
import java.util.Locale;
import java.util.concurrent.Executors;
import java.util.concurrent.Semaphore;
public final class Main {
private final DiskStore store;
private final ObjectStorage store;
private final SigV4 authentication;
private final String bucket;
private final Semaphore slots=new Semaphore(16);
Main(DiskStore store,SigV4 authentication,String bucket){this.store=store;this.authentication=authentication;this.bucket=bucket;}
private final MultipartStorage multipart;
private final Semaphore slots = new Semaphore(16);
Main(DiskStore store, SigV4 authentication, String bucket) throws IOException {
this(store, new MultipartStore(store), authentication, bucket);
}
Main(ObjectStorage store, MultipartStorage multipart, SigV4 authentication, String bucket) {
this.store = store; this.multipart = multipart;
this.authentication = authentication; this.bucket = bucket;
}
void handle(HttpExchange exchange) throws IOException {
boolean admitted=slots.tryAcquire();
String requestId=UUID.randomUUID().toString();
exchange.getResponseHeaders().set("x-amz-request-id",requestId);
exchange.getResponseHeaders().set("X-Content-Type-Options","nosniff");
boolean admitted = slots.tryAcquire();
String requestId = UUID.randomUUID().toString();
exchange.getResponseHeaders().set("x-amz-request-id", requestId);
exchange.getResponseHeaders().set("X-Content-Type-Options", "nosniff");
try {
if(!admitted)throw new StoreException(503,"SlowDown","Too many concurrent requests");
if(exchange.getRequestURI().getRawPath().equals("/health")&&exchange.getRequestMethod().equals("GET")){
byte[] body="{\"status\":\"ok\",\"service\":\"lunarsky-objectstore\"}".getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type","application/json");exchange.sendResponseHeaders(200,body.length);exchange.getResponseBody().write(body);return;
if (!admitted) throw new StoreException(503, "SlowDown", "Too many concurrent requests");
if (exchange.getRequestURI().getRawPath().equals("/health") && exchange.getRequestMethod().equals("GET")) {
byte[] body = "{\"status\":\"ok\",\"service\":\"lunarsky-objectstore\"}".getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type", "application/json");
exchange.sendResponseHeaders(200, body.length);
exchange.getResponseBody().write(body);
return;
}
String hash=authentication.verify(exchange.getRequestMethod(),exchange.getRequestURI(),exchange.getRequestHeaders());
String path=SigV4.decode(exchange.getRequestURI().getRawPath());
String prefix="/"+bucket+"/";
if(!path.startsWith(prefix))throw new StoreException(404,"NoSuchBucket","Bucket not found");
String key=path.substring(prefix.length());
if(key.isEmpty()||key.getBytes(StandardCharsets.UTF_8).length>1024||key.indexOf('\0')>=0)throw new StoreException(400,"InvalidArgument","Invalid object key");
String query=exchange.getRequestURI().getRawQuery();
if(query!=null&&!query.isEmpty()&&!query.matches("x-id=(PutObject|GetObject|HeadObject|DeleteObject)"))unsupported("Query operation");
var headers=exchange.getRequestHeaders();
if(headers.containsKey("range")||headers.containsKey("if-match")||headers.containsKey("if-modified-since")||headers.containsKey("if-unmodified-since"))unsupported("Range or conditional read");
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))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-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");
if (exchange.getRequestURI().getRawPath().equals("/ready") && exchange.getRequestMethod().equals("GET")) {
boolean ready = store.ready();
byte[] body = (ready ? "ready" : "unavailable").getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type", "text/plain; charset=utf-8");
exchange.sendResponseHeaders(ready ? 200 : 503, body.length);
exchange.getResponseBody().write(body);
return;
}
String method=exchange.getRequestMethod();
String algorithm=SigV4.single(headers,"x-amz-sdk-checksum-algorithm");
if(algorithm!=null&&!algorithm.equals("SHA256"))unsupported("Checksum algorithm");
if(!method.equals("PUT")&&(headers.containsKey("transfer-encoding")||(headers.containsKey("content-length")&&!"0".equals(SigV4.single(headers,"content-length")))))throw new StoreException(400,"InvalidRequest","Read/delete requests must have empty bodies");
if(!method.equals("PUT")&&!hash.equals(SigV4.hex(SigV4.hash(new byte[0]))))throw new StoreException(400,"InvalidRequest","Read/delete requests must have empty bodies");
switch(method){
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 + "/")) {
if (!exchange.getRequestMethod().equals("GET") || !"2".equals(query.get("list-type")) ||
!query.keySet().stream().allMatch(java.util.Set.of("list-type", "prefix", "delimiter", "max-keys",
"continuation-token", "start-after", "encoding-type", "x-id")::contains) ||
(query.containsKey("x-id") && !"ListObjectsV2".equals(query.get("x-id"))))
unsupported("Bucket operation");
requireEmptyBody(exchange, hash);
listObjects(exchange, query);
return;
}
String prefix = "/" + bucket + "/";
if (!path.startsWith(prefix)) throw new StoreException(404, "NoSuchBucket", "Bucket not found");
String key = path.substring(prefix.length());
if (key.isEmpty() || key.getBytes(StandardCharsets.UTF_8).length > 1024 || key.indexOf('\0') >= 0)
throw new StoreException(400, "InvalidArgument", "Invalid object key");
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")) ||
"HeadObject".equals(query.get("x-id")) || "DeleteObject".equals(query.get("x-id")))))
unsupported("Query operation");
var headers = exchange.getRequestHeaders();
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))
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-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");
}
String algorithm = SigV4.single(headers, "x-amz-sdk-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) {
case "PUT" -> {
String length=SigV4.single(headers,"content-length"),condition=SigV4.single(headers,"if-none-match");
if(condition!=null&&!condition.equals("*"))unsupported("Write condition");
long bytes;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");
DiskStore.Metadata data=store.put(bucket,key,exchange.getRequestBody(),bytes,hash,SigV4.single(headers,"x-amz-checksum-sha256"),condition!=null);
exchange.getResponseHeaders().set("ETag","\""+data.etag()+"\"");
exchange.getResponseHeaders().set("x-amz-checksum-sha256",java.util.Base64.getEncoder().encodeToString(data.sha256()));
exchange.sendResponseHeaders(200,-1);
String length = SigV4.single(headers, "content-length"), condition = SigV4.single(headers, "if-none-match");
if (condition != null && !condition.equals("*")) unsupported("Write condition");
long bytes;
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");
String contentType = contentType(headers);
ObjectStorage.Metadata data = store.put(bucket, key, exchange.getRequestBody(), bytes, hash,
SigV4.single(headers, "x-amz-checksum-sha256"), condition != null, contentType);
exchange.getResponseHeaders().set("ETag", "\"" + data.etag() + "\"");
exchange.getResponseHeaders().set("x-amz-checksum-sha256", Base64.getEncoder().encodeToString(data.sha256()));
exchange.sendResponseHeaders(200, -1);
}
case "GET", "HEAD" -> {
if(headers.containsKey("if-none-match"))unsupported("Conditional read");
try(var object=store.open(bucket,key)){
var meta=object.metadata();
exchange.getResponseHeaders().set("Content-Type","application/octet-stream");
exchange.getResponseHeaders().set("Content-Length",Long.toString(meta.length()));
exchange.getResponseHeaders().set("ETag","\""+meta.etag()+"\"");
exchange.getResponseHeaders().set("Last-Modified",DateTimeFormatter.RFC_1123_DATE_TIME.withZone(ZoneOffset.UTC).format(Instant.ofEpochMilli(meta.modified())));
if(method.equals("HEAD")||meta.length()==0)exchange.sendResponseHeaders(200,-1);
else{exchange.sendResponseHeaders(200,meta.length());object.stream().transferTo(exchange.getResponseBody());}
}
case "GET", "HEAD" -> readObject(exchange, key);
case "DELETE" -> {
if (headers.containsKey("if-none-match")) unsupported("Conditional delete");
store.delete(bucket, key);
exchange.sendResponseHeaders(204, -1);
}
case "DELETE" -> {if(headers.containsKey("if-none-match"))unsupported("Conditional delete");store.delete(bucket,key);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();}
} 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 void unsupported(String feature){throw new StoreException(501,"NotImplemented",feature+" is not supported in this prototype");}
private static String xml(String text){return text.replace("&","&amp;").replace("<","&lt;").replace(">","&gt;").replace("\"","&quot;");}
private static void sendError(HttpExchange exchange,int status,String code,String message,String id)throws IOException{
if(exchange.getResponseCode()!=-1)return;
byte[] body=("<?xml version=\"1.0\" encoding=\"UTF-8\"?><Error><Code>"+xml(code)+"</Code><Message>"+xml(message)+"</Message><RequestId>"+id+"</RequestId></Error>").getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type","application/xml");exchange.sendResponseHeaders(status,exchange.getRequestMethod().equals("HEAD")?-1:body.length);
if(!exchange.getRequestMethod().equals("HEAD"))exchange.getResponseBody().write(body);
private static String contentType(com.sun.net.httpserver.Headers headers) {
String value = SigV4.single(headers, "content-type");
if (value == null) return "application/octet-stream";
if (value.isBlank() || value.getBytes(StandardCharsets.UTF_8).length > 255 ||
!value.chars().allMatch(c -> c >= 32 && c <= 126))
throw new StoreException(400, "InvalidArgument", "Invalid Content-Type");
return value;
}
public static void main(String[] args)throws Exception{
Map<String,String> env=System.getenv();
String access=required(env,"S3_ACCESS_KEY"),secret=required(env,"S3_SECRET_KEY"),bucket=env.getOrDefault("S3_BUCKET","lunaris-files"),region=env.getOrDefault("S3_REGION","us-east-1");
if(!access.matches("[A-Za-z0-9]{16,128}")||secret.length()<32||!bucket.matches("[a-z0-9][a-z0-9-]{1,61}[a-z0-9]"))throw new IllegalArgumentException("Invalid storage credentials/bucket configuration");
long maxObject=Long.parseLong(env.getOrDefault("MAX_OBJECT_BYTES","10485760")),maxTotal=Long.parseLong(env.getOrDefault("MAX_TOTAL_BYTES","2147483648"));
if(maxObject<1||maxObject>1073741824L||maxTotal<maxObject)throw new IllegalArgumentException("Invalid size limits");
var store=new DiskStore(Path.of(env.getOrDefault("DATA_DIR","/data")),maxObject,maxTotal);
var app=new Main(store,new SigV4(access,secret,region,Clock.systemUTC()),bucket);
var server=HttpServer.create(new InetSocketAddress(Integer.parseInt(env.getOrDefault("PORT","9000"))),64);
var executor=Executors.newVirtualThreadPerTaskExecutor();server.setExecutor(executor);server.createContext("/",app::handle);
Runtime.getRuntime().addShutdownHook(new Thread(()->{
private static void requireEmptyBody(HttpExchange exchange, String hash) {
var headers = exchange.getRequestHeaders();
if (headers.containsKey("transfer-encoding") ||
(headers.containsKey("content-length") && !"0".equals(SigV4.single(headers, "content-length"))) ||
!hash.equals(SigV4.hex(SigV4.hash(new byte[0]))))
throw new StoreException(400, "InvalidRequest", "Request must have an empty body");
}
private static boolean multipartRequest(String method, Map<String, String> query) {
if (query.containsKey("uploads"))
return method.equals("POST") && query.get("uploads").isEmpty() &&
query.keySet().stream().allMatch(java.util.Set.of("uploads", "x-id")::contains) &&
(!query.containsKey("x-id") || query.get("x-id").equals("CreateMultipartUpload"));
if (!query.containsKey("uploadId") ||
!query.keySet().stream().allMatch(java.util.Set.of("uploadId", "partNumber", "x-id")::contains))
return false;
String xId = query.get("x-id");
if (method.equals("PUT")) return query.containsKey("partNumber") &&
(xId == null || xId.equals("UploadPart"));
if (query.containsKey("partNumber")) return false;
return (method.equals("POST") && (xId == null || xId.equals("CompleteMultipartUpload"))) ||
(method.equals("DELETE") && (xId == null || xId.equals("AbortMultipartUpload")));
}
private void handleMultipart(HttpExchange exchange, String method, Map<String, String> query,
String key, String hash) throws IOException {
var headers = exchange.getRequestHeaders();
if (headers.containsKey("content-encoding") || headers.containsKey("if-none-match"))
unsupported("Multipart request header");
if (query.containsKey("uploads")) {
requireEmptyBody(exchange, hash);
String id = multipart.create(bucket, key, contentType(headers));
sendXml(exchange, 200, "<InitiateMultipartUploadResult><Bucket>" + xml(bucket) +
"</Bucket><Key>" + xml(key) + "</Key><UploadId>" + id +
"</UploadId></InitiateMultipartUploadResult>");
return;
}
String id = query.get("uploadId");
switch (method) {
case "PUT" -> {
int number;
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"));
exchange.getResponseHeaders().set("ETag", "\"" + etag + "\"");
exchange.sendResponseHeaders(200, -1);
}
case "POST" -> {
byte[] body = signedBody(exchange, hash, 65536);
List<MultipartStorage.Part> parts = completedParts(body);
var meta = multipart.complete(id, bucket, key, parts);
sendXml(exchange, 200, "<CompleteMultipartUploadResult><Bucket>" + xml(bucket) +
"</Bucket><Key>" + xml(key) + "</Key><ETag>&quot;" + meta.etag() +
"&quot;</ETag></CompleteMultipartUploadResult>");
}
case "DELETE" -> {
requireEmptyBody(exchange, hash);
multipart.abort(id, bucket, key);
exchange.sendResponseHeaders(204, -1);
}
default -> unsupported("Multipart operation");
}
}
private static long contentLength(com.sun.net.httpserver.Headers headers) {
String text = SigV4.single(headers, "content-length");
if (text == null) return -1;
try { return Long.parseLong(text); }
catch (NumberFormatException e) { throw new StoreException(400, "InvalidArgument", "Invalid Content-Length"); }
}
private static byte[] signedBody(HttpExchange exchange, String hash, int limit) throws IOException {
long length = contentLength(exchange.getRequestHeaders());
if (length < 0) throw new StoreException(411, "MissingContentLength", "Content-Length is required");
if (length > limit) throw new StoreException(413, "EntityTooLarge", "Request body is too large");
byte[] body = exchange.getRequestBody().readNBytes(limit + 1);
if (body.length != length) throw new StoreException(400, "IncompleteBody", "Body length does not match Content-Length");
if (!SigV4.hex(SigV4.hash(body)).equals(hash))
throw new StoreException(400, "XAmzContentSHA256Mismatch", "Payload hash mismatch");
return body;
}
private static List<MultipartStorage.Part> completedParts(byte[] body) {
try {
var factory = javax.xml.parsers.DocumentBuilderFactory.newInstance();
factory.setFeature("http://apache.org/xml/features/disallow-doctype-decl", true);
factory.setFeature("http://xml.org/sax/features/external-general-entities", false);
factory.setFeature("http://xml.org/sax/features/external-parameter-entities", false);
factory.setFeature(javax.xml.XMLConstants.FEATURE_SECURE_PROCESSING, true);
factory.setExpandEntityReferences(false);
var document = factory.newDocumentBuilder().parse(new ByteArrayInputStream(body));
if (!document.getDocumentElement().getNodeName().equals("CompleteMultipartUpload"))
throw new IllegalArgumentException();
var nodes = document.getDocumentElement().getChildNodes();
List<MultipartStorage.Part> parts = new ArrayList<>();
for (int i = 0; i < nodes.getLength(); i++) {
if (!(nodes.item(i) instanceof org.w3c.dom.Element element)) continue;
if (!element.getTagName().equals("Part")) throw new IllegalArgumentException();
String number = null, etag = null;
var fields = element.getChildNodes();
for (int j = 0; j < fields.getLength(); j++) {
if (!(fields.item(j) instanceof org.w3c.dom.Element field)) continue;
if (field.getTagName().equals("PartNumber")) number = field.getTextContent().trim();
else if (field.getTagName().equals("ETag")) etag = field.getTextContent().trim();
else throw new IllegalArgumentException();
}
if (number == null || etag == null || !etag.matches("\"?[0-9a-f]{32}\"?"))
throw new IllegalArgumentException();
parts.add(new MultipartStorage.Part(Integer.parseInt(number), etag));
if (parts.size() > 10000) throw new IllegalArgumentException();
}
return parts;
} catch (Exception error) {
throw new StoreException(400, "MalformedXML", "Invalid multipart completion body");
}
}
private static void sendXml(HttpExchange exchange, int status, String xml) throws IOException {
byte[] body = ("<?xml version=\"1.0\" encoding=\"UTF-8\"?>" + xml).getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type", "application/xml");
exchange.sendResponseHeaders(status, body.length);
exchange.getResponseBody().write(body);
}
private void readObject(HttpExchange exchange, String key) throws IOException {
var headers = exchange.getRequestHeaders();
if (headers.containsKey("if-match") || headers.containsKey("if-modified-since") ||
headers.containsKey("if-unmodified-since")) unsupported("Conditional read");
try (var object = store.open(bucket, key)) {
var meta = object.metadata();
String etag = "\"" + meta.etag() + "\"";
String noneMatch = SigV4.single(headers, "if-none-match");
if (noneMatch != null && (noneMatch.equals("*") ||
java.util.Arrays.stream(noneMatch.split(",")).map(String::trim).anyMatch(etag::equals))) {
exchange.getResponseHeaders().set("ETag", etag);
exchange.sendResponseHeaders(304, -1);
return;
}
Range range;
try { range = range(SigV4.single(headers, "range"), meta.length()); }
catch (StoreException error) {
if (error.status == 416) exchange.getResponseHeaders().set("Content-Range", "bytes */" + meta.length());
throw error;
}
var response = exchange.getResponseHeaders();
response.set("Content-Type", meta.contentType());
response.set("Content-Length", Long.toString(range.length()));
response.set("Accept-Ranges", "bytes");
response.set("ETag", etag);
response.set("Last-Modified", DateTimeFormatter.RFC_1123_DATE_TIME.withZone(ZoneOffset.UTC)
.format(Instant.ofEpochMilli(meta.modified())));
if (range.partial()) response.set("Content-Range", "bytes " + range.start() + "-" + range.end() + "/" + meta.length());
int status = range.partial() ? 206 : 200;
if (exchange.getRequestMethod().equals("HEAD") || range.length() == 0)
exchange.sendResponseHeaders(status, -1);
else {
object.stream().skipNBytes(range.start());
exchange.sendResponseHeaders(status, range.length());
byte[] buffer = new byte[65536]; long left = range.length();
while (left > 0) {
int n = object.stream().read(buffer, 0, (int) Math.min(buffer.length, left));
if (n < 0) throw new IOException("Object body ended before its recorded length");
exchange.getResponseBody().write(buffer, 0, n);
left -= n;
}
}
}
}
private record Range(long start, long end, boolean partial) {
long length() { return end < start ? 0 : end - start + 1; }
}
private static Range range(String header, long size) {
if (header == null) return new Range(0, size - 1, false);
if (!header.matches("bytes=[0-9]*-[0-9]*") || header.equals("bytes=-") || size == 0)
throw new StoreException(416, "InvalidRange", "The requested range is not satisfiable");
String[] parts = header.substring(6).split("-", -1);
try {
long start, end;
if (parts[0].isEmpty()) {
long suffix = Long.parseLong(parts[1]);
if (suffix == 0) throw new NumberFormatException();
start = Math.max(0, size - suffix); end = size - 1;
} else {
start = Long.parseLong(parts[0]);
end = parts[1].isEmpty() ? size - 1 : Math.min(Long.parseLong(parts[1]), size - 1);
}
if (start >= size || end < start) throw new NumberFormatException();
return new Range(start, end, true);
} catch (NumberFormatException e) {
throw new StoreException(416, "InvalidRange", "The requested range is not satisfiable");
}
}
private void listObjects(HttpExchange exchange, Map<String, String> query) throws IOException {
String prefix = query.getOrDefault("prefix", ""), delimiter = query.getOrDefault("delimiter", "");
String encoding = query.get("encoding-type");
if (encoding != null && !encoding.equals("url")) unsupported("Encoding type");
if (query.containsKey("continuation-token") && query.containsKey("start-after"))
throw new StoreException(400, "InvalidArgument", "Use either continuation-token or start-after");
int maxKeys;
try { maxKeys = Integer.parseInt(query.getOrDefault("max-keys", "1000")); }
catch (NumberFormatException e) { throw new StoreException(400, "InvalidArgument", "Invalid max-keys"); }
if (maxKeys < 0 || maxKeys > 1000) throw new StoreException(400, "InvalidArgument", "Invalid max-keys");
String after = query.get("start-after");
if (query.containsKey("continuation-token")) {
try {
byte[] decoded = Base64.getUrlDecoder().decode(query.get("continuation-token"));
after = StandardCharsets.UTF_8.newDecoder().onMalformedInput(java.nio.charset.CodingErrorAction.REPORT)
.decode(java.nio.ByteBuffer.wrap(decoded)).toString();
} catch (IllegalArgumentException | java.nio.charset.CharacterCodingException e) {
throw new StoreException(400, "InvalidArgument", "Invalid continuation token");
}
}
var page = store.list(bucket, prefix, delimiter, maxKeys, after);
StringBuilder xml = new StringBuilder("<?xml version=\"1.0\" encoding=\"UTF-8\"?><ListBucketResult xmlns=\"http://s3.amazonaws.com/doc/2006-03-01/\">");
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 (encoding != null) xml.append("<EncodingType>url</EncodingType>");
if (query.containsKey("continuation-token")) xml.append("<ContinuationToken>")
.append(xml(query.get("continuation-token"))).append("</ContinuationToken>");
if (query.containsKey("start-after")) xml.append("<StartAfter>")
.append(xml(listKey(query.get("start-after"), encoding))).append("</StartAfter>");
xml.append("<KeyCount>").append(page.keyCount()).append("</KeyCount><MaxKeys>").append(maxKeys)
.append("</MaxKeys><IsTruncated>").append(page.truncated()).append("</IsTruncated>");
int objectAt = 0, prefixAt = 0;
while (objectAt < page.objects().size() || prefixAt < page.prefixes().size()) {
if (objectAt < page.objects().size() &&
(prefixAt == page.prefixes().size() ||
page.objects().get(objectAt).key().compareTo(page.prefixes().get(prefixAt)) < 0)) {
var entry = page.objects().get(objectAt++);
var meta = entry.metadata();
xml.append("<Contents><Key>").append(xml(listKey(entry.key(), encoding))).append("</Key><LastModified>")
.append(Instant.ofEpochMilli(meta.modified())).append("</LastModified><ETag>&quot;")
.append(meta.etag()).append("&quot;</ETag><Size>").append(meta.length())
.append("</Size><StorageClass>STANDARD</StorageClass></Contents>");
} else {
xml.append("<CommonPrefixes><Prefix>")
.append(xml(listKey(page.prefixes().get(prefixAt++), encoding)))
.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) {
return encoding == null ? key : SigV4.encode(key, false);
}
private static Map<String, String> query(String raw) {
Map<String, String> result = new HashMap<>();
if (raw == null || raw.isEmpty()) return result;
for (String part : raw.split("&", -1)) {
String[] pair = part.split("=", 2);
String name = SigV4.decode(pair[0]);
String value = SigV4.decode(pair.length == 2 ? pair[1] : "");
if (result.put(name, value) != null)
throw new StoreException(400, "InvalidArgument", "Duplicate query parameter");
}
return result;
}
private static void unsupported(String feature) { throw new StoreException(501, "NotImplemented", feature + " is not supported"); }
private static String xml(String text) {
return text.replace("&", "&amp;").replace("<", "&lt;").replace(">", "&gt;").replace("\"", "&quot;");
}
private static void sendError(HttpExchange exchange, int status, String code, String message, String id) throws IOException {
if (exchange.getResponseCode() != -1) return;
byte[] body = ("<?xml version=\"1.0\" encoding=\"UTF-8\"?><Error><Code>" + xml(code) +
"</Code><Message>" + xml(message) + "</Message><RequestId>" + id + "</RequestId></Error>")
.getBytes(StandardCharsets.UTF_8);
exchange.getResponseHeaders().set("Content-Type", "application/xml");
exchange.sendResponseHeaders(status, exchange.getRequestMethod().equals("HEAD") ? -1 : body.length);
if (!exchange.getRequestMethod().equals("HEAD")) exchange.getResponseBody().write(body);
}
public static void main(String[] args) throws Exception {
Map<String, String> env = System.getenv();
String access = required(env, "S3_ACCESS_KEY"), secret = required(env, "S3_SECRET_KEY");
String bucket = env.getOrDefault("S3_BUCKET", "lunaris-files"), region = env.getOrDefault("S3_REGION", "us-east-1");
if (!access.matches("[A-Za-z0-9]{16,128}") || secret.length() < 32 ||
!bucket.matches("[a-z0-9][a-z0-9-]{1,61}[a-z0-9]"))
throw new IllegalArgumentException("Invalid storage credentials/bucket configuration");
long maxObject = Long.parseLong(env.getOrDefault("MAX_OBJECT_BYTES", "134217728"));
long maxTotal = Long.parseLong(env.getOrDefault("MAX_TOTAL_BYTES", "2147483648"));
if (maxObject < 1 || maxObject > 1073741824L || maxTotal < maxObject)
throw new IllegalArgumentException("Invalid size limits");
String mode = env.getOrDefault("STORE_MODE", "disk");
ObjectStorage store;
MultipartStorage multipart;
if (mode.equals("cluster")) {
if (!"true".equals(env.get("CLUSTER_LOCAL_DEV")))
throw new IllegalArgumentException("Cluster mode is local development only; set CLUSTER_LOCAL_DEV=true");
String[] urls = required(env, "CLUSTER_NODES").split(",", -1);
if (urls.length < 2) throw new IllegalArgumentException("CLUSTER_NODES requires at least two URLs");
store = new ClusterStore(required(env, "POSTGRES_JDBC_URL"), required(env, "POSTGRES_USER"),
required(env, "POSTGRES_PASSWORD"), bucket,
java.util.Arrays.stream(urls).map(URI::create).toList(), required(env, "CLUSTER_TOKEN"), null,
maxObject, maxTotal, "true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS")));
multipart = new UnavailableMultipart();
} else if (mode.equals("disk")) {
DiskStore disk = new DiskStore(Path.of(env.getOrDefault("DATA_DIR", "/data")), maxObject, maxTotal);
store = disk;
multipart = new MultipartStore(disk);
} else throw new IllegalArgumentException("Invalid STORE_MODE");
var app = new Main(store, multipart, new SigV4(access, secret, region, Clock.systemUTC()), bucket);
int port = Integer.parseInt(env.getOrDefault("PORT", "9000"));
var server = HttpServer.create(mode.equals("cluster")
? new InetSocketAddress(env.getOrDefault("BIND_ADDRESS", "127.0.0.1"), port)
: new InetSocketAddress(port), 64);
var executor = Executors.newVirtualThreadPerTaskExecutor();
server.setExecutor(executor); server.createContext("/", app::handle);
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
server.stop(5);
executor.close();
try { store.close(); }
catch (IOException error) { System.err.println("Could not release ObjectStore data lock: "+error.getMessage()); }
catch (IOException error) { System.err.println("Could not release ObjectStore data lock: " + error.getMessage()); }
}));
server.start();System.out.println("LunarSky ObjectStore listening; S3 object-operation prototype, bucket="+bucket);
server.start();
if (store instanceof DiskStore disk) printStartup(disk, app.multipart, bucket, region, port);
else System.out.println("ObjectStore local cluster prototype v" + Version.VALUE + " listening on :" + port);
}
private static void printStartup(DiskStore store, MultipartStorage multipart,
String bucket, String region, int port) {
System.out.println(" *");
System.out.println(" / \\ LUNARSKY");
System.out.println(" / L \\ ObjectStore");
System.out.println(" /_____\\ v" + Version.VALUE);
System.out.println();
System.out.println(" Bucket " + bucket + " (" + region + ")");
System.out.println(" Objects " + store.objectCount() + " total, " +
store.indexedObjects() + " indexed, " + size(store.usedBytes()) + " / " + size(store.maxTotal()));
if (store.legacyObjects() > 0)
System.out.println(" Legacy " + store.legacyObjects() + " objects without stored keys");
System.out.println(" Multipart " + multipart.activeUploads() + " active, " +
size(multipart.stagedBytes()) + " staged");
System.out.println(" Listening :" + port);
}
private static String size(long bytes) {
if (bytes < 1024) return bytes + " B";
String[] units = {"KiB", "MiB", "GiB", "TiB"};
double value = bytes;
int unit = -1;
do { value /= 1024; unit++; } while (value >= 1024 && unit < units.length - 1);
return String.format(Locale.ROOT, "%.1f %s", value, units[unit]);
}
private static String required(Map<String, String> env, String key) {
String value = env.get(key);
if (value == null || value.isBlank()) throw new IllegalArgumentException("Missing " + key);
return value;
}
private static String required(Map<String,String> env,String key){String value=env.get(key);if(value==null||value.isBlank())throw new IllegalArgumentException("Missing "+key);return value;}
}
@@ -0,0 +1,17 @@
package cloud.lunarsky.store;
import java.io.IOException;
import java.io.InputStream;
import java.util.List;
interface MultipartStorage {
record Part(int number, String etag) {}
String create(String bucket, String key, String contentType) throws IOException;
String putPart(String id, String bucket, String key, int number, InputStream input,
long length, String expectedHash, String checksum) throws IOException;
ObjectStorage.Metadata complete(String id, String bucket, String key, List<Part> parts) throws IOException;
void abort(String id, String bucket, String key) throws IOException;
int activeUploads();
long stagedBytes();
}
@@ -0,0 +1,217 @@
package cloud.lunarsky.store;
import java.io.*;
import java.nio.file.*;
import java.security.MessageDigest;
import java.util.*;
import cloud.lunarsky.store.MultipartStorage.Part;
final class MultipartStore implements MultipartStorage {
private static final int MAGIC = 0x4c534d50;
private final DiskStore store;
private final Path root;
private long staged;
private int active;
private record Upload(String bucket, String key, String contentType) {}
MultipartStore(DiskStore store) throws IOException {
this.store = store;
root = store.root().resolve("multipart");
Files.createDirectories(root);
DiskStore.syncDirectory(store.root());
try (var uploads = Files.list(root)) {
for (Path dir : uploads.toList()) {
if (dir.getFileName().toString().startsWith(".creating-")) {
discardCreating(dir);
continue;
}
if (!Files.isDirectory(dir)) throw new IOException("Invalid multipart upload entry: " + dir);
readUpload(dir);
active++;
try (var files = Files.list(dir)) {
for (Path file : files.toList()) {
if (file.getFileName().toString().matches("part-[0-9]{5}"))
staged = Math.addExact(staged, Files.size(file));
else if (!file.getFileName().toString().equals("manifest"))
throw new IOException("Invalid multipart upload entry: " + file);
}
}
}
}
if (staged > store.maxTotal()) throw new IOException("Multipart staging limit exceeded");
}
public synchronized String create(String bucket, String key, String contentType) throws IOException {
if (active >= 32) throw new StoreException(503, "SlowDown", "Too many active uploads");
String id = UUID.randomUUID().toString();
Path pending = root.resolve(".creating-" + id), dir = root.resolve(id);
Files.createDirectory(pending);
try {
try (var output = new DataOutputStream(Files.newOutputStream(pending.resolve("manifest"), StandardOpenOption.CREATE_NEW))) {
output.writeInt(MAGIC);
output.writeUTF(bucket);
output.writeUTF(key);
output.writeUTF(contentType);
}
try (var channel = java.nio.channels.FileChannel.open(pending.resolve("manifest"), StandardOpenOption.READ)) {
channel.force(true);
}
DiskStore.syncDirectory(pending);
Files.move(pending, dir, StandardCopyOption.ATOMIC_MOVE);
DiskStore.syncDirectory(root);
} catch (IOException error) {
discardCreating(pending);
discardCreating(dir);
throw error;
}
active++;
return id;
}
private void discardCreating(Path dir) throws IOException {
if (!Files.exists(dir)) return;
if (Files.isDirectory(dir)) {
try (var files = Files.list(dir)) {
for (Path file : files.toList()) Files.delete(file);
}
}
Files.delete(dir);
DiskStore.syncDirectory(root);
}
public synchronized int activeUploads() { return active; }
public synchronized long stagedBytes() { return staged; }
public synchronized String putPart(String id, String bucket, String key, int number, InputStream input,
long length, String expectedHash, String checksum) throws IOException {
if (number < 1 || number > 10000) throw new StoreException(400, "InvalidArgument", "Invalid part number");
if (length < 0) throw new StoreException(411, "MissingContentLength", "Content-Length is required");
if (length > store.maxObject()) throw new StoreException(413, "EntityTooLarge", "Part exceeds the object limit");
Path dir = upload(id, bucket, key), target = part(dir, number);
long previous = Files.exists(target) ? Files.size(target) : 0;
if (staged - previous + length > store.maxTotal())
throw new StoreException(507, "InsufficientStorage", "Multipart staging limit reached");
Path pending = Files.createTempFile(store.root().resolve("pending"), "part-", ".part");
try {
MessageDigest sha = digest("SHA-256"), md5 = digest("MD5");
long count = 0;
try (var output = Files.newOutputStream(pending)) {
byte[] buffer = new byte[65536]; int n;
while ((n = input.read(buffer)) != -1) {
count += n;
if (count > length) throw new StoreException(413, "EntityTooLarge", "Part exceeds declared size");
sha.update(buffer, 0, n); md5.update(buffer, 0, n); output.write(buffer, 0, n);
}
}
if (count != length) throw new StoreException(400, "IncompleteBody", "Part length does not match Content-Length");
byte[] actual = sha.digest();
if (!MessageDigest.isEqual(actual, HexFormat.of().parseHex(expectedHash)))
throw new StoreException(400, "XAmzContentSHA256Mismatch", "Part hash mismatch");
if (checksum != null && !Base64.getEncoder().encodeToString(actual).equals(checksum))
throw new StoreException(400, "BadDigest", "SHA-256 checksum mismatch");
try (var channel = java.nio.channels.FileChannel.open(pending, StandardOpenOption.WRITE)) { channel.force(true); }
Files.move(pending, target, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING);
staged = staged - previous + length;
DiskStore.syncDirectory(dir);
return SigV4.hex(md5.digest());
} finally { Files.deleteIfExists(pending); }
}
public synchronized ObjectStorage.Metadata complete(String id, String bucket, String key, List<Part> parts) throws IOException {
Path dir = upload(id, bucket, key);
if (parts.isEmpty() || parts.size() > 10000)
throw new StoreException(400, "InvalidPart", "No valid parts supplied");
MessageDigest sha = digest("SHA-256");
List<Path> paths = new ArrayList<>();
long total = 0; int last = 0;
for (Part part : parts) {
if (part.number() <= last || part.number() > 10000)
throw new StoreException(400, "InvalidPartOrder", "Parts must be in ascending order");
last = part.number();
Path file = part(dir, part.number());
if (!Files.isRegularFile(file)) throw new StoreException(400, "InvalidPart", "Missing part");
long length = Files.size(file);
total += length;
if (total > store.maxObject()) throw new StoreException(413, "EntityTooLarge", "Object exceeds the configured size limit");
MessageDigest md5 = digest("MD5");
try (var input = Files.newInputStream(file)) {
byte[] buffer = new byte[65536]; int n;
while ((n = input.read(buffer)) != -1) { sha.update(buffer, 0, n); md5.update(buffer, 0, n); }
}
if (!SigV4.hex(md5.digest()).equals(part.etag().replace("\"", "")))
throw new StoreException(400, "InvalidPart", "Part ETag mismatch");
paths.add(file);
}
ObjectStorage.Metadata result;
try (InputStream input = new PartsInput(paths)) {
result = store.put(bucket, key, input, total, SigV4.hex(sha.digest()), null, false,
readUpload(dir).contentType());
}
remove(dir);
return result;
}
public synchronized void abort(String id, String bucket, String key) throws IOException {
remove(upload(id, bucket, key));
}
private Path upload(String id, String bucket, String key) throws IOException {
if (!id.matches("[0-9a-f-]{36}")) throw new StoreException(404, "NoSuchUpload", "Upload not found");
Path dir = root.resolve(id);
if (!Files.isDirectory(dir)) throw new StoreException(404, "NoSuchUpload", "Upload not found");
Upload upload = readUpload(dir);
if (!upload.bucket().equals(bucket) || !upload.key().equals(key))
throw new StoreException(404, "NoSuchUpload", "Upload not found");
return dir;
}
private static Upload readUpload(Path dir) throws IOException {
try (var input = new DataInputStream(Files.newInputStream(dir.resolve("manifest")))) {
if (input.readInt() != MAGIC) throw new IOException("Invalid multipart upload manifest");
Upload upload = new Upload(input.readUTF(), input.readUTF(), input.readUTF());
if (input.read() != -1) throw new IOException("Invalid multipart upload manifest");
return upload;
}
}
private static Path part(Path dir, int number) { return dir.resolve("part-%05d".formatted(number)); }
private void remove(Path dir) throws IOException {
long removed = 0;
try (var files = Files.list(dir)) {
for (Path file : files.toList()) {
if (file.getFileName().toString().matches("part-[0-9]{5}")) removed += Files.size(file);
Files.delete(file);
}
}
DiskStore.syncDirectory(dir);
Files.delete(dir);
DiskStore.syncDirectory(root);
staged -= removed;
active--;
}
private static MessageDigest digest(String algorithm) {
try { return MessageDigest.getInstance(algorithm); }
catch (java.security.NoSuchAlgorithmException e) { throw new IllegalStateException(e); }
}
private static final class PartsInput extends InputStream {
private final Iterator<Path> parts;
private InputStream current;
PartsInput(List<Path> files) { parts = files.iterator(); }
@Override public int read() throws IOException {
byte[] one = new byte[1];
return read(one, 0, 1) < 0 ? -1 : one[0] & 255;
}
@Override public int read(byte[] buffer, int offset, int length) throws IOException {
if (length == 0) return 0;
while (true) {
if (current == null) {
if (!parts.hasNext()) return -1;
current = Files.newInputStream(parts.next());
}
int n = current.read(buffer, offset, length);
if (n >= 0) return n;
current.close(); current = null;
}
}
@Override public void close() throws IOException { if (current != null) current.close(); }
}
}
+147
View File
@@ -0,0 +1,147 @@
package cloud.lunarsky.store;
import java.io.IOException;
import java.io.InputStream;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.security.MessageDigest;
import java.time.Duration;
import java.util.HashSet;
import java.util.HexFormat;
import java.util.List;
import java.util.Set;
import java.util.UUID;
final class NodeClient {
record Node(UUID id, UUID hostId, URI url) {}
private static final HttpClient IDENTITY_HTTP = HttpClient.newBuilder()
.connectTimeout(Duration.ofSeconds(2)).build();
private final List<Node> nodes;
private final String token;
private final String repairToken;
private final HttpClient http = HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(3)).build();
NodeClient(List<Node> nodes, String token, String repairToken) {
if (nodes.isEmpty() || nodes.stream().map(Node::id).distinct().count() != nodes.size() ||
nodes.stream().map(Node::url).distinct().count() != nodes.size())
throw new IllegalArgumentException("Cluster node IDs and URLs must be unique");
if (token == null || token.length() < 32) throw new IllegalArgumentException("Invalid cluster token");
for (Node node : nodes) {
if (node.id() == null || node.hostId() == null)
throw new IllegalArgumentException("Invalid storage node identity");
validateUrl(node.url());
}
this.nodes = List.copyOf(nodes);
this.token = token;
this.repairToken = repairToken;
}
int count() { return nodes.size(); }
List<Node> nodes() { return nodes; }
Node node(int index) { return nodes.get(index); }
int index(UUID id) {
for (int i = 0; i < nodes.size(); i++) if (nodes.get(i).id().equals(id)) return i;
return -1;
}
UUID faultDomain(int index, boolean testNodeDomains) {
Node node = nodes.get(index);
return testNodeDomains ? node.id() : node.hostId();
}
static NodeIdentity probe(URI url, String token) throws IOException {
validateUrl(url);
if (token == null || token.length() < 32) throw new IllegalArgumentException("Invalid cluster token");
HttpRequest request = HttpRequest.newBuilder(url.resolve("/identity"))
.timeout(Duration.ofSeconds(2)).header("X-Cluster-Token", token).GET().build();
try {
HttpResponse<InputStream> response = IDENTITY_HTTP.send(request, HttpResponse.BodyHandlers.ofInputStream());
try (InputStream body = response.body()) {
if (response.statusCode() != 200) throw new IOException("Node identity request failed: " + response.statusCode());
byte[] bytes = body.readNBytes(128);
if (bytes.length == 128) throw new IOException("Node identity response is too large");
String[] parts = new String(bytes, java.nio.charset.StandardCharsets.US_ASCII).trim().split(" ", -1);
if (parts.length != 2) throw new IOException("Invalid node identity response");
return new NodeIdentity(UUID.fromString(parts[0]), UUID.fromString(parts[1]));
}
} catch (InterruptedException error) {
Thread.currentThread().interrupt();
throw new IOException("Interrupted during node identity request", error);
} catch (IllegalArgumentException error) {
throw new IOException("Invalid node identity response", error);
}
}
static void validateUrl(URI url) {
if (url == null || !"http".equals(url.getScheme()) || url.getHost() == null ||
url.getPort() < 1 || url.getRawUserInfo() != null ||
(url.getRawPath() != null && !url.getRawPath().isEmpty()) ||
url.getRawQuery() != null || url.getRawFragment() != null)
throw new IllegalArgumentException("Invalid private storage node URL");
}
boolean availableHostsAtLeast(int required, boolean testNodeDomains) {
Set<UUID> healthy = new HashSet<>();
for (int i = 0; i < nodes.size(); i++) {
Node node = nodes.get(i);
try {
NodeIdentity actual = probe(node.url(), token);
if (actual.nodeId().equals(node.id()) && actual.hostId().equals(node.hostId()))
healthy.add(faultDomain(i, testNodeDomains));
if (healthy.size() >= required) return true;
} catch (IOException error) { }
}
return false;
}
void put(int index, UUID id, byte[] data, byte[] sha256) throws IOException {
put(index, id, data, sha256, false);
}
void repair(int index, UUID id, byte[] data, byte[] sha256) throws IOException {
put(index, id, data, sha256, true);
}
private void put(int index, UUID id, byte[] data, byte[] sha256, boolean repair) throws IOException {
Node node = nodes.get(index);
HttpRequest.Builder builder = HttpRequest.newBuilder(node.url().resolve("/segments/" + id))
.timeout(Duration.ofSeconds(30)).header("X-Cluster-Token", token)
.header("X-Cluster-Expected-Node", node.id().toString())
.header("X-Cluster-Sha256", HexFormat.of().formatHex(sha256));
if (repair) {
if (repairToken == null || repairToken.length() < 32)
throw new IOException("Repair authority is not available to this process");
builder.header("X-Cluster-Repair", "true").header("X-Cluster-Repair-Token", repairToken);
}
HttpResponse<Void> response = send(builder.PUT(HttpRequest.BodyPublishers.ofByteArray(data)).build(),
HttpResponse.BodyHandlers.discarding());
if (response.statusCode() != 200) throw new IOException("Node " + node.id() + " rejected segment: " + response.statusCode());
}
byte[] get(int index, UUID id, int length, byte[] sha256) throws IOException {
if (length < 1 || length > ClusterNode.MAX_SEGMENT) throw new IOException("Invalid segment length");
Node node = nodes.get(index);
HttpRequest request = HttpRequest.newBuilder(node.url().resolve("/segments/" + id))
.timeout(Duration.ofSeconds(30)).header("X-Cluster-Token", token)
.header("X-Cluster-Expected-Node", node.id().toString()).GET().build();
HttpResponse<InputStream> response = send(request, HttpResponse.BodyHandlers.ofInputStream());
try (InputStream body = response.body()) {
if (response.statusCode() != 200)
throw new IOException("Node " + node.id() + " has no verified copy of segment " + id);
byte[] bytes = body.readNBytes(length + 1);
if (bytes.length != length || !MessageDigest.isEqual(SigV4.hash(bytes), sha256))
throw new IOException("Node " + node.id() + " has no verified copy of segment " + id);
return bytes;
}
}
private <T> HttpResponse<T> send(HttpRequest request, HttpResponse.BodyHandler<T> handler) throws IOException {
try { return http.send(request, handler); }
catch (InterruptedException error) {
Thread.currentThread().interrupt();
throw new IOException("Interrupted during node request", error);
}
}
}
@@ -0,0 +1,42 @@
package cloud.lunarsky.store;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.nio.channels.FileChannel;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.nio.file.StandardOpenOption;
import java.util.List;
import java.util.UUID;
record NodeIdentity(UUID nodeId, UUID hostId) {
static NodeIdentity open(Path root, UUID expectedHost) throws IOException {
Path file = root.resolve("node-identity");
if (Files.exists(file)) {
List<String> lines = Files.readAllLines(file, StandardCharsets.UTF_8);
if (lines.size() != 2) throw new IOException("Invalid node identity file");
try {
NodeIdentity identity = new NodeIdentity(UUID.fromString(lines.get(0)), UUID.fromString(lines.get(1)));
if (!identity.hostId().equals(expectedHost))
throw new IOException("Node volume belongs to a different storage host");
return identity;
} catch (IllegalArgumentException error) {
throw new IOException("Invalid node identity file", error);
}
}
NodeIdentity identity = new NodeIdentity(UUID.randomUUID(), expectedHost);
Path temporary = Files.createTempFile(root, ".node-identity-", ".pending");
try {
Files.writeString(temporary, identity.nodeId() + "\n" + identity.hostId() + "\n", StandardCharsets.UTF_8);
try (FileChannel channel = FileChannel.open(temporary, StandardOpenOption.WRITE)) {
channel.force(true);
}
Files.move(temporary, file, StandardCopyOption.ATOMIC_MOVE);
DiskStore.syncDirectory(root);
return identity;
} finally {
Files.deleteIfExists(temporary);
}
}
}
+135
View File
@@ -0,0 +1,135 @@
package cloud.lunarsky.store;
import java.io.IOException;
import java.net.URI;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
final class NodeRegistry {
private NodeRegistry() {}
static NodeClient.Node join(Connection connection, URI url, UUID expectedHost, String token) throws IOException {
NodeIdentity identity = NodeClient.probe(url, token);
if (!identity.hostId().equals(expectedHost))
throw new IOException("The node reported a different physical host ID");
try {
connection.setAutoCommit(false);
try (Statement statement = connection.createStatement()) {
statement.execute("SELECT pg_advisory_xact_lock(6834071092781)");
}
try (PreparedStatement query = connection.prepareStatement(
"SELECT node_id, host_id, endpoint FROM cluster_nodes WHERE node_id=? OR endpoint=?")) {
query.setObject(1, identity.nodeId());
query.setString(2, url.toString());
try (ResultSet result = query.executeQuery()) {
if (result.next()) {
if (!identity.nodeId().equals(result.getObject(1)) ||
!identity.hostId().equals(result.getObject(2)) || !url.toString().equals(result.getString(3)))
throw new IOException("Node ID, host ID, or endpoint conflicts with an existing registration");
if (result.next()) throw new IOException("Conflicting node registrations");
connection.commit();
return new NodeClient.Node(identity.nodeId(), identity.hostId(), url);
}
}
}
try (PreparedStatement insert = connection.prepareStatement(
"INSERT INTO cluster_nodes (node_id, host_id, endpoint, state) VALUES (?, ?, ?, 'active')")) {
insert.setObject(1, identity.nodeId());
insert.setObject(2, identity.hostId());
insert.setString(3, url.toString());
insert.executeUpdate();
}
connection.commit();
return new NodeClient.Node(identity.nodeId(), identity.hostId(), url);
} catch (SQLException | IOException error) {
try { connection.rollback(); } catch (SQLException rollback) { error.addSuppressed(rollback); }
if (error instanceof IOException io) throw io;
throw new IOException("Node registration failed", error);
} finally {
try { connection.setAutoCommit(true); }
catch (SQLException error) { throw new IOException("Could not restore metadata connection", error); }
}
}
static NodeClient load(Connection connection, List<URI> urls, String token, String repairToken) throws IOException {
if (urls.size() < 2 || urls.stream().distinct().count() != urls.size())
throw new IllegalArgumentException("At least two distinct node URLs are required");
try {
connection.setAutoCommit(false);
try (Statement statement = connection.createStatement()) {
statement.execute("SELECT pg_advisory_xact_lock(6834071092781)");
}
Map<String, NodeClient.Node> stored = new HashMap<>();
try (Statement statement = connection.createStatement();
ResultSet result = statement.executeQuery("SELECT node_id, host_id, endpoint FROM cluster_nodes WHERE state <> 'retired'")) {
while (result.next()) {
URI url = URI.create(result.getString(3));
stored.put(url.toString(), new NodeClient.Node((UUID) result.getObject(1),
(UUID) result.getObject(2), url));
}
}
if (stored.isEmpty()) {
try (Statement statement = connection.createStatement();
ResultSet result = statement.executeQuery("SELECT EXISTS (SELECT 1 FROM cluster_segments)")) {
result.next();
if (result.getBoolean(1)) throw new IOException("Existing segments have no registered node identities");
}
for (URI url : urls) {
NodeIdentity identity = NodeClient.probe(url, token);
try (PreparedStatement insert = connection.prepareStatement(
"INSERT INTO cluster_nodes (node_id, host_id, endpoint, state) VALUES (?, ?, ?, 'active')")) {
insert.setObject(1, identity.nodeId());
insert.setObject(2, identity.hostId());
insert.setString(3, url.toString());
insert.executeUpdate();
}
stored.put(url.toString(), new NodeClient.Node(identity.nodeId(), identity.hostId(), url));
}
}
List<NodeClient.Node> configured = new ArrayList<>();
Set<UUID> configuredIds = new HashSet<>();
for (URI url : urls) {
NodeClient.Node node = stored.get(url.toString());
if (node == null) throw new IOException("Unregistered storage node URL: " + url);
NodeIdentity actual = null;
try {
actual = NodeClient.probe(url, token);
} catch (IOException offline) { }
if (actual != null && (!actual.nodeId().equals(node.id()) || !actual.hostId().equals(node.hostId())))
throw new IOException("Storage node identity changed at " + url);
configured.add(node);
configuredIds.add(node.id());
}
try (Statement statement = connection.createStatement();
ResultSet result = statement.executeQuery(
"SELECT DISTINCT unnest(s.replica_ids) FROM cluster_segments s JOIN cluster_objects o ON o.generation=s.generation")) {
while (result.next()) {
UUID id = (UUID) result.getObject(1);
if (!configuredIds.contains(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); }
}
}
}
@@ -0,0 +1,26 @@
package cloud.lunarsky.store;
import java.io.IOException;
import java.io.InputStream;
import java.util.List;
/** Storage operations shared by the local and cluster gateways. */
interface ObjectStorage extends AutoCloseable {
record Metadata(long length, long modified, String etag, byte[] sha256,
String bucket, String key, String contentType) {}
record OpenObject(Metadata metadata, InputStream stream) implements AutoCloseable {
public void close() throws IOException { stream.close(); }
}
record ListedObject(String key, Metadata metadata) {}
record ListPage(List<ListedObject> objects, List<String> prefixes, String nextKey, boolean truncated) {
int keyCount() { return objects.size() + prefixes.size(); }
}
Metadata put(String bucket, String key, InputStream input, long length, String expectedHash,
String checksum, boolean createOnly, String contentType) throws IOException;
OpenObject open(String bucket, String key) throws IOException;
void delete(String bucket, String key) throws IOException;
ListPage list(String bucket, String prefix, String delimiter, int maxKeys, String after) throws IOException;
default boolean ready() { return true; }
void close() throws IOException;
}
@@ -0,0 +1,38 @@
package cloud.lunarsky.store;
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;
final class PlacementPolicy {
private PlacementPolicy() {}
static List<Integer> candidates(UUID segmentId, NodeClient nodes, boolean testNodeDomains) {
Map<UUID, List<Integer>> byHost = new HashMap<>();
for (int i = 0; i < nodes.count(); i++)
byHost.computeIfAbsent(nodes.faultDomain(i, testNodeDomains), ignored -> new ArrayList<>()).add(i);
List<UUID> hosts = new ArrayList<>(byHost.keySet());
hosts.sort(Comparator.comparingLong((UUID host) -> score(segmentId, host)).reversed());
int longest = 0;
for (List<Integer> group : byHost.values()) {
group.sort(Comparator.comparingLong((Integer index) -> score(segmentId, nodes.node(index).id())).reversed());
longest = Math.max(longest, group.size());
}
List<Integer> order = new ArrayList<>(nodes.count());
for (int round = 0; round < longest; round++)
for (UUID host : hosts)
if (round < byHost.get(host).size()) order.add(byHost.get(host).get(round));
return order;
}
private static long score(UUID segment, UUID candidate) {
ByteBuffer bytes = ByteBuffer.allocate(32);
bytes.putLong(segment.getMostSignificantBits()).putLong(segment.getLeastSignificantBits());
bytes.putLong(candidate.getMostSignificantBits()).putLong(candidate.getLeastSignificantBits());
return ByteBuffer.wrap(SigV4.hash(bytes.array())).getLong();
}
}
@@ -0,0 +1,71 @@
package cloud.lunarsky.store;
import java.io.IOException;
import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
final class SchemaMigrator {
private SchemaMigrator() {}
static int prepare(Connection connection, String bucket) throws IOException {
try {
connection.setAutoCommit(false);
try (Statement statement = connection.createStatement()) {
statement.execute("SELECT pg_advisory_xact_lock(6834071092781)");
statement.execute("CREATE TABLE IF NOT EXISTS cluster_schema_migrations (version integer PRIMARY KEY)");
int version;
try (ResultSet result = statement.executeQuery("SELECT COALESCE(MAX(version), 0) FROM cluster_schema_migrations")) {
result.next();
version = result.getInt(1);
}
if (version > 2) 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))");
statement.execute("CREATE TABLE IF NOT EXISTS cluster_segments (generation uuid NOT NULL, ordinal integer NOT NULL, segment_id uuid NOT NULL, length integer NOT NULL, sha256 bytea NOT NULL, replicas text NOT NULL, PRIMARY KEY (generation, ordinal))");
statement.execute("CREATE TABLE IF NOT EXISTS cluster_tombstones (bucket text NOT NULL, object_key text COLLATE \"C\" NOT NULL, generation uuid NOT NULL, deleted_at bigint NOT NULL, PRIMARY KEY (bucket, object_key))");
statement.execute("INSERT INTO cluster_schema_migrations VALUES (1)");
}
if (version < 2) {
statement.execute("ALTER TABLE cluster_segments ADD COLUMN IF NOT EXISTS replica_ids uuid[]");
statement.execute("ALTER TABLE cluster_segments ADD COLUMN IF NOT EXISTS placement_version bigint NOT NULL DEFAULT 0");
statement.execute("ALTER TABLE cluster_segments ADD CONSTRAINT cluster_replica_ids_required CHECK (replica_ids IS NOT NULL) NOT VALID");
statement.execute("CREATE TABLE IF NOT EXISTS cluster_nodes (node_id uuid PRIMARY KEY, host_id uuid NOT NULL, endpoint text NOT NULL UNIQUE, legacy_index integer UNIQUE, state text NOT NULL CHECK (state IN ('joining','active','draining','offline','retired')))");
statement.execute("CREATE TABLE IF NOT EXISTS cluster_format (singleton integer PRIMARY KEY CHECK (singleton=1), version integer NOT NULL)");
statement.execute("INSERT INTO cluster_schema_migrations VALUES (2)");
}
statement.execute("INSERT INTO cluster_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")) {
insert.setString(1, bucket);
insert.executeUpdate();
}
int format;
try (Statement statement = connection.createStatement();
ResultSet result = statement.executeQuery("SELECT version FROM cluster_format WHERE singleton=1")) {
if (!result.next()) throw new IOException("Missing cluster format marker");
format = result.getInt(1);
}
if (format < 1 || format > 2) throw new IOException("Unsupported cluster data format " + format);
if (format == 2) {
try (Statement statement = connection.createStatement();
ResultSet result = statement.executeQuery("SELECT EXISTS (SELECT 1 FROM cluster_segments WHERE replica_ids IS NULL)")) {
result.next();
if (result.getBoolean(1)) throw new IOException("Cluster format has unmigrated segment replicas");
}
}
connection.commit();
return format;
} catch (SQLException | IOException error) {
try { connection.rollback(); } catch (SQLException rollback) { error.addSuppressed(rollback); }
if (error instanceof IOException io) throw io;
throw new IOException("Metadata schema migration failed", error);
} finally {
try { connection.setAutoCommit(true); }
catch (SQLException error) { throw new IOException("Could not restore metadata connection", error); }
}
}
}
+4 -1
View File
@@ -4,7 +4,10 @@ final class StoreException extends RuntimeException {
final int status;
final String code;
StoreException(int status, String code, String message) {
super(message);
this(status, code, message, null);
}
StoreException(int status, String code, String message, Throwable cause) {
super(message, cause);
this.status = status;
this.code = code;
}
@@ -0,0 +1,17 @@
package cloud.lunarsky.store;
import java.io.InputStream;
import java.util.List;
final class UnavailableMultipart implements MultipartStorage {
private StoreException unavailable() {
return new StoreException(501, "NotImplemented", "Multipart uploads are unavailable in the local cluster prototype");
}
@Override public String create(String bucket, String key, String contentType) { throw unavailable(); }
@Override public String putPart(String id, String bucket, String key, int number, InputStream input,
long length, String expectedHash, String checksum) { throw unavailable(); }
@Override public ObjectStorage.Metadata complete(String id, String bucket, String key, List<Part> parts) { throw unavailable(); }
@Override public void abort(String id, String bucket, String key) { throw unavailable(); }
@Override public int activeUploads() { return 0; }
@Override public long stagedBytes() { return 0; }
}
+7
View File
@@ -0,0 +1,7 @@
package cloud.lunarsky.store;
final class Version {
static final String VALUE = "0.0.2";
private Version() {}
}