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.Comparator; import java.util.HexFormat; import java.util.List; 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; } if (!authorized(exchange)) return; if (path.equals("/identity") && exchange.getRequestMethod().equals("GET")) { respond(exchange, 200, identity.nodeId() + " " + identity.hostId()); return; } if (!expectedNodeAndRepairAuthorized(exchange)) return; if (path.equals("/segments") && exchange.getRequestMethod().equals("GET")) { if (maintenanceAuthorized(exchange)) inventory(exchange); return; } String id = segmentId(exchange, path); if (id == null) return; switch (exchange.getRequestMethod()) { case "PUT" -> put(exchange, segmentPath(id, true)); case "GET" -> get(exchange, segmentPath(id, false)); case "DELETE" -> { if (maintenanceAuthorized(exchange)) delete(exchange, segmentPath(id, false)); } default -> respond(exchange, 405, "Method not allowed"); } } catch (IllegalArgumentException error) { respond(exchange, 400, "Invalid request"); } catch (IOException error) { if (exchange.getResponseCode() == -1) respond(exchange, 500, "Storage failure"); throw error; } finally { exchange.close(); } } private boolean authorized(HttpExchange exchange) throws IOException { byte[] supplied = exchange.getRequestHeaders().getFirst("X-Cluster-Token") == null ? new byte[0] : exchange.getRequestHeaders().getFirst("X-Cluster-Token") .getBytes(java.nio.charset.StandardCharsets.UTF_8); if (!MessageDigest.isEqual(token, supplied)) { respond(exchange, 403, "Forbidden"); return false; } return true; } private boolean expectedNodeAndRepairAuthorized(HttpExchange exchange) throws IOException { if (!identity.nodeId().toString().equals(exchange.getRequestHeaders().getFirst("X-Cluster-Expected-Node"))) { respond(exchange, 409, "Wrong storage node"); return false; } if (exchange.getRequestMethod().equals("PUT") && "true".equals(exchange.getRequestHeaders().getFirst("X-Cluster-Repair"))) { String suppliedRepair = exchange.getRequestHeaders().getFirst("X-Cluster-Repair-Token"); byte[] suppliedBytes = suppliedRepair == null ? new byte[0] : suppliedRepair.getBytes(java.nio.charset.StandardCharsets.UTF_8); if (!MessageDigest.isEqual(repairToken, suppliedBytes)) { respond(exchange, 403, "Repair authority required"); return false; } } return true; } private boolean maintenanceAuthorized(HttpExchange exchange) throws IOException { String supplied = exchange.getRequestHeaders().getFirst("X-Cluster-Repair-Token"); byte[] value = supplied == null ? new byte[0] : supplied.getBytes(java.nio.charset.StandardCharsets.UTF_8); if (MessageDigest.isEqual(repairToken, value)) return true; respond(exchange, 403, "Repair authority required"); return false; } private void inventory(HttpExchange exchange) throws IOException { String query = exchange.getRequestURI().getRawQuery(); if (query == null || !query.matches("shard=[0-9a-f]{2}(&after=[0-9a-f-]{36})?")) { respond(exchange, 400, "Invalid inventory request"); return; } String shard = query.substring(6, 8); String after = query.length() > 8 ? query.substring(15) : ""; if (!after.isEmpty() && (!after.startsWith(shard) || !UUID.fromString(after).toString().equals(after))) { respond(exchange, 400, "Invalid inventory cursor"); return; } Path directory = segments.resolve(shard); if (!Files.isDirectory(directory)) { respond(exchange, 200, ""); return; } List files; try (var entries = Files.list(directory)) { files = entries.filter(Files::isRegularFile).sorted(Comparator.comparing(path -> path.getFileName().toString())).toList(); } StringBuilder body = new StringBuilder(); int count = 0; for (Path file : files) { String id = file.getFileName().toString(); if (id.compareTo(after) <= 0 || !id.matches("[0-9a-f-]{36}")) continue; if (!UUID.fromString(id).toString().equals(id)) continue; body.append(id).append(' ').append(Files.getLastModifiedTime(file).toMillis()).append('\n'); if (++count == 1000) break; } respond(exchange, 200, body.toString()); } private synchronized void delete(HttpExchange exchange, Path target) throws IOException { String expected = exchange.getRequestHeaders().getFirst("X-Cluster-Expected-Mtime"); String age = exchange.getRequestHeaders().getFirst("X-Cluster-Gc-Min-Age-Millis"); long expectedTime, minimumAge; try { expectedTime = Long.parseLong(expected); minimumAge = Long.parseLong(age); } catch (NumberFormatException error) { respond(exchange, 400, "Invalid deletion guard"); return; } if (minimumAge < 0 || minimumAge > System.currentTimeMillis()) { respond(exchange, 400, "Invalid deletion age"); return; } if (!Files.isRegularFile(target)) { exchange.sendResponseHeaders(404, -1); return; } long modified = Files.getLastModifiedTime(target).toMillis(); if (modified != expectedTime || modified > System.currentTimeMillis() - minimumAge) { exchange.sendResponseHeaders(409, -1); return; } Files.delete(target); DiskStore.syncDirectory(target.getParent()); exchange.sendResponseHeaders(204, -1); } private static String segmentId(HttpExchange exchange, String path) throws IOException { if (!path.matches("/segments/[0-9a-f-]{36}")) { respond(exchange, 404, "Not found"); return null; } String id = path.substring("/segments/".length()); if (!UUID.fromString(id).toString().equals(id)) { respond(exchange, 400, "Invalid segment ID"); return null; } return id; } private synchronized Path segmentPath(String id, boolean createShard) throws IOException { 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 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); } }