Add ObjectStore source

This commit is contained in:
LunarSkyOSS committed 2026-10-09 15:23:16 +02:00
1 parent 8a3bee5dcc
commit 51e206f8df
14 files changed
+673

No files matched your search

+137
View File
@@ -0,0 +1,137 @@
package cloud.lunarsky.store;
import java.io.*;
import java.nio.ByteBuffer;
import java.nio.channels.FileChannel;
import java.nio.channels.FileLock;
import java.nio.channels.OverlappingFileLockException;
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;
final class DiskStore implements AutoCloseable {
private static final long MAGIC = 0x4c534f424a303031L;
private static final int HEADER = 72;
private final Path objects, temporary;
private final FileChannel lockChannel;
private final FileLock processLock;
private final long maxObject, maxTotal;
private final Object[] locks = new Object[128];
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(); }
}
DiskStore(Path root, long maxObject, long maxTotal) throws IOException {
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);
FileLock acquired = null;
boolean ready = false;
try {
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);
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);
}
}
}
ready = true;
} finally {
if (!ready) {
if (acquired != null) acquired.release();
channel.close();
}
}
lockChannel = channel;
processLock = acquired;
}
@Override public void close() throws IOException {
processLock.release();
lockChannel.close();
}
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 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");
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;
}
}
return new Metadata(count,modified,SigV4.hex(etag),hash);
} finally {Files.deleteIfExists(pending);}
}
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;}
}
}
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;}
}
}
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);
}
private static MessageDigest digest(String algorithm){try{return MessageDigest.getInstance(algorithm);}catch(java.security.NoSuchAlgorithmException e){throw new IllegalStateException(e);}}
}
+114
View File
@@ -0,0 +1,114 @@
package cloud.lunarsky.store;
import com.sun.net.httpserver.HttpExchange;
import com.sun.net.httpserver.HttpServer;
import java.io.IOException;
import java.net.InetSocketAddress;
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.Map;
import java.util.UUID;
import java.util.concurrent.Executors;
import java.util.concurrent.Semaphore;
public final class Main {
private final DiskStore 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;}
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");
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;
}
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");
}
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){
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);
}
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 "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();}
}
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);
}
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(()->{
server.stop(5);
executor.close();
try { store.close(); }
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);
}
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;}
}
+129
View File
@@ -0,0 +1,129 @@
package cloud.lunarsky.store;
import com.sun.net.httpserver.Headers;
import java.net.URI;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.time.Clock;
import java.time.Duration;
import java.time.Instant;
import java.time.ZoneOffset;
import java.time.format.DateTimeFormatter;
import java.util.Arrays;
import java.util.HexFormat;
import java.util.Map;
import java.util.TreeMap;
import java.util.regex.Pattern;
import javax.crypto.Mac;
import javax.crypto.spec.SecretKeySpec;
final class SigV4 {
private static final DateTimeFormatter DATE = DateTimeFormatter.ofPattern("uuuuMMdd'T'HHmmss'Z'").withZone(ZoneOffset.UTC);
private static final Pattern HEX = Pattern.compile("[0-9a-f]{64}");
private final String accessKey, secretKey, region;
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;
}
String verify(String method, URI uri, Headers headers) {
String authorization = single(headers, "authorization");
if (authorization == null || !authorization.startsWith("AWS4-HMAC-SHA256 ")) denied("Signed requests are required");
Map<String,String> fields = new TreeMap<>();
for (String part : authorization.substring(17).split(",")) {
String[] pair = part.trim().split("=", 2);
if (pair.length != 2 || fields.put(pair[0], pair[1]) != null) denied("Invalid authorization header");
}
if (!fields.keySet().equals(java.util.Set.of("Credential", "SignedHeaders", "Signature"))) denied("Invalid authorization fields");
String[] credential = fields.get("Credential").split("/", -1);
if (credential.length != 5 || !credential[0].equals(accessKey) || !credential[2].equals(region)
|| !credential[3].equals("s3") || !credential[4].equals("aws4_request")) denied("Invalid credential scope");
String date = single(headers, "x-amz-date"), payload = single(headers, "x-amz-content-sha256");
if (date == null || !credential[1].matches("[0-9]{8}") || !date.matches("[0-9]{8}T[0-9]{6}Z") || !date.startsWith(credential[1])) denied("Invalid signing date");
try {
Instant signed = Instant.from(DATE.parse(date));
if (Duration.between(signed, clock.instant()).abs().compareTo(Duration.ofMinutes(5)) > 0)
throw new StoreException(403, "RequestTimeTooSkewed", "Request timestamp is outside the permitted window");
} catch (java.time.DateTimeException e) { denied("Invalid signing date"); }
if (payload == null || !HEX.matcher(payload).matches())
throw new StoreException(400, "NotImplemented", "A hexadecimal SHA-256 payload hash is required; unsigned and chunk-signed payloads are unsupported");
if (headers.containsKey("x-amz-security-token")) denied("Temporary credentials are unsupported");
String signedHeaders = fields.get("SignedHeaders");
String[] names = signedHeaders.split(";", -1);
if (names.length > 32 || !signedHeaders.equals(String.join(";", Arrays.stream(names).distinct().sorted().toList()))) denied("Signed headers must be unique and sorted");
var namesSet = java.util.Set.copyOf(Arrays.asList(names));
if (!namesSet.containsAll(java.util.Set.of("host", "x-amz-date", "x-amz-content-sha256"))) denied("Missing signed headers");
for (String key : headers.keySet()) {
String lower = key.toLowerCase(java.util.Locale.ROOT);
if (lower.startsWith("x-amz-") && !namesSet.contains(lower)) denied("Unsigned Amazon header");
}
if (headers.containsKey("if-none-match") && !namesSet.contains("if-none-match")) denied("Unsigned write condition");
StringBuilder canonicalHeaders = new StringBuilder();
for (String name : names) {
if (!name.matches("[a-z0-9-]+")) denied("Invalid signed header name");
String value = single(headers, name);
if (value == null) denied("Missing signed header");
canonicalHeaders.append(name).append(':').append(value.trim().replaceAll("[\\t ]+", " ")).append('\n');
}
String canonical = method + "\n" + encode(decode(uri.getRawPath()), true) + "\n"
+ canonicalQuery(uri.getRawQuery()) + "\n" + canonicalHeaders + "\n" + signedHeaders + "\n" + payload;
String scope = String.join("/", Arrays.copyOfRange(credential, 1, 5));
String toSign = "AWS4-HMAC-SHA256\n" + date + "\n" + scope + "\n" + hex(hash(canonical.getBytes(StandardCharsets.UTF_8)));
byte[] signingKey = signingKey(secretKey, credential[1], region);
String signature = fields.get("Signature");
if (!HEX.matcher(signature).matches() || !MessageDigest.isEqual(hmac(signingKey, toSign), HexFormat.of().parseHex(signature))) denied("Signature mismatch");
return payload;
}
static String single(Headers headers, String name) {
var values = headers.get(name);
if (values == null) return null;
if (values.size() != 1) denied("Duplicate security-relevant header");
return values.getFirst();
}
static String decode(String value) {
try {
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),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");}
}
static String encode(String value, boolean keepSlash) {
StringBuilder result = new StringBuilder();
for (byte b : value.getBytes(StandardCharsets.UTF_8)) {
int c = b & 255;
if ((c >= 'A' && c <= 'Z') || (c >= 'a' && c <= 'z') || (c >= '0' && c <= '9') || c == '-' || c == '_' || c == '.' || c == '~' || (c == '/' && keepSlash)) result.append((char)c);
else result.append('%').append("0123456789ABCDEF".charAt(c >> 4)).append("0123456789ABCDEF".charAt(c & 15));
}
return result.toString();
}
static String canonicalQuery(String raw) {
if (raw == null || raw.isEmpty()) return "";
return Arrays.stream(raw.split("&", -1)).map(part -> {
String[] pair = part.split("=", 2);
return encode(decode(pair[0]), false) + "=" + encode(decode(pair.length == 2 ? pair[1] : ""), false);
}).sorted().collect(java.util.stream.Collectors.joining("&"));
}
static byte[] signingKey(String secret, String date, String region) {
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)); }
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 String hex(byte[] data) { return HexFormat.of().formatHex(data); }
private static void denied(String message) { throw new StoreException(403,"AccessDenied",message); }
}
@@ -0,0 +1,11 @@
package cloud.lunarsky.store;
final class StoreException extends RuntimeException {
final int status;
final String code;
StoreException(int status, String code, String message) {
super(message);
this.status = status;
this.code = code;
}
}