125 lines
7.3 KiB
Java
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);
|
|
}
|
|
}
|