From 510e519bd186059da172fb3d618e6445e267964a Mon Sep 17 00:00:00 2001 From: LunarSkyOSS Date: Sun, 11 Oct 2026 01:24:09 +0200 Subject: [PATCH] Configure gateway limits and encrypted storage --- .env.cluster.example | 8 ++ .env.example | 6 + Dockerfile | 1 + README.md | 26 +++- TESTS.md | 5 +- compose.cluster.encrypted.yaml | 114 ++++++++++++++++++ compose.cluster.yaml | 7 ++ compose.encrypted.yaml | 11 ++ compose.yaml | 3 + scripts/prepare-encrypted-storage.sh | 69 +++++++++++ scripts/test.sh | 1 + src/cloud/lunarsky/store/ClusterGc.java | 4 +- src/cloud/lunarsky/store/ClusterJoin.java | 2 + src/cloud/lunarsky/store/ClusterMigrate.java | 2 + src/cloud/lunarsky/store/ClusterNode.java | 1 + src/cloud/lunarsky/store/ClusterRepair.java | 4 +- src/cloud/lunarsky/store/ClusterStore.java | 2 + src/cloud/lunarsky/store/DiskStore.java | 1 + src/cloud/lunarsky/store/EncryptedVolume.java | 29 +++++ src/cloud/lunarsky/store/Main.java | 70 +++++++++-- src/cloud/lunarsky/store/StorageLimits.java | 19 +++ .../lunarsky/store/EncryptedVolumeTest.java | 42 +++++++ test/cloud/lunarsky/store/HttpTest.java | 70 ++++++++++- 23 files changed, 480 insertions(+), 17 deletions(-) create mode 100644 compose.cluster.encrypted.yaml create mode 100644 compose.encrypted.yaml create mode 100644 scripts/prepare-encrypted-storage.sh create mode 100644 src/cloud/lunarsky/store/EncryptedVolume.java create mode 100644 src/cloud/lunarsky/store/StorageLimits.java create mode 100644 test/cloud/lunarsky/store/EncryptedVolumeTest.java diff --git a/.env.cluster.example b/.env.cluster.example index cba32b5..a01387c 100644 --- a/.env.cluster.example +++ b/.env.cluster.example @@ -5,3 +5,11 @@ CLUSTER_REPAIR_TOKEN= POSTGRES_PASSWORD=[PG_Pass] S3_BUCKET=objects CLUSTER_HOST_PORT=9001 +MAX_OBJECT_BYTES=134217728 +MAX_TOTAL_BYTES=2147483648 +MAX_IN_FLIGHT_REQUESTS=16 +HTTP_BACKLOG=64 +S3_VIRTUAL_HOST_SUFFIX= +# For compose.cluster.encrypted.yaml only; provision and mount LUKS before setting these. +ENCRYPTED_STORAGE_ROOT= +ENCRYPTED_VOLUME_ID= diff --git a/.env.example b/.env.example index ae964f9..895f400 100644 --- a/.env.example +++ b/.env.example @@ -13,3 +13,9 @@ PUBLIC_BYTES_PER_SECOND=0 PUBLIC_BYTE_BURST= PUBLIC_MAX_IN_FLIGHT_PER_IP= PUBLIC_TRUSTED_PROXY_IPS= +MAX_IN_FLIGHT_REQUESTS=16 +HTTP_BACKLOG=64 +S3_VIRTUAL_HOST_SUFFIX= +# For compose.encrypted.yaml only; provision and mount LUKS before setting these. +ENCRYPTED_STORAGE_ROOT= +ENCRYPTED_VOLUME_ID= diff --git a/Dockerfile b/Dockerfile index a7f7a65..e4e2974 100644 --- a/Dockerfile +++ b/Dockerfile @@ -11,6 +11,7 @@ RUN java --add-modules jdk.httpserver -cp /out:/tmp/hash4j.jar cloud.lunarsky.st RUN java --add-modules jdk.httpserver -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.ConcurrencyTest RUN java --add-modules jdk.httpserver,java.net.http -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.HttpTest RUN java --add-modules jdk.httpserver,java.net.http -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.ClientLimitsTest +RUN java -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.EncryptedVolumeTest RUN java --add-modules jdk.httpserver,java.net.http -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.ClusterNodeTest RUN java --add-modules jdk.httpserver,java.net.http -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.ClusterTlsTest RUN java -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.CliTest diff --git a/README.md b/README.md index a582352..4115b9b 100644 --- a/README.md +++ b/README.md @@ -16,6 +16,7 @@ Source: [GitHub](https://github.com/LunarSkyOSS/ObjectStore) · [Gitea mirror](h - [S3 API support checklist](#s3-api-support-checklist) - [Tests](TESTS.md) - [Single-node setup](#single-node-setup) +- [Encrypted storage](#encrypted-storage) - [Access keys and ACLs](#access-keys-and-acls) - [CLI and tests](#cli-and-tests) - [Capability discovery](#capability-discovery) @@ -35,6 +36,7 @@ Source: [GitHub](https://github.com/LunarSkyOSS/ObjectStore) · [Gitea mirror](h ObjectStore creates the configured default bucket at startup. Additional buckets share the configured capacity limit. - ✅ Persistent single-node storage with checksum verification +- ✅ Optional dm-crypt-backed data mounts with startup guards - ✅ Configurable per-object and total logical size limits - ✅ CLI status, version, and full payload verification - ✅ Local cluster prototype with stable node IDs and host-aware placement code @@ -49,6 +51,7 @@ New objects retain their content type and key. Objects written by the earlier si ## S3 API support checklist - ✅ Header-based and presigned-query AWS Signature Version 4 authentication +- ✅ Path-style and configurable virtual-hosted-style bucket addresses - ✅ `PutObject`, `GetObject`, `HeadObject`, and `DeleteObject` in both modes - ✅ Single-range GET and `ListObjectsV2` in both modes - ✅ SHA-256 payload verification and `x-amz-checksum-sha256` in both modes @@ -84,6 +87,23 @@ Port 9000 binds to localhost. Data stays in the `object-data` Docker volume. `do The standalone defaults are 128 MiB per object and 2 GiB total. Set `MAX_OBJECT_BYTES` and `MAX_TOTAL_BYTES` in `.env` to change them. Incomplete multipart uploads consume space until aborted. +## Encrypted storage + +At-rest encryption is an **opt-in deployment setting** backed by a host-mounted dm-crypt/LUKS filesystem. ObjectStore does not encrypt bytes itself or store an encryption key in its environment. Prepare and unlock the LUKS filesystem outside Docker, then set `ENCRYPTED_STORAGE_ROOT` to its exact mount point and `ENCRYPTED_VOLUME_ID` to a stable 32-character lowercase hex identifier in your private environment file. Generate the identifier with `openssl rand -hex 16`; it is a mount guard, **not** a cryptographic key. Keep the LUKS recovery key separately from the data and backups. + +With the encrypted filesystem mounted, prepare empty directories using `sh scripts/prepare-encrypted-storage.sh single /path/to/objectstore.env` or `sh scripts/prepare-encrypted-storage.sh cluster /path/to/cluster.env`. The script reads only the two encryption settings from that file; it does not execute it. It requires an active dm-crypt-backed mount at `ENCRYPTED_STORAGE_ROOT` and refuses to mark a nonempty directory. Run it with permission to create the directories and set container ownership. Enable the matching Compose override: + +```sh +docker compose --env-file /path/to/objectstore.env -f compose.yaml -f compose.encrypted.yaml up -d --build +docker compose --env-file /path/to/cluster.env -f compose.cluster.yaml -f compose.cluster.encrypted.yaml up -d --build +``` + +Use the first command for single-node mode or the second for the **local development cluster**, not both. The overrides bind encrypted directories for object data, cluster nodes, PostgreSQL, and cluster staging files. Each process checks a marker stored on its mounted directory before opening data; PostgreSQL checks before starting. If a mount is missing after reboot, the service refuses to start instead of silently creating a new plaintext store. Provision an unlock-and-mount service before starting Compose; a container restart cannot unlock LUKS. The preparation script checks the host mount, while the runtime marker is only a guard against an absent or wrong mount. + +**Existing Docker volumes are not migrated automatically.** Stop the stack, take and verify a backup, prepare the empty encrypted directories, then copy the stopped volumes into their corresponding directories before starting with the override. For PostgreSQL, copy the existing data directory into `cluster/metadata/pgdata` because the encrypted override changes `PGDATA`. Verify the restored objects and database before retiring the original volumes. The cluster overlay places all local test containers on one encrypted host; a future multi-host deployment needs an encrypted data mount and key recovery plan on each host. + +This covers the specified data and staging mounts, but not Docker logs, host swap, other temporary files, external backups, or a compromised running host. Encrypt those separately as appropriate, use TLS for network paths, restrict access, and test backup recovery. Encryption at rest is one security measure; enabling it does not by itself establish GDPR compliance. + ## Access keys and ACLs `S3_ACCESS_KEY` is the owner identity. Its secret is `S3_SECRET_KEY`. Additional keys are optional: place one `ACCESSKEY:secret` pair per line in a file mounted read-only inside the container, and set `S3_CREDENTIALS_FILE=/run/secrets/s3-credentials` in the environment file. Use a Compose override to mount an absolute host path at `/run/secrets/s3-credentials` for the `objectstore` service (or `gateway` in cluster mode). Each access key needs 16–128 alphanumeric characters and each secret at least 32 characters. A restart loads changes to that file. Keep it outside Git and protect it as a secret. Additional identities have no access until the owner grants it. @@ -191,7 +211,11 @@ The maintenance service is opt-in. It repairs missing or corrupt replicas and re ## Limits and safety -Public-client limits are **disabled by default**. The localhost Compose setup is unchanged. For an endpoint that accepts external clients, set one or both of `PUBLIC_REQUESTS_PER_SECOND` and `PUBLIC_BYTES_PER_SECOND` to a positive number in `.env`. The first limits requests per client IP with a token bucket; the second paces upload and download bytes through one shared per-IP budget. `PUBLIC_REQUEST_BURST` and `PUBLIC_BYTE_BURST` default to one second of their respective rates. `PUBLIC_MAX_IN_FLIGHT_PER_IP` defaults to 8 when either rate is enabled. Requests above the rate or concurrency limit receive S3 `503 SlowDown` and `Retry-After: 1`; an admitted transfer is paced rather than cut off. These are per-gateway limits, not cluster-wide quotas. The existing 16-request gateway cap and storage limits still apply. +Public-client limits are **disabled by default**. The localhost Compose setup is unchanged. For an endpoint that accepts external clients, set one or both of `PUBLIC_REQUESTS_PER_SECOND` and `PUBLIC_BYTES_PER_SECOND` to a positive number in `.env`. The first limits requests per client IP with a token bucket; the second paces upload and download bytes through one shared per-IP budget. `PUBLIC_REQUEST_BURST` and `PUBLIC_BYTE_BURST` default to one second of their respective rates. `PUBLIC_MAX_IN_FLIGHT_PER_IP` defaults to 8 when either rate is enabled. Requests above the rate or concurrency limit receive S3 `503 SlowDown` and `Retry-After: 1`; an admitted transfer is paced rather than cut off. These are per-gateway limits, not cluster-wide quotas. + +`MAX_IN_FLIGHT_REQUESTS` controls the per-gateway concurrency cap (default 16, range 1–1024); an additional request receives `503 SlowDown`. `HTTP_BACKLOG` controls the listening socket backlog (default 64, range 1–4096). Raise either only after a mixed upload/download load test, because each active request can hold memory, disk space, and database connections. The configured object and total storage limits still apply. `MAX_OBJECT_BYTES` currently cannot exceed 1 GiB. + +Set `S3_VIRTUAL_HOST_SUFFIX` to a DNS suffix such as `s3.example.com` to accept `bucket.s3.example.com/key` alongside path-style `/bucket/key`. Configure DNS and TLS for the bucket hostnames, and preserve the client's original `Host` header through the proxy; Signature V4 signs it. The suffix is empty by default. This changes request parsing, not bucket naming rules or DNS configuration. For example, to start with 100 requests per second and 16 MiB/s combined upload and download per IP, set `PUBLIC_REQUESTS_PER_SECOND=100` and `PUBLIC_BYTES_PER_SECOND=16777216`. Both start with a one-second burst. Adjust these numbers after measuring the actual workload; do not copy them as a universal production policy. diff --git a/TESTS.md b/TESTS.md index 44a542a..65d891f 100644 --- a/TESTS.md +++ b/TESTS.md @@ -18,8 +18,9 @@ The script compiles the source and test programs into `out/classes`, then runs: | --- | --- | | `StoreTest` | Signature V4 and signed-stream vectors, CRC64NVME and XXHash reference vectors, local writes and reads, quotas, restart persistence for objects, checksums, ACLs, attributes, buckets, and versions, multipart recovery, legacy reads, locking, and corruption rejection. | | `ConcurrencyTest` | Atomic local overwrites and consistent reads, listings, and deletes during concurrent access. | -| `HttpTest` | Signed capability discovery, presigned URLs, streaming uploads and trailers, object and bucket operations, ACL grants with a second access key and public reads, copies, checksum persistence and rejection, ranges, listing, metadata, tags, multipart uploads, and versioning in single-node mode. | +| `HttpTest` | Signed capability discovery, presigned URLs, streaming uploads and trailers, object and bucket operations, ACL grants with a second access key and public reads, copies, checksum persistence and rejection, ranges, listing, metadata, tags, multipart uploads, versioning, and signed virtual-hosted requests in single-node mode. | | `ClientLimitsTest` | Disabled defaults, trusted-proxy address validation, ignored untrusted headers, per-IP request refusal, and paced response bytes. | +| `EncryptedVolumeTest` | Marker acceptance and refusal when the configured encrypted directory is missing, mismatched, or a symlink. | | `ClusterNodeTest` | Node identity and locking, authenticated segment transfers, checksum rejection, repair authorization, inventory and guarded deletion, and restart cleanup. | | `ClusterTlsTest` | HTTPS node identity and segment roundtrip with a trusted certificate, plus rejection of untrusted and wrong-host certificates. | | `CliTest` | Version, status, verification, and a nonzero result for corrupt data. | @@ -27,6 +28,8 @@ The script compiles the source and test programs into `out/classes`, then runs: The script exits nonzero on failure. The test programs use temporary local directories and loopback HTTP or HTTPS ports; they do not use an existing ObjectStore volume. `ClusterTlsTest` uses the JDK's `keytool` to create disposable test certificates. +The encrypted Compose overlays can be checked with `docker compose config` without starting the stack. This verifies their bind mounts and startup settings, not that a host filesystem is actually encrypted. `scripts/prepare-encrypted-storage.sh` must also verify an active dm-crypt-backed mount. A complete deployment drill must unlock it, prepare directories, start with the overlay, write data, restart, restore from backup, and confirm that the stack refuses to start while the mount is unavailable. Do not point that drill at existing data. + ## Disposable metadata routing test Run `sh scripts/test-metadata.sh` with Docker Compose. It creates a separate, temporary PostgreSQL container and two in-process storage nodes. The test connects through a two-host JDBC URL whose first host is unavailable, checks writable readiness, then makes the test database read-only and verifies that readiness drops. The script removes its test containers and temporary database afterward. It does not promote a standby or test automatic failover. diff --git a/compose.cluster.encrypted.yaml b/compose.cluster.encrypted.yaml new file mode 100644 index 0000000..7846bc5 --- /dev/null +++ b/compose.cluster.encrypted.yaml @@ -0,0 +1,114 @@ +x-encrypted-id: &encrypted-id ${ENCRYPTED_VOLUME_ID:?Set ENCRYPTED_VOLUME_ID} +x-encrypted-data: &encrypted-data + ENCRYPTED_VOLUME_ID: *encrypted-id + ENCRYPTED_VOLUME_MARKER_FILE: /data/.objectstore-encrypted +x-encrypted-staging: &encrypted-staging + ENCRYPTED_VOLUME_ID: *encrypted-id + ENCRYPTED_VOLUME_MARKER_FILE: /tmp/.objectstore-encrypted +x-metadata-entrypoint: &metadata-entrypoint + - /bin/sh + - -ec + - | + test -f /var/lib/postgresql/data/.objectstore-encrypted && + test ! -L /var/lib/postgresql/data/.objectstore-encrypted && + printf '%s\n' "$(printenv ENCRYPTED_VOLUME_ID)" | grep -Eq '^[a-f0-9]{32}$' && + grep -Fxq -e "$(printenv ENCRYPTED_VOLUME_ID)" /var/lib/postgresql/data/.objectstore-encrypted || + { echo 'Encrypted metadata volume unavailable' >&2; exit 1; } + exec docker-entrypoint.sh postgres + +services: + gateway: + environment: *encrypted-staging + volumes: + - type: bind + source: ${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT}/cluster/gateway-staging + target: /tmp + bind: + create_host_path: false + + metadata: + environment: + ENCRYPTED_VOLUME_ID: *encrypted-id + PGDATA: /var/lib/postgresql/data/pgdata + entrypoint: *metadata-entrypoint + volumes: + - type: bind + source: ${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT}/cluster/metadata + target: /var/lib/postgresql/data + bind: + create_host_path: false + + metadata-recovery: + environment: + ENCRYPTED_VOLUME_ID: *encrypted-id + PGDATA: /var/lib/postgresql/data/pgdata + entrypoint: *metadata-entrypoint + volumes: + - type: bind + source: ${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT}/cluster/metadata-recovery + target: /var/lib/postgresql/data + bind: + create_host_path: false + + repair: + environment: *encrypted-staging + volumes: + - type: bind + source: ${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT}/cluster/repair-staging + target: /tmp + bind: + create_host_path: false + + gc: + environment: *encrypted-staging + volumes: + - type: bind + source: ${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT}/cluster/gc-staging + target: /tmp + bind: + create_host_path: false + + maintenance: + environment: *encrypted-staging + volumes: + - type: bind + source: ${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT}/cluster/maintenance-staging + target: /tmp + bind: + create_host_path: false + + node-a: + environment: *encrypted-data + volumes: + - type: bind + source: ${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT}/cluster/node-a + target: /data + bind: + create_host_path: false + + node-b: + environment: *encrypted-data + volumes: + - type: bind + source: ${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT}/cluster/node-b + target: /data + bind: + create_host_path: false + + node-c: + environment: *encrypted-data + volumes: + - type: bind + source: ${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT}/cluster/node-c + target: /data + bind: + create_host_path: false + + node-d: + environment: *encrypted-data + volumes: + - type: bind + source: ${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT}/cluster/node-d + target: /data + bind: + create_host_path: false diff --git a/compose.cluster.yaml b/compose.cluster.yaml index 20555d5..d3c4e71 100644 --- a/compose.cluster.yaml +++ b/compose.cluster.yaml @@ -15,6 +15,9 @@ services: S3_CREDENTIALS_FILE: ${S3_CREDENTIALS_FILE:-} MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728} MAX_TOTAL_BYTES: ${MAX_TOTAL_BYTES:-2147483648} + MAX_IN_FLIGHT_REQUESTS: ${MAX_IN_FLIGHT_REQUESTS:-16} + HTTP_BACKLOG: ${HTTP_BACKLOG:-64} + S3_VIRTUAL_HOST_SUFFIX: ${S3_VIRTUAL_HOST_SUFFIX:-} PUBLIC_REQUESTS_PER_SECOND: ${PUBLIC_REQUESTS_PER_SECOND:-0} PUBLIC_REQUEST_BURST: ${PUBLIC_REQUEST_BURST:-} PUBLIC_BYTES_PER_SECOND: ${PUBLIC_BYTES_PER_SECOND:-0} @@ -99,6 +102,8 @@ services: POSTGRES_JDBC_URL: ${POSTGRES_JDBC_URL:-jdbc:postgresql://metadata:5432/objectstore?connectTimeout=3&socketTimeout=10&targetServerType=primary} POSTGRES_USER: objectstore POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD} + MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728} + MAX_TOTAL_BYTES: ${MAX_TOTAL_BYTES:-2147483648} CLUSTER_GC_MIN_AGE_SECONDS: ${CLUSTER_GC_MIN_AGE_SECONDS:-1209600} CLUSTER_BACKUP_RETENTION_SECONDS: ${CLUSTER_BACKUP_RETENTION_SECONDS:-0} CLUSTER_GC_TEST_MODE: ${CLUSTER_GC_TEST_MODE:-false} @@ -134,6 +139,8 @@ services: POSTGRES_JDBC_URL: ${POSTGRES_JDBC_URL:-jdbc:postgresql://metadata:5432/objectstore?connectTimeout=3&socketTimeout=10&targetServerType=primary} POSTGRES_USER: objectstore POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD} + MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728} + MAX_TOTAL_BYTES: ${MAX_TOTAL_BYTES:-2147483648} depends_on: metadata: condition: service_healthy diff --git a/compose.encrypted.yaml b/compose.encrypted.yaml new file mode 100644 index 0000000..cd30fe6 --- /dev/null +++ b/compose.encrypted.yaml @@ -0,0 +1,11 @@ +services: + objectstore: + environment: + ENCRYPTED_VOLUME_ID: ${ENCRYPTED_VOLUME_ID:?Set ENCRYPTED_VOLUME_ID} + ENCRYPTED_VOLUME_MARKER_FILE: /data/.objectstore-encrypted + volumes: + - type: bind + source: ${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT}/single + target: /data + bind: + create_host_path: false diff --git a/compose.yaml b/compose.yaml index 9c6004f..1563686 100644 --- a/compose.yaml +++ b/compose.yaml @@ -11,6 +11,9 @@ services: S3_CREDENTIALS_FILE: ${S3_CREDENTIALS_FILE:-} MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728} MAX_TOTAL_BYTES: ${MAX_TOTAL_BYTES:-2147483648} + MAX_IN_FLIGHT_REQUESTS: ${MAX_IN_FLIGHT_REQUESTS:-16} + HTTP_BACKLOG: ${HTTP_BACKLOG:-64} + S3_VIRTUAL_HOST_SUFFIX: ${S3_VIRTUAL_HOST_SUFFIX:-} PUBLIC_REQUESTS_PER_SECOND: ${PUBLIC_REQUESTS_PER_SECOND:-0} PUBLIC_REQUEST_BURST: ${PUBLIC_REQUEST_BURST:-} PUBLIC_BYTES_PER_SECOND: ${PUBLIC_BYTES_PER_SECOND:-0} diff --git a/scripts/prepare-encrypted-storage.sh b/scripts/prepare-encrypted-storage.sh new file mode 100644 index 0000000..8ec7b6f --- /dev/null +++ b/scripts/prepare-encrypted-storage.sh @@ -0,0 +1,69 @@ +#!/bin/sh +set -eu + +mode=${1:-} +case "$mode" in + single) directories='single' ;; + cluster) directories='cluster/gateway-staging cluster/metadata cluster/metadata-recovery cluster/repair-staging cluster/gc-staging cluster/maintenance-staging cluster/node-a cluster/node-b cluster/node-c cluster/node-d' ;; + *) echo 'Usage: prepare-encrypted-storage.sh single|cluster [environment-file]' >&2; exit 2 ;; +esac + +if [ "$#" -gt 2 ]; then + echo 'Too many arguments' >&2 + exit 2 +fi +if [ "$#" -eq 2 ]; then + [ -r "$2" ] || { echo 'Environment file is not readable' >&2; exit 2; } + while IFS='=' read -r name value || [ -n "${name:-}" ]; do + case "$name" in + ENCRYPTED_STORAGE_ROOT) ENCRYPTED_STORAGE_ROOT=$value ;; + ENCRYPTED_VOLUME_ID) ENCRYPTED_VOLUME_ID=$value ;; + esac + done < "$2" +fi + +root=${ENCRYPTED_STORAGE_ROOT:?Set ENCRYPTED_STORAGE_ROOT to an existing mounted LUKS directory} +id=${ENCRYPTED_VOLUME_ID:?Set ENCRYPTED_VOLUME_ID to 32 lowercase hex characters} +case "$id" in *[!0-9a-f]*|'') echo 'ENCRYPTED_VOLUME_ID must be lowercase hex' >&2; exit 2 ;; esac +[ "${#id}" -eq 32 ] || { echo 'ENCRYPTED_VOLUME_ID must be 32 characters' >&2; exit 2; } +[ -d "$root" ] && [ ! -L "$root" ] || { echo 'Encrypted root is missing or a symlink' >&2; exit 1; } +root=$(realpath -e "$root") +source=$(findmnt -n -M "$root" -o SOURCE) || { echo 'Encrypted root is not a mount point' >&2; exit 1; } +device=$(printf '%s\n' "$source" | sed 's/\[.*$//') +case "$device" in /dev/*) ;; *) echo 'Encrypted root must be backed by a block device' >&2; exit 1 ;; esac +lsblk -s -n -r -o TYPE "$device" | grep -Fxq crypt || { + echo 'Encrypted root is not backed by an active dm-crypt mapping' >&2 + exit 1 +} + +for relative in $directories; do + path=$root/$relative + [ ! -L "$root/cluster" ] || { echo 'Refusing symlink under encrypted root' >&2; exit 1; } + [ ! -L "$path" ] || { echo "Refusing symlink: $path" >&2; exit 1; } + if [ -d "$path" ] && [ ! -e "$path/.objectstore-encrypted" ] && + [ -n "$(find "$path" -mindepth 1 -maxdepth 1 -print -quit)" ]; then + echo "Refusing to mark nonempty directory: $path" >&2 + exit 1 + fi + install -d -m 0700 "$path" + [ "$(findmnt -n -T "$path" -o TARGET)" = "$root" ] || { + echo "Directory is not on the encrypted mount: $path" >&2 + exit 1 + } + marker=$path/.objectstore-encrypted + if [ -e "$marker" ]; then + [ ! -L "$marker" ] && [ "$(cat "$marker")" = "$id" ] || { + echo "Encrypted volume marker mismatch: $path" >&2 + exit 1 + } + else + printf '%s\n' "$id" > "$marker" + fi + chmod 0600 "$marker" + case "$relative" in + cluster/metadata|cluster/metadata-recovery) chown 70:70 "$path" "$marker" ;; + *) chown 10001:10001 "$path" "$marker" ;; + esac +done + +echo "Prepared $mode paths on the active encrypted mount: $root" diff --git a/scripts/test.sh b/scripts/test.sh index 76b42a7..a123e37 100644 --- a/scripts/test.sh +++ b/scripts/test.sh @@ -8,6 +8,7 @@ java --add-modules jdk.httpserver -cp out/classes:lib/hash4j-0.30.0.jar cloud.lu java --add-modules jdk.httpserver -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.ConcurrencyTest java --add-modules jdk.httpserver,java.net.http -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.HttpTest java --add-modules jdk.httpserver,java.net.http -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.ClientLimitsTest +java -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.EncryptedVolumeTest java --add-modules jdk.httpserver,java.net.http -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.ClusterNodeTest java --add-modules jdk.httpserver,java.net.http -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.ClusterTlsTest java -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.CliTest diff --git a/src/cloud/lunarsky/store/ClusterGc.java b/src/cloud/lunarsky/store/ClusterGc.java index b1aa4a4..6d1f1c1 100644 --- a/src/cloud/lunarsky/store/ClusterGc.java +++ b/src/cloud/lunarsky/store/ClusterGc.java @@ -15,6 +15,7 @@ public final class ClusterGc { if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV"))) throw new IllegalArgumentException("Cluster garbage collection is only enabled in local cluster mode"); boolean testDomains = "true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS")); + ObjectStorage.Limits limits = StorageLimits.fromEnvironment(env); boolean disposableTest = testDomains && "true".equals(env.get("CLUSTER_GC_TEST_MODE")); long age = Long.parseLong(env.getOrDefault("CLUSTER_GC_MIN_AGE_SECONDS", "1209600")); long backupRetention = Long.parseLong(env.getOrDefault("CLUSTER_BACKUP_RETENTION_SECONDS", "0")); @@ -25,7 +26,8 @@ public final class ClusterGc { try (ClusterStore store = new ClusterStore(env.get("POSTGRES_JDBC_URL"), env.get("POSTGRES_USER"), env.get("POSTGRES_PASSWORD"), env.get("S3_BUCKET"), Arrays.stream(env.get("CLUSTER_NODES").split(",")).map(URI::create).toList(), - env.get("CLUSTER_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"), 134217728, 2147483648L, + env.get("CLUSTER_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"), + limits.maxObjectBytes(), limits.maxTotalBytes(), testDomains)) { if (apply) { var repair = store.repairOnce(); diff --git a/src/cloud/lunarsky/store/ClusterJoin.java b/src/cloud/lunarsky/store/ClusterJoin.java index 7fc8e03..fad23d8 100644 --- a/src/cloud/lunarsky/store/ClusterJoin.java +++ b/src/cloud/lunarsky/store/ClusterJoin.java @@ -12,6 +12,8 @@ public final class ClusterJoin { if (args.length != 2) throw new IllegalArgumentException("Usage: objectstore cluster-join node-url expected-host-uuid"); Map env = System.getenv(); + EncryptedVolume.requireConfigured(env, + java.nio.file.Path.of(System.getProperty("java.io.tmpdir"))); if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV"))) throw new IllegalArgumentException("Node registration is only enabled in local cluster mode"); URI url = URI.create(args[0]); diff --git a/src/cloud/lunarsky/store/ClusterMigrate.java b/src/cloud/lunarsky/store/ClusterMigrate.java index 073622b..115fe3b 100644 --- a/src/cloud/lunarsky/store/ClusterMigrate.java +++ b/src/cloud/lunarsky/store/ClusterMigrate.java @@ -24,6 +24,8 @@ public final class ClusterMigrate { (args.length != 2 || !args[0].equals("--apply"))) throw new IllegalArgumentException("Usage: objectstore cluster-migrate --check | --apply node-id-0,node-id-1,node-id-2"); Map env = System.getenv(); + EncryptedVolume.requireConfigured(env, + java.nio.file.Path.of(System.getProperty("java.io.tmpdir"))); if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV"))) throw new IllegalArgumentException("Migration is enabled only in local cluster mode"); List urls = Arrays.stream(env.get("CLUSTER_NODES").split(",", -1)).map(URI::create).toList(); diff --git a/src/cloud/lunarsky/store/ClusterNode.java b/src/cloud/lunarsky/store/ClusterNode.java index 9acab3e..b392de0 100644 --- a/src/cloud/lunarsky/store/ClusterNode.java +++ b/src/cloud/lunarsky/store/ClusterNode.java @@ -32,6 +32,7 @@ public final class ClusterNode implements AutoCloseable { private final FileLock lock; ClusterNode(Path root, String token, String repairToken, UUID hostId) throws IOException { + EncryptedVolume.requireConfigured(System.getenv(), root); 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"); diff --git a/src/cloud/lunarsky/store/ClusterRepair.java b/src/cloud/lunarsky/store/ClusterRepair.java index 15cda68..b2a4eab 100644 --- a/src/cloud/lunarsky/store/ClusterRepair.java +++ b/src/cloud/lunarsky/store/ClusterRepair.java @@ -12,6 +12,7 @@ public final class ClusterRepair { if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV"))) throw new IllegalArgumentException("Cluster repair is only enabled in local cluster mode"); boolean loop = args.length == 1; + ObjectStorage.Limits limits = StorageLimits.fromEnvironment(env); long seconds = Long.parseLong(env.getOrDefault("CLUSTER_MAINTENANCE_INTERVAL_SECONDS", "60")); if (seconds < 1 || seconds > 3600) throw new IllegalArgumentException("Invalid maintenance interval"); boolean gcEnabled = loop && "true".equals(env.get("CLUSTER_GC_ENABLED")); @@ -23,7 +24,8 @@ public final class ClusterRepair { try (ClusterStore store = new ClusterStore(env.get("POSTGRES_JDBC_URL"), env.get("POSTGRES_USER"), env.get("POSTGRES_PASSWORD"), env.get("S3_BUCKET"), Arrays.stream(env.get("CLUSTER_NODES").split(",")).map(URI::create).toList(), - env.get("CLUSTER_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"), 134217728, 2147483648L, + env.get("CLUSTER_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"), + limits.maxObjectBytes(), limits.maxTotalBytes(), "true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS")))) { var report = store.repairOnce(); System.out.println("segments_scanned=" + report.scanned()); diff --git a/src/cloud/lunarsky/store/ClusterStore.java b/src/cloud/lunarsky/store/ClusterStore.java index 036146d..d2d497b 100644 --- a/src/cloud/lunarsky/store/ClusterStore.java +++ b/src/cloud/lunarsky/store/ClusterStore.java @@ -43,6 +43,8 @@ final class ClusterStore implements ObjectStorage, MultipartStorage { ClusterStore(String jdbcUrl, String user, String password, String bucket, List nodeUrls, String token, String repairToken, long maxObject, long maxTotal, boolean testNodeDomains) throws IOException { + EncryptedVolume.requireConfigured(System.getenv(), + java.nio.file.Path.of(System.getProperty("java.io.tmpdir"))); if (jdbcUrl == null || !jdbcUrl.startsWith("jdbc:postgresql://") || user == null || password == null) throw new IllegalArgumentException("Invalid metadata database configuration"); this.jdbcUrl = jdbcUrl; diff --git a/src/cloud/lunarsky/store/DiskStore.java b/src/cloud/lunarsky/store/DiskStore.java index e859ebc..fd0486d 100644 --- a/src/cloud/lunarsky/store/DiskStore.java +++ b/src/cloud/lunarsky/store/DiskStore.java @@ -48,6 +48,7 @@ final class DiskStore implements ObjectStorage { private record VersionRecord(String id, String storageId, boolean marker, long modified) {} DiskStore(Path root, long maxObject, long maxTotal) throws IOException { + EncryptedVolume.requireConfigured(System.getenv(), root); this.root = root; objects = root.resolve("objects"); temporary = root.resolve("pending"); diff --git a/src/cloud/lunarsky/store/EncryptedVolume.java b/src/cloud/lunarsky/store/EncryptedVolume.java new file mode 100644 index 0000000..dc03ef7 --- /dev/null +++ b/src/cloud/lunarsky/store/EncryptedVolume.java @@ -0,0 +1,29 @@ +package cloud.lunarsky.store; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.Map; + +final class EncryptedVolume { + private EncryptedVolume() {} + + static void requireConfigured(Map environment) throws IOException { + requireConfigured(environment, null); + } + + static void requireConfigured(Map environment, Path dataDirectory) throws IOException { + String marker = environment.getOrDefault("ENCRYPTED_VOLUME_MARKER_FILE", ""); + String expected = environment.getOrDefault("ENCRYPTED_VOLUME_ID", ""); + if (marker.isEmpty() && expected.isEmpty()) return; + if (!expected.matches("[a-f0-9]{32}") || marker.isEmpty()) + throw new IOException("Encrypted storage marker configuration is incomplete"); + Path path = Path.of(marker); + if (!path.isAbsolute() || dataDirectory != null && + !path.normalize().equals(dataDirectory.toAbsolutePath().normalize().resolve(".objectstore-encrypted")) || + Files.isSymbolicLink(path) || !Files.isRegularFile(path) || + Files.size(path) > 33 || !Files.readString(path, StandardCharsets.US_ASCII).trim().equals(expected)) + throw new IOException("Encrypted storage volume is unavailable or does not match its marker"); + } +} diff --git a/src/cloud/lunarsky/store/Main.java b/src/cloud/lunarsky/store/Main.java index 0a47d44..34641a6 100644 --- a/src/cloud/lunarsky/store/Main.java +++ b/src/cloud/lunarsky/store/Main.java @@ -36,7 +36,8 @@ public final class Main { private final String bucket; private final MultipartStorage multipart; private final ClientLimits clientLimits; - private final Semaphore slots = new Semaphore(16); + private final Semaphore slots; + private final String virtualHostSuffix; Main(DiskStore store, SigV4 authentication, String bucket) throws IOException { this(store, new MultipartStore(store), authentication, bucket); @@ -48,12 +49,21 @@ public final class Main { Main(ObjectStorage store, MultipartStorage multipart, SigV4 authentication, String bucket, ClientLimits clientLimits) throws IOException { + this(store, multipart, authentication, bucket, clientLimits, 16, ""); + } + + Main(ObjectStorage store, MultipartStorage multipart, SigV4 authentication, String bucket, + ClientLimits clientLimits, int maxInFlight, String virtualHostSuffix) throws IOException { + if (maxInFlight < 1 || maxInFlight > 1024) + throw new IllegalArgumentException("MAX_IN_FLIGHT_REQUESTS must be between 1 and 1024"); this.store = store; this.multipart = multipart; this.authentication = authentication; this.owner = authentication.root(); this.bucket = bucket; this.clientLimits = clientLimits; + this.slots = new Semaphore(maxInFlight); + this.virtualHostSuffix = validateVirtualHostSuffix(virtualHostSuffix); store.ensureBucket(bucket); } @@ -102,7 +112,8 @@ public final class Main { capabilities(exchange); return; } - if (path.equals("/")) { + String hostBucket = virtualHostBucket(exchange); + if (path.equals("/") && hostBucket == null) { requireOwner(verified.principal()); if (!exchange.getRequestMethod().equals("GET") || !(query.isEmpty() || query.size() == 1 && "ListBuckets".equals(query.get("x-id")))) @@ -112,16 +123,41 @@ public final class Main { return; } int slash = path.indexOf('/', 1); - String requestedBucket = slash < 0 ? path.substring(1) : path.substring(1, slash); + String requestedBucket = hostBucket == null ? + (slash < 0 ? path.substring(1) : path.substring(1, slash)) : hostBucket; if (requestedBucket.isEmpty()) throw new StoreException(404, "NoSuchBucket", "Bucket not found"); - if (slash < 0 || slash == path.length() - 1) { + if (hostBucket != null ? path.equals("/") : slash < 0 || slash == path.length() - 1) { handleBucket(exchange, query, verified.payload(), requestedBucket, verified.principal()); } else { store.bucket(requestedBucket); - handleObject(exchange, path.substring(slash + 1), query, verified, requestedBucket); + handleObject(exchange, hostBucket == null ? path.substring(slash + 1) : path.substring(1), + query, verified, requestedBucket); } } + private static String validateVirtualHostSuffix(String value) { + if (value == null || value.isBlank()) return ""; + String suffix = value.toLowerCase(Locale.ROOT); + if (suffix.length() > 253 || !suffix.matches("[a-z0-9-]+(?:\\.[a-z0-9-]+)*") || + java.util.Arrays.stream(suffix.split("\\.")).anyMatch( + label -> label.length() > 63 || label.startsWith("-") || label.endsWith("-"))) + throw new IllegalArgumentException("Invalid S3_VIRTUAL_HOST_SUFFIX"); + return suffix; + } + + private String virtualHostBucket(HttpExchange exchange) { + if (virtualHostSuffix.isEmpty()) return null; + String host = exchange.getRequestHeaders().getFirst("Host"); + if (host == null) throw new StoreException(400, "InvalidRequest", "Missing Host header"); + host = host.toLowerCase(Locale.ROOT).replaceFirst(":\\d+$", ""); + String ending = "." + virtualHostSuffix; + if (!host.endsWith(ending)) return null; + String bucketName = host.substring(0, host.length() - ending.length()); + if (!bucketName.matches("[a-z0-9][a-z0-9-]{1,61}[a-z0-9]")) + throw new StoreException(400, "InvalidRequest", "Invalid bucket hostname"); + return bucketName; + } + private void sendStoreError(HttpExchange exchange, StoreException error, String requestId) throws IOException { if (error.status == 503 && error.code.equals("SlowDown")) exchange.getResponseHeaders().set("Retry-After", "1"); @@ -1197,10 +1233,8 @@ public final class Main { 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", "134217728")); - long maxTotal = Long.parseLong(env.getOrDefault("MAX_TOTAL_BYTES", "2147483648")); - if (maxObject < 1 || maxObject > 1073741824L || maxTotal < maxObject) - throw new IllegalArgumentException("Invalid size limits"); + ObjectStorage.Limits limits = StorageLimits.fromEnvironment(env); + long maxObject = limits.maxObjectBytes(), maxTotal = limits.maxTotalBytes(); String mode = env.getOrDefault("STORE_MODE", "disk"); ObjectStorage store; MultipartStorage multipart; @@ -1220,11 +1254,13 @@ public final class Main { multipart = new MultipartStore(disk); } else throw new IllegalArgumentException("Invalid STORE_MODE"); var app = new Main(store, multipart, new SigV4(identities(env, access, secret), - access, region, Clock.systemUTC()), bucket, ClientLimits.fromEnvironment(env)); + access, region, Clock.systemUTC()), bucket, ClientLimits.fromEnvironment(env), + positiveInt(env, "MAX_IN_FLIGHT_REQUESTS", 16, 1024), + env.getOrDefault("S3_VIRTUAL_HOST_SUFFIX", "")); int port = Integer.parseInt(env.getOrDefault("PORT", "9000")); var server = HttpServer.create(mode.equals("cluster") ? new InetSocketAddress(env.getOrDefault("BIND_ADDRESS", "127.0.0.1"), port) - : new InetSocketAddress(port), 64); + : new InetSocketAddress(port), positiveInt(env, "HTTP_BACKLOG", 64, 4096)); var executor = Executors.newVirtualThreadPerTaskExecutor(); server.setExecutor(executor); server.createContext("/", app::handle); @@ -1291,4 +1327,16 @@ public final class Main { if (value == null || value.isBlank()) throw new IllegalArgumentException("Missing " + key); return value; } + + private static int positiveInt(Map env, String key, int fallback, int maximum) { + String configured = env.get(key); + if (configured == null || configured.isBlank()) return fallback; + try { + int value = Integer.parseInt(configured); + if (value < 1 || value > maximum) throw new NumberFormatException(); + return value; + } catch (NumberFormatException error) { + throw new IllegalArgumentException(key + " must be between 1 and " + maximum, error); + } + } } diff --git a/src/cloud/lunarsky/store/StorageLimits.java b/src/cloud/lunarsky/store/StorageLimits.java new file mode 100644 index 0000000..b8a7323 --- /dev/null +++ b/src/cloud/lunarsky/store/StorageLimits.java @@ -0,0 +1,19 @@ +package cloud.lunarsky.store; + +import java.util.Map; + +final class StorageLimits { + private StorageLimits() {} + + static ObjectStorage.Limits fromEnvironment(Map environment) { + try { + long maxObject = Long.parseLong(environment.getOrDefault("MAX_OBJECT_BYTES", "134217728")); + long maxTotal = Long.parseLong(environment.getOrDefault("MAX_TOTAL_BYTES", "2147483648")); + if (maxObject < 1 || maxObject > 1073741824L || maxTotal < maxObject) + throw new NumberFormatException(); + return new ObjectStorage.Limits(maxObject, maxTotal); + } catch (NumberFormatException error) { + throw new IllegalArgumentException("Invalid storage size limits", error); + } + } +} diff --git a/test/cloud/lunarsky/store/EncryptedVolumeTest.java b/test/cloud/lunarsky/store/EncryptedVolumeTest.java new file mode 100644 index 0000000..5390e7d --- /dev/null +++ b/test/cloud/lunarsky/store/EncryptedVolumeTest.java @@ -0,0 +1,42 @@ +package cloud.lunarsky.store; + +import java.io.IOException; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.Map; + +public final class EncryptedVolumeTest { + private static void rejected(Map configuration) throws Exception { + try { + EncryptedVolume.requireConfigured(configuration); + throw new AssertionError("Unavailable encrypted storage was accepted"); + } catch (IOException expected) { } + } + + public static void main(String[] args) throws Exception { + Path directory = Files.createTempDirectory("objectstore-encrypted-volume-"); + Path marker = directory.resolve(".objectstore-encrypted"); + String id = "0123456789abcdef0123456789abcdef"; + Map configuration = Map.of( + "ENCRYPTED_VOLUME_MARKER_FILE", marker.toString(), "ENCRYPTED_VOLUME_ID", id); + try { + EncryptedVolume.requireConfigured(Map.of()); + rejected(configuration); + Files.writeString(marker, id + "\n"); + EncryptedVolume.requireConfigured(configuration); + rejected(Map.of("ENCRYPTED_VOLUME_MARKER_FILE", marker.toString())); + rejected(Map.of("ENCRYPTED_VOLUME_MARKER_FILE", marker.toString(), + "ENCRYPTED_VOLUME_ID", "fedcba9876543210fedcba9876543210")); + Files.delete(marker); + Path elsewhere = directory.resolve("elsewhere"); + Files.writeString(elsewhere, id + "\n"); + Files.createSymbolicLink(marker, elsewhere); + rejected(configuration); + } finally { + Files.deleteIfExists(marker); + Files.deleteIfExists(directory.resolve("elsewhere")); + Files.delete(directory); + } + System.out.println("Encrypted volume marker tests passed"); + } +} diff --git a/test/cloud/lunarsky/store/HttpTest.java b/test/cloud/lunarsky/store/HttpTest.java index 1e4a804..17d2e78 100644 --- a/test/cloud/lunarsky/store/HttpTest.java +++ b/test/cloud/lunarsky/store/HttpTest.java @@ -806,15 +806,70 @@ public final class HttpTest { throw new AssertionError("Completed multipart version was not retained"); } + private static int virtualHostPut(int port, String host, String signedHost, byte[] body) throws Exception { + URI signed = URI.create("http://" + signedHost + ":" + port + "/photos/cat.jpg"); + HttpRequest request = signedUri(signed, "PUT", body, Map.of()); + String headers = "PUT /photos/cat.jpg HTTP/1.1\r\nHost: " + host + ":" + port + + "\r\nAuthorization: " + request.headers().firstValue("authorization").orElseThrow() + + "\r\nx-amz-date: " + request.headers().firstValue("x-amz-date").orElseThrow() + + "\r\nx-amz-content-sha256: " + request.headers().firstValue("x-amz-content-sha256").orElseThrow() + + "\r\nContent-Length: " + body.length + "\r\nConnection: close\r\n\r\n"; + try (var socket = new java.net.Socket("127.0.0.1", port)) { + socket.setSoTimeout(5000); + socket.getOutputStream().write(headers.getBytes(StandardCharsets.US_ASCII)); + socket.getOutputStream().write(body); + String response = new String(socket.getInputStream().readAllBytes(), StandardCharsets.ISO_8859_1); + return Integer.parseInt(response.split(" ", 3)[1]); + } + } + + private static void testGatewayConcurrency(DiskStore store, HttpClient client, + java.util.concurrent.Executor executor) throws Exception { + HttpServer limited = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 16); + var app = new Main(store, new MultipartStore(store), + new SigV4(ACCESS, SECRET, REGION, Clock.systemUTC()), "objects", + ClientLimits.disabled(), 1, ""); + limited.setExecutor(executor); + limited.createContext("/", app::handle); + limited.start(); + int port = limited.getAddress().getPort(); + String base = "http://127.0.0.1:" + port; + byte[] body = {42}; + HttpRequest request = signedUri(URI.create(base + "/objects/held"), "PUT", body, Map.of()); + try (var socket = new java.net.Socket("127.0.0.1", port)) { + socket.setSoTimeout(5000); + String headers = "PUT /objects/held HTTP/1.1\r\nHost: 127.0.0.1:" + port + + "\r\nAuthorization: " + request.headers().firstValue("authorization").orElseThrow() + + "\r\nx-amz-date: " + request.headers().firstValue("x-amz-date").orElseThrow() + + "\r\nx-amz-content-sha256: " + request.headers().firstValue("x-amz-content-sha256").orElseThrow() + + "\r\nContent-Length: 1\r\nConnection: close\r\n\r\n"; + socket.getOutputStream().write(headers.getBytes(StandardCharsets.US_ASCII)); + boolean refused = false; + for (int attempt = 0; attempt < 50 && !refused; attempt++) { + int status = client.send(HttpRequest.newBuilder(URI.create(base + "/health")).GET().build(), + HttpResponse.BodyHandlers.discarding()).statusCode(); + refused = status == 503; + if (!refused) Thread.sleep(10); + } + if (!refused) throw new AssertionError("Configured gateway concurrency cap was not enforced"); + socket.getOutputStream().write(body); + String finished = new String(socket.getInputStream().readAllBytes(), StandardCharsets.ISO_8859_1); + if (!finished.startsWith("HTTP/1.1 200 ")) + throw new AssertionError("Held request did not complete after releasing its body"); + } finally { + limited.stop(0); + } + } + public static void main(String[] args) throws Exception { Path root = Files.createTempDirectory("store-http-test-"); var executor = Executors.newVirtualThreadPerTaskExecutor(); HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 16); DiskStore store = new DiskStore(root, 1024, 4096); try { - var app = new Main(store, + var app = new Main(store, new MultipartStore(store), new SigV4(Map.of(ACCESS, SECRET, SECONDARY, SECONDARY_SECRET), - ACCESS, REGION, Clock.systemUTC()), "objects"); + ACCESS, REGION, Clock.systemUTC()), "objects", ClientLimits.disabled(), 16, "s3.test"); server.setExecutor(executor); server.createContext("/", app::handle); server.start(); @@ -834,6 +889,17 @@ public final class HttpTest { testBuckets(client, base); testVersioning(client, base); testAcl(client, base); + byte[] hostedBody = "virtual host".getBytes(StandardCharsets.UTF_8); + int port = server.getAddress().getPort(); + if (virtualHostPut(port, "objects.s3.test", "objects.s3.test", hostedBody) != 200) + throw new AssertionError("Signed virtual-hosted upload failed"); + try (var opened = store.open("objects", "photos/cat.jpg")) { + if (!java.util.Arrays.equals(hostedBody, opened.stream().readAllBytes())) + throw new AssertionError("Virtual-hosted bucket or key was parsed incorrectly"); + } + if (virtualHostPut(port, "objects.s3.test", "other.s3.test", hostedBody) != 403) + throw new AssertionError("Changing a signed virtual hostname was accepted"); + testGatewayConcurrency(store, client, executor); System.out.println("HTTP tests passed: capabilities, objects, copy, checksums, listing, multipart, attributes, buckets, versioning, ACLs"); } finally { server.stop(0);