Files
2026-10-09 21:37:52 +02:00

125 lines
7.3 KiB
Java

package cloud.lunarsky.store;
import com.sun.net.httpserver.HttpServer;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.nio.file.Files;
import java.nio.file.Path;
import java.time.Duration;
import java.util.List;
import java.util.UUID;
public final class ClusterNodeTest {
public static void main(String[] args) throws Exception {
Path root = Files.createTempDirectory("objectstore-node-");
String token = "local-cluster-test-token-0123456789";
String repairToken = "local-repair-test-token-0123456789";
UUID id = UUID.randomUUID();
UUID hostId = UUID.randomUUID();
UUID nodeId;
byte[] value = "a durable segment".getBytes(java.nio.charset.StandardCharsets.UTF_8);
try (ClusterNode node = new ClusterNode(root, token, repairToken, hostId)) {
try {
new ClusterNode(root, token, repairToken, hostId);
throw new AssertionError("Second writer opened a locked node directory");
} catch (IOException expected) { }
HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
server.createContext("/", node::handle);
server.start();
try {
URI uri = URI.create("http://127.0.0.1:" + server.getAddress().getPort());
NodeIdentity identity = NodeClient.probe(uri, token);
nodeId = identity.nodeId();
require(identity.hostId().equals(hostId), "Node reported the wrong storage host");
NodeClient client = new NodeClient(List.of(new NodeClient.Node(identity.nodeId(), hostId, uri)),
token, repairToken);
client.put(0, id, value, SigV4.hash(value));
client.put(0, id, value, SigV4.hash(value));
require(java.util.Arrays.equals(value, client.get(0, id, value.length, SigV4.hash(value))), "Roundtrip failed");
var request = HttpRequest.newBuilder(uri.resolve("/segments/" + id)).timeout(Duration.ofSeconds(3))
.header("X-Cluster-Token", token).header("X-Cluster-Expected-Node", nodeId.toString())
.header("X-Cluster-Sha256", SigV4.hex(SigV4.hash(value)))
.PUT(HttpRequest.BodyPublishers.ofString("wrong")).build();
var response = HttpClient.newHttpClient().send(request, HttpResponse.BodyHandlers.ofString());
require(response.statusCode() == 400, "Wrong checksum accepted");
var unauthenticated = HttpRequest.newBuilder(uri.resolve("/segments/" + id)).GET().build();
require(HttpClient.newHttpClient().send(unauthenticated, HttpResponse.BodyHandlers.ofString())
.statusCode() == 403, "Unauthenticated read accepted");
var wrongNode = HttpRequest.newBuilder(uri.resolve("/segments/" + id)).timeout(Duration.ofSeconds(3))
.header("X-Cluster-Token", token).header("X-Cluster-Expected-Node", UUID.randomUUID().toString())
.GET().build();
require(HttpClient.newHttpClient().send(wrongNode, HttpResponse.BodyHandlers.ofString())
.statusCode() == 409, "Wrong-node request was accepted");
Path segment = root.resolve("segments").resolve(id.toString().substring(0, 2)).resolve(id.toString());
Files.writeString(segment, "corrupted");
try {
client.get(0, id, value.length, SigV4.hash(value));
throw new AssertionError("Corrupted replica passed verification");
} catch (IOException expected) { }
NodeClient gateway = new NodeClient(List.of(new NodeClient.Node(identity.nodeId(), hostId, uri)),
token, null);
try {
gateway.repair(0, id, value, SigV4.hash(value));
throw new AssertionError("Gateway was allowed to repair a replica");
} catch (IOException expected) { }
var forgedRepair = HttpRequest.newBuilder(uri.resolve("/segments/" + id))
.header("X-Cluster-Token", token).header("X-Cluster-Expected-Node", nodeId.toString())
.header("X-Cluster-Repair", "true").header("X-Cluster-Sha256", SigV4.hex(SigV4.hash(value)))
.PUT(HttpRequest.BodyPublishers.ofByteArray(value)).build();
require(HttpClient.newHttpClient().send(forgedRepair, HttpResponse.BodyHandlers.discarding())
.statusCode() == 403, "Shared node token was allowed to repair a replica");
client.repair(0, id, value, SigV4.hash(value));
require(java.util.Arrays.equals(value, client.get(0, id, value.length, SigV4.hash(value))),
"Repair did not restore the original bytes");
} finally { server.stop(0); }
}
Path pending = root.resolve("pending").resolve("unfinished.part");
Files.writeString(pending, "unfinished");
try (ClusterNode restarted = new ClusterNode(root, token, repairToken, hostId)) {
require(!Files.exists(pending), "Incomplete segment survived restart");
Path segment = root.resolve("segments").resolve(id.toString().substring(0, 2)).resolve(id.toString());
require(java.util.Arrays.equals(value, Files.readAllBytes(segment)), "Committed segment did not survive restart");
}
require(NodeIdentity.open(root, hostId).nodeId().equals(nodeId),
"Node identity changed after restart");
try {
new ClusterNode(root, token, repairToken, UUID.randomUUID());
throw new AssertionError("Node volume accepted a changed host identity");
} catch (IOException expected) { }
HttpServer oversized = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
oversized.createContext("/identity", exchange -> {
byte[] bytes = new byte[1024];
exchange.sendResponseHeaders(200, 0);
exchange.getResponseBody().write(bytes);
exchange.close();
});
oversized.createContext("/segments/", exchange -> {
byte[] bytes = new byte[1024];
exchange.sendResponseHeaders(200, 0);
exchange.getResponseBody().write(bytes);
exchange.close();
});
oversized.start();
try {
URI uri = URI.create("http://127.0.0.1:" + oversized.getAddress().getPort());
try {
NodeClient.probe(uri, token);
throw new AssertionError("Oversized identity response was accepted");
} catch (IOException expected) { }
NodeClient client = new NodeClient(List.of(new NodeClient.Node(nodeId, hostId, uri)), token, null);
try {
client.get(0, id, value.length, SigV4.hash(value));
throw new AssertionError("Oversized segment response was accepted");
} catch (IOException expected) { }
} finally { oversized.stop(0); }
System.out.println("Cluster node tests passed: lock, authenticated roundtrip, checksums, restart cleanup");
}
private static void require(boolean condition, String message) {
if (!condition) throw new AssertionError(message);
}
}