From 11bfe71124b7e5f43ceb9f7e7c38ff93a686b1b3 Mon Sep 17 00:00:00 2001 From: LunarSkyOSS Date: Sat, 10 Oct 2026 08:30:28 +0200 Subject: [PATCH] Put Java statements on separate lines --- src/cloud/lunarsky/store/Cli.java | 4 +- src/cloud/lunarsky/store/ClusterMigrate.java | 3 +- src/cloud/lunarsky/store/ClusterNode.java | 3 +- src/cloud/lunarsky/store/ClusterStore.java | 86 ++++++++++++++----- src/cloud/lunarsky/store/DiskStore.java | 43 +++++++--- src/cloud/lunarsky/store/Main.java | 25 ++++-- src/cloud/lunarsky/store/MultipartStore.java | 21 +++-- src/cloud/lunarsky/store/NodeClient.java | 4 +- src/cloud/lunarsky/store/NodeRegistry.java | 6 +- src/cloud/lunarsky/store/SchemaMigrator.java | 3 +- src/cloud/lunarsky/store/SigV4.java | 42 ++++++--- .../store/ClusterIntegrationTest.java | 25 ++++-- .../lunarsky/store/ClusterMigrationTest.java | 26 ++++-- test/cloud/lunarsky/store/HttpTest.java | 3 +- test/cloud/lunarsky/store/StoreTest.java | 19 ++-- 15 files changed, 228 insertions(+), 85 deletions(-) diff --git a/src/cloud/lunarsky/store/Cli.java b/src/cloud/lunarsky/store/Cli.java index dc690ec..c021f58 100644 --- a/src/cloud/lunarsky/store/Cli.java +++ b/src/cloud/lunarsky/store/Cli.java @@ -109,7 +109,9 @@ public final class Cli { } if (verify) { MessageDigest sha = digest("SHA-256"), md5 = digest("MD5"); - byte[] buffer = new byte[65536]; long count = 0; int n; + byte[] buffer = new byte[65536]; + long count = 0; + int n; while ((n = input.read(buffer)) != -1) { count += n; sha.update(buffer, 0, n); diff --git a/src/cloud/lunarsky/store/ClusterMigrate.java b/src/cloud/lunarsky/store/ClusterMigrate.java index c8ace99..073622b 100644 --- a/src/cloud/lunarsky/store/ClusterMigrate.java +++ b/src/cloud/lunarsky/store/ClusterMigrate.java @@ -148,7 +148,8 @@ public final class ClusterMigrate { } connection.commit(); } catch (SQLException | IOException | RuntimeException error) { - try { connection.rollback(); } catch (SQLException rollback) { error.addSuppressed(rollback); } + 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; diff --git a/src/cloud/lunarsky/store/ClusterNode.java b/src/cloud/lunarsky/store/ClusterNode.java index 22d67be..de0b9ae 100644 --- a/src/cloud/lunarsky/store/ClusterNode.java +++ b/src/cloud/lunarsky/store/ClusterNode.java @@ -252,7 +252,8 @@ public final class ClusterNode implements AutoCloseable { Runtime.getRuntime().addShutdownHook(new Thread(() -> { server.stop(5); executor.close(); - try { node.close(); } catch (IOException error) { System.err.println("Node close failed: " + error); } + try { node.close(); } + catch (IOException error) { System.err.println("Node close failed: " + error); } })); server.start(); System.out.println("ObjectStore cluster node listening on :" + port); diff --git a/src/cloud/lunarsky/store/ClusterStore.java b/src/cloud/lunarsky/store/ClusterStore.java index c878328..c6ebafd 100644 --- a/src/cloud/lunarsky/store/ClusterStore.java +++ b/src/cloud/lunarsky/store/ClusterStore.java @@ -38,8 +38,12 @@ final class ClusterStore implements ObjectStorage { 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.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); @@ -88,7 +92,8 @@ final class ClusterStore implements ObjectStorage { 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); + sha.update(buffer, 0, count); + md5.update(buffer, 0, count); output.write(buffer, 0, count); remaining -= count; } @@ -169,8 +174,11 @@ final class ClusterStore implements ObjectStorage { "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.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(); } @@ -178,13 +186,18 @@ final class ClusterStore implements ObjectStorage { } 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(); + 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(); + 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(); + delete.setString(1, bucket); + delete.setString(2, key); + delete.executeUpdate(); } connection.commit(); } catch (SQLException | RuntimeException error) { @@ -204,7 +217,8 @@ final class ClusterStore implements ObjectStorage { 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); + 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); @@ -245,17 +259,23 @@ final class ClusterStore implements ObjectStorage { 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(); + 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.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(); + update.setLong(1, used - previous); + update.setString(2, bucket); + update.executeUpdate(); } } connection.commit(); @@ -297,10 +317,18 @@ final class ClusterStore implements ObjectStorage { if (!key.startsWith(prefix)) break; if (after != null && key.compareTo(after) <= 0) continue; String group = commonPrefix(key, prefix, delimiter); - 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 { + 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; @@ -327,14 +355,20 @@ final class ClusterStore implements ObjectStorage { } 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); + 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()); + 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 replicaIds(ResultSet result, int column) throws SQLException, IOException { java.sql.Array value = result.getArray(column); @@ -406,7 +440,10 @@ final class ClusterStore implements ObjectStorage { healthy.add(id); healthyHosts.add(nodes.faultDomain(node, testNodeDomains)); } - if (copy == null) { unrecoverable++; continue; } + 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; @@ -484,6 +521,9 @@ final class ClusterStore implements ObjectStorage { } return current.read(buffer, offset, length); } - @Override public void close() { closed = true; current = null; } + @Override public void close() { + closed = true; + current = null; + } } } diff --git a/src/cloud/lunarsky/store/DiskStore.java b/src/cloud/lunarsky/store/DiskStore.java index d5bf450..bb04d72 100644 --- a/src/cloud/lunarsky/store/DiskStore.java +++ b/src/cloud/lunarsky/store/DiskStore.java @@ -34,8 +34,10 @@ final class DiskStore implements ObjectStorage { 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; + objects = root.resolve("objects"); + temporary = root.resolve("pending"); + this.maxObject = maxObject; + this.maxTotal = maxTotal; Arrays.setAll(locks, i -> new Object()); Files.createDirectories(root); FileChannel channel = FileChannel.open(root.resolve(".process.lock"), StandardOpenOption.CREATE, StandardOpenOption.WRITE); @@ -45,7 +47,8 @@ final class DiskStore implements ObjectStorage { try { acquired = channel.tryLock(); } 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); + Files.createDirectories(objects); + Files.createDirectories(temporary); syncDirectory(root); try (var paths = Files.list(temporary)) { for (Path p : paths.toList()) if (p.getFileName().toString().endsWith(".part")) Files.delete(p); @@ -140,12 +143,15 @@ final class DiskStore implements ObjectStorage { long count = 0; try (OutputStream out = Files.newOutputStream(pending)) { out.write(new byte[headerLength]); - byte[] buffer = new byte[65536]; int n; + 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); + 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"); @@ -200,7 +206,10 @@ final class DiskStore implements ObjectStorage { 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; } + catch (IOException e) { + input.close(); + throw e; + } } } @@ -242,10 +251,21 @@ final class DiskStore implements ObjectStorage { 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; } + 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); @@ -256,7 +276,8 @@ final class DiskStore implements ObjectStorage { 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); + 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, diff --git a/src/cloud/lunarsky/store/Main.java b/src/cloud/lunarsky/store/Main.java index 867798f..afac05a 100644 --- a/src/cloud/lunarsky/store/Main.java +++ b/src/cloud/lunarsky/store/Main.java @@ -34,8 +34,10 @@ public final class Main { } Main(ObjectStorage store, MultipartStorage multipart, SigV4 authentication, String bucket) { - this.store = store; this.multipart = multipart; - this.authentication = authentication; this.bucket = bucket; + this.store = store; + this.multipart = multipart; + this.authentication = authentication; + this.bucket = bucket; } void handle(HttpExchange exchange) throws IOException { @@ -58,7 +60,10 @@ public final class Main { 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(); } + } finally { + if (admitted) slots.release(); + exchange.close(); + } } private boolean handleStatus(HttpExchange exchange) throws IOException { @@ -321,7 +326,8 @@ public final class Main { else { object.stream().skipNBytes(range.start()); exchange.sendResponseHeaders(status, range.length()); - byte[] buffer = new byte[65536]; long left = 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"); @@ -345,7 +351,8 @@ public final class Main { 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; + 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); @@ -494,7 +501,8 @@ public final class Main { ? 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); + server.setExecutor(executor); + server.createContext("/", app::handle); Runtime.getRuntime().addShutdownHook(new Thread(() -> { server.stop(5); executor.close(); @@ -526,7 +534,10 @@ public final class Main { String[] units = {"KiB", "MiB", "GiB", "TiB"}; double value = bytes; int unit = -1; - do { value /= 1024; unit++; } while (value >= 1024 && unit < units.length - 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 env, String key) { diff --git a/src/cloud/lunarsky/store/MultipartStore.java b/src/cloud/lunarsky/store/MultipartStore.java index 5762d67..a0068aa 100644 --- a/src/cloud/lunarsky/store/MultipartStore.java +++ b/src/cloud/lunarsky/store/MultipartStore.java @@ -97,11 +97,14 @@ final class MultipartStore implements MultipartStorage { MessageDigest sha = digest("SHA-256"), md5 = digest("MD5"); long count = 0; try (var output = Files.newOutputStream(pending)) { - byte[] buffer = new byte[65536]; int n; + 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); + 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"); @@ -124,7 +127,8 @@ final class MultipartStore implements MultipartStorage { throw new StoreException(400, "InvalidPart", "No valid parts supplied"); MessageDigest sha = digest("SHA-256"); List paths = new ArrayList<>(); - long total = 0; int last = 0; + 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"); @@ -136,8 +140,12 @@ final class MultipartStore implements MultipartStorage { 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); } + 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"); @@ -209,7 +217,8 @@ final class MultipartStore implements MultipartStorage { } int n = current.read(buffer, offset, length); if (n >= 0) return n; - current.close(); current = null; + current.close(); + current = null; } } @Override public void close() throws IOException { if (current != null) current.close(); } diff --git a/src/cloud/lunarsky/store/NodeClient.java b/src/cloud/lunarsky/store/NodeClient.java index 952f254..91e12b5 100644 --- a/src/cloud/lunarsky/store/NodeClient.java +++ b/src/cloud/lunarsky/store/NodeClient.java @@ -43,7 +43,9 @@ final class NodeClient { List 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; + 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) { diff --git a/src/cloud/lunarsky/store/NodeRegistry.java b/src/cloud/lunarsky/store/NodeRegistry.java index 931a223..804d303 100644 --- a/src/cloud/lunarsky/store/NodeRegistry.java +++ b/src/cloud/lunarsky/store/NodeRegistry.java @@ -52,7 +52,8 @@ final class NodeRegistry { connection.commit(); return new NodeClient.Node(identity.nodeId(), identity.hostId(), url); } catch (SQLException | IOException error) { - try { connection.rollback(); } catch (SQLException rollback) { error.addSuppressed(rollback); } + try { connection.rollback(); } + catch (SQLException rollback) { error.addSuppressed(rollback); } if (error instanceof IOException io) throw io; throw new IOException("Node registration failed", error); } finally { @@ -77,7 +78,8 @@ final class NodeRegistry { connection.commit(); return nodes; } catch (SQLException | IOException | RuntimeException error) { - try { connection.rollback(); } catch (SQLException rollback) { error.addSuppressed(rollback); } + 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; diff --git a/src/cloud/lunarsky/store/SchemaMigrator.java b/src/cloud/lunarsky/store/SchemaMigrator.java index 40ef88d..f52398a 100644 --- a/src/cloud/lunarsky/store/SchemaMigrator.java +++ b/src/cloud/lunarsky/store/SchemaMigrator.java @@ -60,7 +60,8 @@ final class SchemaMigrator { connection.commit(); return format; } catch (SQLException | IOException error) { - try { connection.rollback(); } catch (SQLException rollback) { error.addSuppressed(rollback); } + 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 { diff --git a/src/cloud/lunarsky/store/SigV4.java b/src/cloud/lunarsky/store/SigV4.java index 8cde1d7..c78eaf7 100644 --- a/src/cloud/lunarsky/store/SigV4.java +++ b/src/cloud/lunarsky/store/SigV4.java @@ -24,7 +24,10 @@ final class SigV4 { private final Clock clock; SigV4(String accessKey, String secretKey, String region, Clock clock) { - this.accessKey = accessKey; this.secretKey = secretKey; this.region = region; this.clock = clock; + this.accessKey = accessKey; + this.secretKey = secretKey; + this.region = region; + this.clock = clock; } String verify(String method, URI uri, Headers headers) { @@ -111,17 +114,25 @@ final class SigV4 { static String decode(String value) { try { - var bytes=new java.io.ByteArrayOutputStream(); - for(int i=0;i=value.length())throw new IllegalArgumentException(); - int hi=Character.digit(value.charAt(i+1),16),lo=Character.digit(value.charAt(i+2),16); - if(hi<0||lo<0)throw new IllegalArgumentException(); - bytes.write((hi<<4)|lo);i+=3; - }else{int point=value.codePointAt(i);bytes.writeBytes(new String(Character.toChars(point)).getBytes(StandardCharsets.UTF_8));i+=Character.charCount(point);} + var bytes = new java.io.ByteArrayOutputStream(); + for (int i = 0; i < value.length();) { + if (value.charAt(i) == '%') { + if (i + 2 >= value.length()) throw new IllegalArgumentException(); + int hi = Character.digit(value.charAt(i + 1), 16); + int lo = Character.digit(value.charAt(i + 2), 16); + if (hi < 0 || lo < 0) throw new IllegalArgumentException(); + bytes.write((hi << 4) | lo); + i += 3; + } else { + int point = value.codePointAt(i); + bytes.writeBytes(new String(Character.toChars(point)).getBytes(StandardCharsets.UTF_8)); + i += Character.charCount(point); + } } return StandardCharsets.UTF_8.newDecoder().onMalformedInput(java.nio.charset.CodingErrorAction.REPORT).decode(java.nio.ByteBuffer.wrap(bytes.toByteArray())).toString(); - }catch(IllegalArgumentException|java.nio.charset.CharacterCodingException e){throw new StoreException(400,"InvalidURI","Malformed URI encoding");} + } catch (IllegalArgumentException | java.nio.charset.CharacterCodingException e) { + throw new StoreException(400, "InvalidURI", "Malformed URI encoding"); + } } static String encode(String value, boolean keepSlash) { @@ -146,10 +157,17 @@ final class SigV4 { return hmac(hmac(hmac(hmac(("AWS4"+secret).getBytes(StandardCharsets.UTF_8),date),region),"s3"),"aws4_request"); } static byte[] hmac(byte[] key, String text) { - try { Mac mac=Mac.getInstance("HmacSHA256");mac.init(new SecretKeySpec(key,"HmacSHA256"));return mac.doFinal(text.getBytes(StandardCharsets.UTF_8)); } + try { + Mac mac = Mac.getInstance("HmacSHA256"); + mac.init(new SecretKeySpec(key, "HmacSHA256")); + return mac.doFinal(text.getBytes(StandardCharsets.UTF_8)); + } catch (java.security.GeneralSecurityException e) { throw new IllegalStateException(e); } } - static byte[] hash(byte[] data) { try {return MessageDigest.getInstance("SHA-256").digest(data);}catch(java.security.NoSuchAlgorithmException e){throw new IllegalStateException(e);} } + static byte[] hash(byte[] data) { + try { return MessageDigest.getInstance("SHA-256").digest(data); } + catch (java.security.NoSuchAlgorithmException e) { throw new IllegalStateException(e); } + } static String hex(byte[] data) { return HexFormat.of().formatHex(data); } private static void denied(String message) { throw new StoreException(403,"AccessDenied",message); } } diff --git a/test/cloud/lunarsky/store/ClusterIntegrationTest.java b/test/cloud/lunarsky/store/ClusterIntegrationTest.java index 61161b1..f93c8ab 100644 --- a/test/cloud/lunarsky/store/ClusterIntegrationTest.java +++ b/test/cloud/lunarsky/store/ClusterIntegrationTest.java @@ -21,7 +21,9 @@ public final class ClusterIntegrationTest { switch (args[0]) { case "basic" -> { byte[] large = new byte[ClusterNode.MAX_SEGMENT + 37]; - for (int i = 0; i < large.length; i++) large[i] = (byte) (i * 31); + for (int i = 0; i < large.length; i++) { + large[i] = (byte) (i * 31); + } String largeKey = "cluster-test/large"; put(store, bucket, largeKey, large, false); try (var opened = store.open(bucket, largeKey)) { @@ -39,7 +41,10 @@ public final class ClusterIntegrationTest { "Overwrite was not visible"); } store.delete(bucket, largeKey); - try { store.open(bucket, largeKey); throw new AssertionError("Deleted object remained visible"); } + try { + store.open(bucket, largeKey); + throw new AssertionError("Deleted object remained visible"); + } catch (StoreException error) { require(error.status == 404, "Wrong missing-object status"); } put(store, bucket, KEY, stable, false); require(store.ready(), "Healthy cluster is not ready"); @@ -63,7 +68,10 @@ public final class ClusterIntegrationTest { put(store, bucket, "cluster-test/rejected", new byte[]{1}, false); throw new AssertionError("Write succeeded with only one node"); } catch (StoreException error) { require(error.status == 503, "Wrong unavailable status"); } - try { store.open(bucket, "cluster-test/rejected"); throw new AssertionError("Failed write became visible"); } + try { + store.open(bucket, "cluster-test/rejected"); + throw new AssertionError("Failed write became visible"); + } catch (StoreException error) { require(error.status == 404, "Partial object became visible"); } System.out.println("Cluster quorum-loss test passed"); } @@ -82,13 +90,18 @@ public final class ClusterIntegrationTest { var executor = java.util.concurrent.Executors.newFixedThreadPool(2); try { var a = executor.submit(() -> { - start.await(); put(store, bucket, key, first, false); return null; + start.await(); + put(store, bucket, key, first, false); + return null; }); var b = executor.submit(() -> { - start.await(); put(store, bucket, key, second, false); return null; + start.await(); + put(store, bucket, key, second, false); + return null; }); start.countDown(); - a.get(); b.get(); + a.get(); + b.get(); try (var opened = store.open(bucket, key)) { byte[] actual = opened.stream().readAllBytes(); require(Arrays.equals(actual, first) || Arrays.equals(actual, second), diff --git a/test/cloud/lunarsky/store/ClusterMigrationTest.java b/test/cloud/lunarsky/store/ClusterMigrationTest.java index f07454d..70b226b 100644 --- a/test/cloud/lunarsky/store/ClusterMigrationTest.java +++ b/test/cloud/lunarsky/store/ClusterMigrationTest.java @@ -31,7 +31,9 @@ public final class ClusterMigrationTest { if (args[0].equals("create")) { UUID segment = UUID.randomUUID(); byte[] hash = SigV4.hash(DATA); - for (int i = 0; i < nodes.count(); i++) nodes.put(i, segment, DATA, hash); + for (int i = 0; i < nodes.count(); i++) { + nodes.put(i, segment, DATA, hash); + } try (var connection = DriverManager.getConnection(env.get("POSTGRES_JDBC_URL"), env.get("POSTGRES_USER"), env.get("POSTGRES_PASSWORD")); var statement = connection.createStatement()) { @@ -40,18 +42,28 @@ public final class ClusterMigrationTest { statement.execute("CREATE TABLE 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 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))"); try (var insert = connection.prepareStatement("INSERT INTO cluster_usage VALUES (?, ?)")) { - insert.setString(1, bucket); insert.setLong(2, DATA.length); insert.executeUpdate(); + insert.setString(1, bucket); + insert.setLong(2, DATA.length); + insert.executeUpdate(); } UUID generation = UUID.randomUUID(); try (var insert = connection.prepareStatement("INSERT INTO cluster_objects VALUES (?, ?, ?, ?, ?, ?, ?, ?)")) { - insert.setString(1, bucket); insert.setString(2, KEY); insert.setObject(3, generation); - insert.setLong(4, DATA.length); insert.setLong(5, System.currentTimeMillis()); + insert.setString(1, bucket); + insert.setString(2, KEY); + insert.setObject(3, generation); + insert.setLong(4, DATA.length); + insert.setLong(5, System.currentTimeMillis()); insert.setString(6, HexFormat.of().formatHex(MessageDigest.getInstance("MD5").digest(DATA))); - insert.setBytes(7, hash); insert.setString(8, "text/plain"); insert.executeUpdate(); + insert.setBytes(7, hash); + insert.setString(8, "text/plain"); + insert.executeUpdate(); } try (var insert = connection.prepareStatement("INSERT INTO cluster_segments VALUES (?, 0, ?, ?, ?, '0,1,2')")) { - insert.setObject(1, generation); insert.setObject(2, segment); - insert.setInt(3, DATA.length); insert.setBytes(4, hash); insert.executeUpdate(); + insert.setObject(1, generation); + insert.setObject(2, segment); + insert.setInt(3, DATA.length); + insert.setBytes(4, hash); + insert.executeUpdate(); } } System.out.println("Legacy cluster fixture created"); diff --git a/test/cloud/lunarsky/store/HttpTest.java b/test/cloud/lunarsky/store/HttpTest.java index 234d20a..96bd107 100644 --- a/test/cloud/lunarsky/store/HttpTest.java +++ b/test/cloud/lunarsky/store/HttpTest.java @@ -142,7 +142,8 @@ public final class HttpTest { "?partNumber=1&uploadId=" + upload), "PUT", first, Map.of()), HttpResponse.BodyHandlers.ofByteArray()); var partTwo = client.send(signedUri(URI.create(base + "/objects/" + movie + "?partNumber=2&uploadId=" + upload), "PUT", second, Map.of()), HttpResponse.BodyHandlers.ofByteArray()); - status(200, partOne); status(200, partTwo); + status(200, partOne); + status(200, partTwo); String completion = "1" + partOne.headers().firstValue("etag").orElseThrow() + "2" + diff --git a/test/cloud/lunarsky/store/StoreTest.java b/test/cloud/lunarsky/store/StoreTest.java index 9b1b94e..3dff481 100644 --- a/test/cloud/lunarsky/store/StoreTest.java +++ b/test/cloud/lunarsky/store/StoreTest.java @@ -9,25 +9,33 @@ import java.util.List; public final class StoreTest { interface Operation {void run() throws Exception;} static void fails(int status,Operation operation)throws Exception{ - try{operation.run();throw new AssertionError("Expected "+status);}catch(StoreException error){if(error.status!=status)throw error;} + try { + operation.run(); + throw new AssertionError("Expected " + status); + } catch (StoreException error) { + if (error.status != status) throw error; + } } static ObjectStorage.Metadata put(DiskStore store,String key,byte[] body,boolean only)throws Exception{ return store.put("test",key,new ByteArrayInputStream(body),body.length,SigV4.hex(SigV4.hash(body)),null,only,"application/octet-stream"); } private static void testSignature() throws Exception { var headers=new com.sun.net.httpserver.Headers(); - headers.set("host","examplebucket.s3.amazonaws.com");headers.set("range","bytes=0-9"); + headers.set("host","examplebucket.s3.amazonaws.com"); + headers.set("range","bytes=0-9"); headers.set("x-amz-date","20130524T000000Z"); headers.set("x-amz-content-sha256","e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"); headers.set("authorization","AWS4-HMAC-SHA256 Credential=AKIAIOSFODNN7EXAMPLE/20130524/us-east-1/s3/aws4_request,SignedHeaders=host;range;x-amz-content-sha256;x-amz-date,Signature=f0e8bdb87c964420e857bd35b5d6ed310bd44f0170aba48dd91039c6036bdb41"); var clock=java.time.Clock.fixed(java.time.Instant.parse("2013-05-24T00:00:00Z"),java.time.ZoneOffset.UTC); var auth=new SigV4("AKIAIOSFODNN7EXAMPLE","wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY","us-east-1",clock); - var uri=java.net.URI.create("/test.txt");auth.verify("GET",uri,headers); + var uri=java.net.URI.create("/test.txt"); + auth.verify("GET",uri,headers); fails(403,()->auth.verify("GET",java.net.URI.create("/other.txt"),headers)); fails(403,()->auth.verify("DELETE",uri,headers)); fails(403,()->new SigV4("AKIAIOSFODNN7EXAMPLE","wrong","us-east-1",clock).verify("GET",uri,headers)); fails(403,()->new SigV4("AKIAIOSFODNN7EXAMPLE","wrong","us-east-1",java.time.Clock.systemUTC()).verify("GET",uri,headers)); - headers.add("host","duplicate");fails(403,()->auth.verify("GET",uri,headers)); + headers.add("host","duplicate"); + fails(403,()->auth.verify("GET",uri,headers)); System.out.println("SigV4 official vector and tampering tests passed"); } @@ -59,7 +67,8 @@ public final class StoreTest { try(var restarted=new DiskStore(root,8,10)){ try(var obj=restarted.open("test","../nested/☾")){if(obj.stream().read()!=9)throw new AssertionError("Persistence");} if(restarted.list("test","","",100,null).objects().size()!=2)throw new AssertionError("Index persistence"); - restarted.delete("test","../nested/☾");restarted.delete("test","../nested/☾"); + restarted.delete("test","../nested/☾"); + restarted.delete("test","../nested/☾"); fails(404,()->restarted.open("test","../nested/☾")); var uploads=new MultipartStore(restarted); String upload=uploads.create("test","from-parts","text/plain");