6 Commits
Author SHA1 Message Date
admin 510e519bd1 Configure gateway limits and encrypted storage 2026-10-11 01:24:09 +02:00
admin 43987bbadb Route cluster metadata to writable primary 2026-10-11 00:36:16 +02:00
admin a71c4fa5f1 Support TLS for cluster node transport 2026-10-10 23:57:38 +02:00
Electricane 1d1f329220 Update Main.java
To-do Update Main.java to reduce Complexity and readability.
2026-10-10 23:04:16 +02:00
Solunex 91f6825726 Update .env.cluster.example
Container Image will refuse to start without Access and Secret Key.
Replacing empty values with placeholders.
2026-10-10 22:55:24 +02:00
admin 885239be1c Address ObjectStore quality findings 2026-10-10 16:07:13 +02:00
38 changed files with 1175 additions and 194 deletions

No files matched your search

+11 -3
View File
@@ -1,7 +1,15 @@
S3_ACCESS_KEY= S3_ACCESS_KEY=[Replace_with_your_S3_Access_Key]
S3_SECRET_KEY= S3_SECRET_KEY=[Replace_with_your_S3_Secret_Key]
CLUSTER_TOKEN= CLUSTER_TOKEN=
CLUSTER_REPAIR_TOKEN= CLUSTER_REPAIR_TOKEN=
POSTGRES_PASSWORD= POSTGRES_PASSWORD=[PG_Pass]
S3_BUCKET=objects S3_BUCKET=objects
CLUSTER_HOST_PORT=9001 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=
+6
View File
@@ -13,3 +13,9 @@ PUBLIC_BYTES_PER_SECOND=0
PUBLIC_BYTE_BURST= PUBLIC_BYTE_BURST=
PUBLIC_MAX_IN_FLIGHT_PER_IP= PUBLIC_MAX_IN_FLIGHT_PER_IP=
PUBLIC_TRUSTED_PROXY_IPS= 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=
+2
View File
@@ -11,7 +11,9 @@ 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 -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.HttpTest
RUN java --add-modules jdk.httpserver,java.net.http -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.ClientLimitsTest 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.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 RUN java -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.CliTest
FROM eclipse-temurin:21-jre-alpine FROM eclipse-temurin:21-jre-alpine
+57 -3
View File
@@ -16,11 +16,14 @@ Source: [GitHub](https://github.com/LunarSkyOSS/ObjectStore) · [Gitea mirror](h
- [S3 API support checklist](#s3-api-support-checklist) - [S3 API support checklist](#s3-api-support-checklist)
- [Tests](TESTS.md) - [Tests](TESTS.md)
- [Single-node setup](#single-node-setup) - [Single-node setup](#single-node-setup)
- [Encrypted storage](#encrypted-storage)
- [Access keys and ACLs](#access-keys-and-acls) - [Access keys and ACLs](#access-keys-and-acls)
- [CLI and tests](#cli-and-tests) - [CLI and tests](#cli-and-tests)
- [Capability discovery](#capability-discovery) - [Capability discovery](#capability-discovery)
- [Java client](#java-client) - [Java client](#java-client)
- [Local cluster prototype](#local-cluster-prototype) - [Local cluster prototype](#local-cluster-prototype)
- [Node transport TLS](#node-transport-tls)
- [Metadata primary routing](#metadata-primary-routing)
- [Migrating a local cluster](#migrating-a-local-cluster) - [Migrating a local cluster](#migrating-a-local-cluster)
- [Adding a cluster node](#adding-a-cluster-node) - [Adding a cluster node](#adding-a-cluster-node)
- [Cluster maintenance and recovery](#cluster-maintenance-and-recovery) - [Cluster maintenance and recovery](#cluster-maintenance-and-recovery)
@@ -33,9 +36,11 @@ 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. ObjectStore creates the configured default bucket at startup. Additional buckets share the configured capacity limit.
- ✅ Persistent single-node storage with checksum verification - ✅ Persistent single-node storage with checksum verification
- ✅ Optional dm-crypt-backed data mounts with startup guards
- ✅ Configurable per-object and total logical size limits - ✅ Configurable per-object and total logical size limits
- ✅ CLI status, version, and full payload verification - ✅ CLI status, version, and full payload verification
- ✅ Local cluster prototype with stable node IDs and host-aware placement code - ✅ Local cluster prototype with stable node IDs and host-aware placement code
- ✅ Optional HTTPS between cluster processes and storage nodes
- ✅ Opt-in automatic repair, rebalance, and guarded garbage collection in the local cluster - ✅ Opt-in automatic repair, rebalance, and guarded garbage collection in the local cluster
- ✅ Metadata backup and tested restore to a separate local PostgreSQL instance - ✅ Metadata backup and tested restore to a separate local PostgreSQL instance
- ✅ Manual [two-machine durability and metadata-restore drill](tests/two-host/README.md) - ✅ Manual [two-machine durability and metadata-restore drill](tests/two-host/README.md)
@@ -46,6 +51,7 @@ New objects retain their content type and key. Objects written by the earlier si
## S3 API support checklist ## S3 API support checklist
- ✅ Header-based and presigned-query AWS Signature Version 4 authentication - ✅ 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 - ✅ `PutObject`, `GetObject`, `HeadObject`, and `DeleteObject` in both modes
- ✅ Single-range GET and `ListObjectsV2` in both modes - ✅ Single-range GET and `ListObjectsV2` in both modes
- ✅ SHA-256 payload verification and `x-amz-checksum-sha256` in both modes - ✅ SHA-256 payload verification and `x-amz-checksum-sha256` in both modes
@@ -81,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. 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 ## 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. `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.
@@ -117,7 +140,7 @@ The [JDK-only Java client](client/README.md) works with ObjectStore and other S3
## Local cluster prototype ## Local cluster prototype
The local cluster prototype starts three segment containers and one PostgreSQL container on the same Docker host. Copy `.env.cluster.example` to a private environment file, replace all four credentials, and run: The local cluster prototype starts three segment containers and one PostgreSQL container on the same Docker host. Copy `.env.cluster.example` to a private environment file, replace all five credential values, and run:
```sh ```sh
docker compose --env-file /path/to/cluster.env -f compose.cluster.yaml up -d --build docker compose --env-file /path/to/cluster.env -f compose.cluster.yaml up -d --build
@@ -131,6 +154,33 @@ Multipart parts are stored on cluster nodes and indexed in PostgreSQL. Incomplet
Node UUIDs persist on their volumes, and replica manifests use those UUIDs so reordering configured URLs cannot move an existing replica. Each node also has an operator-assigned physical host UUID. New writes require acknowledgements from two different host UUIDs. The optional `CLUSTER_TEST_NODE_DOMAINS=true` override counts containers instead, solely for local process tests; all containers in this Compose file share one physical host. Node UUIDs persist on their volumes, and replica manifests use those UUIDs so reordering configured URLs cannot move an existing replica. Each node also has an operator-assigned physical host UUID. New writes require acknowledgements from two different host UUIDs. The optional `CLUSTER_TEST_NODE_DOMAINS=true` override counts containers instead, solely for local process tests; all containers in this Compose file share one physical host.
## Node transport TLS
The default local Compose cluster uses HTTP inside its private Docker network. For an HTTPS test, give each node a PKCS#12 keystore containing its private key and a certificate whose DNS subject alternative name matches its `CLUSTER_NODES` hostname. Give the gateway, repair, garbage collection, and maintenance processes a PKCS#12 truststore containing the issuing CA or each node certificate. Mount the files read-only and keep the keystores and password files outside Git.
Set `NODE_TLS_KEYSTORE` and `NODE_TLS_PASSWORD_FILE` on each node. Set `CLUSTER_TLS_TRUSTSTORE` and `CLUSTER_TLS_PASSWORD_FILE` on every process that contacts nodes, and change each node URL to `https://`. With a truststore configured, HTTP node URLs are rejected. The client verifies the certificate chain and hostname; a failed handshake does not fall back to HTTP. Both settings in each pair are required. Restart affected processes after rotating certificates or truststores.
The optional [Compose TLS overlay](compose.cluster.tls.yaml) expects `node-a.p12` through `node-d.p12`, matching `.pass` files, and `trust.p12` with `trust.pass` in `CLUSTER_TLS_DIR`. Set that variable to a private certificate directory and include both Compose files:
```sh
CLUSTER_TLS_DIR=/private/objectstore-certs docker compose --env-file /path/to/cluster.env \
-f compose.cluster.yaml -f compose.cluster.tls.yaml up -d --build
```
The files must be readable by container UID 10001 without making private keys or passwords world-readable. Add HTTPS URLs for additional nodes when expanding the cluster.
This secures node traffic only. The local cluster still lacks automatic PostgreSQL failover, database TLS configuration, encryption at rest, and production multi-server validation. Its HTTP S3 gateway remains bound to localhost; use a separate trusted proxy for external TLS. Do not treat the TLS overlay as a production deployment.
## Metadata primary routing
`POSTGRES_JDBC_URL` can override the database URL for the gateway, repair, garbage collection, and maintenance processes. Its local Compose default now uses `targetServerType=primary`, and `/ready` returns unavailable when the database is read-only or in recovery. The base Compose file still starts and waits for its own single PostgreSQL container; it is not an HA deployment. In an independently managed deployment, list the PostgreSQL hosts in the JDBC URL and keep `targetServerType=primary`:
```text
jdbc:postgresql://db-a:5432,db-b:5432/objectstore?targetServerType=primary&connectTimeout=3&socketTimeout=10
```
ObjectStore opens a new database connection for each operation, so the JDBC driver can select a promoted primary after the old one is stopped. This does **not** promote a standby, fence the old primary, configure synchronous replication, or guarantee that a recently acknowledged write reached the standby. Those jobs belong to a separately operated PostgreSQL HA system. Never allow two writable metadata databases: they can diverge while serving different ObjectStore requests. A request interrupted during failover has an uncertain outcome; verify it before retrying a non-idempotent operation. For remote database connections, configure PostgreSQL TLS and use JDBC `sslmode=verify-full` with a mounted CA certificate. The provided local Compose database does not enable TLS.
## Migrating a local cluster ## Migrating a local cluster
For an existing **local** three-node cluster that stores replicas by URL position: For an existing **local** three-node cluster that stores replicas by URL position:
@@ -161,7 +211,11 @@ The maintenance service is opt-in. It repairs missing or corrupt replicas and re
## Limits and safety ## 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. 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.
@@ -169,7 +223,7 @@ ObjectStore uses the socket peer as the client IP and ignores forwarded-IP heade
Enable these limits only on a deliberately public endpoint. They also apply to direct localhost storage calls to that same endpoint; leave them disabled for the local-only setup or run a separate local-only instance if local storage calls must be exempt. A direct loopback `/health` probe without a forwarded-IP header remains exempt. The byte limit is aggregate ingress plus egress for each IP, and several users behind one NAT share it. If a proxy buffers complete uploads before forwarding them, the upload byte limit controls proxy-to-ObjectStore traffic, not the client's initial upload speed; disable request buffering when end-to-end upload pacing is required. It is a fairness control, not a defense against connection floods before the Java handler runs. Put an internet-facing proxy or firewall in front of the gateway for TLS, connection limits, request timeouts, and buffering controls. Do not expose storage nodes or PostgreSQL publicly. Enable these limits only on a deliberately public endpoint. They also apply to direct localhost storage calls to that same endpoint; leave them disabled for the local-only setup or run a separate local-only instance if local storage calls must be exempt. A direct loopback `/health` probe without a forwarded-IP header remains exempt. The byte limit is aggregate ingress plus egress for each IP, and several users behind one NAT share it. If a proxy buffers complete uploads before forwarding them, the upload byte limit controls proxy-to-ObjectStore traffic, not the client's initial upload speed; disable request buffering when end-to-end upload pacing is required. It is a fairness control, not a defense against connection floods before the Java handler runs. Put an internet-facing proxy or firewall in front of the gateway for TLS, connection limits, request timeouts, and buffering controls. Do not expose storage nodes or PostgreSQL publicly.
The local cluster has no automatic metadata failover, private-network TLS, scoped credentials, or physical host verification. Garbage collection can remove data required by an older metadata backup, so its retention guard is essential. Host UUIDs are operator labels, not proof that machines have separate power, disks, or network paths. Keep `CLUSTER_LOCAL_DEV=true` limited to local tests. The default local cluster still uses plaintext node and PostgreSQL connections. The optional TLS overlay secures node traffic, but PostgreSQL TLS, automatic metadata failover, scoped credentials, and physical host verification remain absent. Garbage collection can remove data required by an older metadata backup, so its retention guard is essential. Host UUIDs are operator labels, not proof that machines have separate power, disks, or network paths. Keep `CLUSTER_LOCAL_DEV=true` limited to local tests.
The standalone cluster node binds to localhost by default. Set `NODE_BIND` only for a private test network; the Compose file binds inside its private Docker network. PostgreSQL JDBC 42.7.14 is bundled in the image with its license inside the JAR. The standalone cluster node binds to localhost by default. Set `NODE_BIND` only for a private test network; the Compose file binds inside its private Docker network. PostgreSQL JDBC 42.7.14 is bundled in the image with its license inside the JAR.
+12 -4
View File
@@ -18,17 +18,25 @@ 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. | | `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. | | `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. | | `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. | | `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. | | `CliTest` | Version, status, verification, and a nonzero result for corrupt data. |
| `ClientTest` and `MultipartClientTest` | Java client request signing, capability discovery, error handling, and multipart operations. | | `ClientTest` and `MultipartClientTest` | Java client request signing, capability discovery, error handling, and multipart operations. |
The script exits nonzero on failure. The test programs use temporary local directories and loopback HTTP ports; they do not use an existing ObjectStore volume. 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.
## Disposable Docker cluster tests ## Disposable Docker cluster tests
Requires Docker with Compose, Python 3, `curl`, and a free local port 9001. Make a test-only environment file from `.env.cluster.example` and fill in all five blank credentials with test-only values. Keep that file private and out of Git. Requires Docker with Compose, Python 3.9 or newer, `curl`, and a free local port 9001. Make a test-only environment file from `.env.cluster.example` and replace all five credential values with test-only values. Keep that file private and out of Git.
```sh ```sh
cp .env.cluster.example /tmp/objectstore-cluster-tests.env cp .env.cluster.example /tmp/objectstore-cluster-tests.env
@@ -49,7 +57,7 @@ COMPOSE_PROJECT_NAME=objectstore-tests docker compose --env-file /tmp/objectstor
If port 9001 is occupied, set `CLUSTER_HOST_PORT` to the same free port in both the environment file and the shell before running the script. The script reads that port from the shell; Compose reads it from the file. If port 9001 is occupied, set `CLUSTER_HOST_PORT` to the same free port in both the environment file and the shell before running the script. The script reads that port from the shell; Compose reads it from the file.
The Docker suite checks signed capability discovery and S3 operations, bucket creation and deletion, metadata and tags, public-read ACLs and anonymous access, copies, upload checksums, multi-segment objects, concurrent overwrites, multipart staging and listings, versioned reads and delete markers, versioned multipart completion, completion after a gateway restart and node loss, reads and writes with a node stopped, refusal to write without a storage quorum, restart recovery, corrupt-replica repair including staged parts, metadata unavailability, and placement on a newly joined node. It then checks rebalance to the fourth node, automatic repair, garbage collection dry run and delayed deletion, and a metadata backup restored to a separate PostgreSQL instance while the primary is stopped. Historical regular and multipart versions are checked after repair and cleanup. It also checks that containers labeled as one physical host cannot satisfy the normal host quorum. Its local-only override permits the remaining phases to use containers as separate test domains. The Docker suite checks signed capability discovery and S3 operations, bucket creation and deletion, metadata and tags, public-read ACLs and anonymous access, copies, upload checksums, multi-segment objects, concurrent overwrites, multipart staging and listings, versioned reads and delete markers, versioned multipart completion, completion after a gateway restart and node loss, reads and writes with a node stopped, refusal to write without a storage quorum, restart recovery, corrupt-replica repair including staged parts, metadata unavailability, read-only metadata readiness rejection, and placement on a newly joined node. It then checks rebalance to the fourth node, automatic repair, garbage collection dry run and delayed deletion, and a metadata backup restored to a separate PostgreSQL instance while the original database is stopped. A two-host JDBC URL selects that restored writable database. Historical regular and multipart versions are checked after repair and cleanup. It also checks that containers labeled as one physical host cannot satisfy the normal host quorum. Its local-only override permits the remaining phases to use containers as separate test domains.
`ClusterMigrationTest` is a separate legacy-format fixture and is **not** run by either test script. Do not run its `create` phase against a populated metadata database. The migration procedure is in the [README](README.md#migrating-a-local-cluster). `ClusterMigrationTest` is a separate legacy-format fixture and is **not** run by either test script. Do not run its `create` phase against a populated metadata database. The migration procedure is in the [README](README.md#migrating-a-local-cluster).
+7 -2
View File
@@ -371,7 +371,9 @@ public final class ImageManager extends Frame {
Object value = event.getTransferable().getTransferData(DataFlavor.javaFileListFlavor); Object value = event.getTransferable().getTransferData(DataFlavor.javaFileListFlavor);
List<?> dropped = (List<?>) value; List<?> dropped = (List<?>) value;
List<java.io.File> files = new ArrayList<>(); List<java.io.File> files = new ArrayList<>();
for (Object item : dropped) if (item instanceof java.io.File file) files.add(file); for (Object item : dropped) {
if (item instanceof java.io.File file) files.add(file);
}
event.dropComplete(true); event.dropComplete(true);
EventQueue.invokeLater(() -> uploadImages(files)); EventQueue.invokeLater(() -> uploadImages(files));
} catch (Exception error) { } catch (Exception error) {
@@ -391,7 +393,10 @@ public final class ImageManager extends Frame {
Button cancel = new Button("Cancel"); Button cancel = new Button("Cancel");
Button proceed = new Button("Continue"); Button proceed = new Button("Continue");
cancel.addActionListener(event -> dialog.dispose()); cancel.addActionListener(event -> dialog.dispose());
proceed.addActionListener(event -> { accepted[0] = true; dialog.dispose(); }); proceed.addActionListener(event -> {
accepted[0] = true;
dialog.dispose();
});
buttons.add(cancel); buttons.add(cancel);
buttons.add(proceed); buttons.add(proceed);
dialog.add(buttons, BorderLayout.SOUTH); dialog.add(buttons, BorderLayout.SOUTH);
@@ -124,7 +124,10 @@ final class Json {
return result.toString(); return result.toString();
} }
if (current < 0x20) throw new ProtocolException("Unescaped control character in JSON string"); if (current < 0x20) throw new ProtocolException("Unescaped control character in JSON string");
if (current != '\\') { result.append(current); continue; } if (current != '\\') {
result.append(current);
continue;
}
if (index >= source.length()) throw new ProtocolException("Incomplete JSON escape"); if (index >= source.length()) throw new ProtocolException("Incomplete JSON escape");
char escaped = source.charAt(index++); char escaped = source.charAt(index++);
switch (escaped) { switch (escaped) {
@@ -159,31 +162,42 @@ final class Json {
private BigDecimal number() throws ProtocolException { private BigDecimal number() throws ProtocolException {
int start = index; int start = index;
if (take('-') && index >= source.length()) throw new ProtocolException("Invalid JSON number"); if (take('-') && index >= source.length()) throw new ProtocolException("Invalid JSON number");
integerDigits();
if (take('.')) requireDigits("Invalid JSON fraction");
if (take('e') || take('E')) {
if (!take('+')) take('-');
requireDigits("Invalid JSON exponent");
}
try { return new BigDecimal(source.substring(start, index)); }
catch (NumberFormatException e) { throw new ProtocolException("Invalid JSON number", e); }
}
private void integerDigits() throws ProtocolException {
if (take('0')) { if (take('0')) {
if (index < source.length() && Character.isDigit(source.charAt(index))) if (index < source.length() && Character.isDigit(source.charAt(index)))
throw new ProtocolException("Invalid JSON number"); throw new ProtocolException("Invalid JSON number");
} else { } else {
if (index >= source.length() || source.charAt(index) < '1' || source.charAt(index) > '9') if (index >= source.length() || source.charAt(index) < '1' || source.charAt(index) > '9')
throw new ProtocolException("Invalid JSON number"); throw new ProtocolException("Invalid JSON number");
while (index < source.length() && source.charAt(index) >= '0' && source.charAt(index) <= '9') index++; scanDigits();
} }
if (take('.')) { }
int first = index;
while (index < source.length() && source.charAt(index) >= '0' && source.charAt(index) <= '9') index++; private void requireDigits(String message) throws ProtocolException {
if (first == index) throw new ProtocolException("Invalid JSON fraction"); int first = index;
} scanDigits();
if (take('e') || take('E')) { if (first == index) throw new ProtocolException(message);
if (!take('+')) take('-'); }
int first = index;
while (index < source.length() && source.charAt(index) >= '0' && source.charAt(index) <= '9') index++; private void scanDigits() {
if (first == index) throw new ProtocolException("Invalid JSON exponent"); while (index < source.length() && source.charAt(index) >= '0' && source.charAt(index) <= '9') index++;
}
try { return new BigDecimal(source.substring(start, index)); }
catch (NumberFormatException e) { throw new ProtocolException("Invalid JSON number", e); }
} }
private boolean take(char value) { private boolean take(char value) {
if (index < source.length() && source.charAt(index) == value) { index++; return true; } if (index < source.length() && source.charAt(index) == value) {
index++;
return true;
}
return false; return false;
} }
+114
View File
@@ -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
+58
View File
@@ -0,0 +1,58 @@
x-client-tls: &client-tls
CLUSTER_NODES: https://node-a:9100,https://node-b:9100,https://node-c:9100
CLUSTER_TLS_TRUSTSTORE: /run/objectstore-tls/trust.p12
CLUSTER_TLS_PASSWORD_FILE: /run/objectstore-tls/trust.pass
x-tls-volume: &tls-volume
- ${CLUSTER_TLS_DIR:?Set CLUSTER_TLS_DIR to a private certificate directory}:/run/objectstore-tls:ro
x-node-health: &node-health
test: ["CMD", "nc", "-z", "-w", "2", "127.0.0.1", "9100"]
interval: 10s
timeout: 3s
retries: 3
services:
gateway:
environment: *client-tls
volumes: *tls-volume
repair:
environment: *client-tls
volumes: *tls-volume
gc:
environment: *client-tls
volumes: *tls-volume
maintenance:
environment: *client-tls
volumes: *tls-volume
node-a:
environment:
NODE_TLS_KEYSTORE: /run/objectstore-tls/node-a.p12
NODE_TLS_PASSWORD_FILE: /run/objectstore-tls/node-a.pass
volumes: *tls-volume
healthcheck: *node-health
node-b:
environment:
NODE_TLS_KEYSTORE: /run/objectstore-tls/node-b.p12
NODE_TLS_PASSWORD_FILE: /run/objectstore-tls/node-b.pass
volumes: *tls-volume
healthcheck: *node-health
node-c:
environment:
NODE_TLS_KEYSTORE: /run/objectstore-tls/node-c.p12
NODE_TLS_PASSWORD_FILE: /run/objectstore-tls/node-c.pass
volumes: *tls-volume
healthcheck: *node-health
node-d:
environment:
NODE_TLS_KEYSTORE: /run/objectstore-tls/node-d.p12
NODE_TLS_PASSWORD_FILE: /run/objectstore-tls/node-d.pass
volumes: *tls-volume
healthcheck: *node-health
+10 -3
View File
@@ -15,6 +15,9 @@ services:
S3_CREDENTIALS_FILE: ${S3_CREDENTIALS_FILE:-} S3_CREDENTIALS_FILE: ${S3_CREDENTIALS_FILE:-}
MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728} MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728}
MAX_TOTAL_BYTES: ${MAX_TOTAL_BYTES:-2147483648} 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_REQUESTS_PER_SECOND: ${PUBLIC_REQUESTS_PER_SECOND:-0}
PUBLIC_REQUEST_BURST: ${PUBLIC_REQUEST_BURST:-} PUBLIC_REQUEST_BURST: ${PUBLIC_REQUEST_BURST:-}
PUBLIC_BYTES_PER_SECOND: ${PUBLIC_BYTES_PER_SECOND:-0} PUBLIC_BYTES_PER_SECOND: ${PUBLIC_BYTES_PER_SECOND:-0}
@@ -23,7 +26,7 @@ services:
PUBLIC_TRUSTED_PROXY_IPS: ${PUBLIC_TRUSTED_PROXY_IPS:-} PUBLIC_TRUSTED_PROXY_IPS: ${PUBLIC_TRUSTED_PROXY_IPS:-}
CLUSTER_TOKEN: ${CLUSTER_TOKEN:?Set CLUSTER_TOKEN} CLUSTER_TOKEN: ${CLUSTER_TOKEN:?Set CLUSTER_TOKEN}
CLUSTER_NODES: ${CLUSTER_NODES:-http://node-a:9100,http://node-b:9100,http://node-c:9100} CLUSTER_NODES: ${CLUSTER_NODES:-http://node-a:9100,http://node-b:9100,http://node-c:9100}
POSTGRES_JDBC_URL: jdbc:postgresql://metadata:5432/objectstore?connectTimeout=3&socketTimeout=10 POSTGRES_JDBC_URL: ${POSTGRES_JDBC_URL:-jdbc:postgresql://metadata:5432/objectstore?connectTimeout=3&socketTimeout=10&targetServerType=primary}
POSTGRES_USER: objectstore POSTGRES_USER: objectstore
POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD} POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD}
ports: ports:
@@ -96,9 +99,11 @@ services:
CLUSTER_TOKEN: ${CLUSTER_TOKEN:?Set CLUSTER_TOKEN} CLUSTER_TOKEN: ${CLUSTER_TOKEN:?Set CLUSTER_TOKEN}
CLUSTER_REPAIR_TOKEN: ${CLUSTER_REPAIR_TOKEN:?Set CLUSTER_REPAIR_TOKEN} CLUSTER_REPAIR_TOKEN: ${CLUSTER_REPAIR_TOKEN:?Set CLUSTER_REPAIR_TOKEN}
S3_BUCKET: ${S3_BUCKET:-objects} S3_BUCKET: ${S3_BUCKET:-objects}
POSTGRES_JDBC_URL: jdbc:postgresql://metadata:5432/objectstore?connectTimeout=3&socketTimeout=10 POSTGRES_JDBC_URL: ${POSTGRES_JDBC_URL:-jdbc:postgresql://metadata:5432/objectstore?connectTimeout=3&socketTimeout=10&targetServerType=primary}
POSTGRES_USER: objectstore POSTGRES_USER: objectstore
POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD} 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_GC_MIN_AGE_SECONDS: ${CLUSTER_GC_MIN_AGE_SECONDS:-1209600}
CLUSTER_BACKUP_RETENTION_SECONDS: ${CLUSTER_BACKUP_RETENTION_SECONDS:-0} CLUSTER_BACKUP_RETENTION_SECONDS: ${CLUSTER_BACKUP_RETENTION_SECONDS:-0}
CLUSTER_GC_TEST_MODE: ${CLUSTER_GC_TEST_MODE:-false} CLUSTER_GC_TEST_MODE: ${CLUSTER_GC_TEST_MODE:-false}
@@ -131,9 +136,11 @@ services:
CLUSTER_BACKUP_RETENTION_SECONDS: ${CLUSTER_BACKUP_RETENTION_SECONDS:-0} CLUSTER_BACKUP_RETENTION_SECONDS: ${CLUSTER_BACKUP_RETENTION_SECONDS:-0}
CLUSTER_GC_TEST_MODE: ${CLUSTER_GC_TEST_MODE:-false} CLUSTER_GC_TEST_MODE: ${CLUSTER_GC_TEST_MODE:-false}
S3_BUCKET: ${S3_BUCKET:-objects} S3_BUCKET: ${S3_BUCKET:-objects}
POSTGRES_JDBC_URL: jdbc:postgresql://metadata:5432/objectstore?connectTimeout=3&socketTimeout=10 POSTGRES_JDBC_URL: ${POSTGRES_JDBC_URL:-jdbc:postgresql://metadata:5432/objectstore?connectTimeout=3&socketTimeout=10&targetServerType=primary}
POSTGRES_USER: objectstore POSTGRES_USER: objectstore
POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD} POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD}
MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728}
MAX_TOTAL_BYTES: ${MAX_TOTAL_BYTES:-2147483648}
depends_on: depends_on:
metadata: metadata:
condition: service_healthy condition: service_healthy
+11
View File
@@ -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
+3
View File
@@ -11,6 +11,9 @@ services:
S3_CREDENTIALS_FILE: ${S3_CREDENTIALS_FILE:-} S3_CREDENTIALS_FILE: ${S3_CREDENTIALS_FILE:-}
MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728} MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728}
MAX_TOTAL_BYTES: ${MAX_TOTAL_BYTES:-2147483648} 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_REQUESTS_PER_SECOND: ${PUBLIC_REQUESTS_PER_SECOND:-0}
PUBLIC_REQUEST_BURST: ${PUBLIC_REQUEST_BURST:-} PUBLIC_REQUEST_BURST: ${PUBLIC_REQUEST_BURST:-}
PUBLIC_BYTES_PER_SECOND: ${PUBLIC_BYTES_PER_SECOND:-0} PUBLIC_BYTES_PER_SECOND: ${PUBLIC_BYTES_PER_SECOND:-0}
+69
View File
@@ -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"
+26 -7
View File
@@ -22,12 +22,31 @@ port = int(values.get("CLUSTER_HOST_PORT", "9001"))
if not 1 <= port <= 65535: if not 1 <= port <= 65535:
raise ValueError("CLUSTER_HOST_PORT must be between 1 and 65535") raise ValueError("CLUSTER_HOST_PORT must be between 1 and 65535")
host = f"127.0.0.1:{port}" host = f"127.0.0.1:{port}"
MAX_RESPONSE_BYTES = 1024 * 1024
def sign(key, message): def sign(key, message):
return hmac.new(key, message.encode(), hashlib.sha256).digest() return hmac.new(key, message.encode(), hashlib.sha256).digest()
def parse_xml(content):
if len(content) > MAX_RESPONSE_BYTES:
raise ValueError("Unsafe XML response from test server")
text = content.decode("utf-8")
if "<!DOCTYPE" in text or "<!ENTITY" in text:
raise ValueError("Unsafe XML response from test server")
parser = ET.XMLParser()
parser.feed(text)
return parser.close()
def read_response(response):
content = response.read(MAX_RESPONSE_BYTES + 1)
if len(content) > MAX_RESPONSE_BYTES:
raise ValueError("Oversized response from test server")
return content
def request(method, path, body=b"", extra=None): def request(method, path, body=b"", extra=None):
extra = extra or {} extra = extra or {}
date = datetime.datetime.now(datetime.timezone.utc).strftime("%Y%m%dT%H%M%SZ") date = datetime.datetime.now(datetime.timezone.utc).strftime("%Y%m%dT%H%M%SZ")
@@ -50,7 +69,7 @@ def request(method, path, body=b"", extra=None):
try: try:
connection.request(method, path, body=body if method in ("PUT", "POST") else None, headers=headers) connection.request(method, path, body=body if method in ("PUT", "POST") else None, headers=headers)
response = connection.getresponse() response = connection.getresponse()
return response.status, response.read(), response.headers return response.status, read_response(response), response.headers
finally: finally:
connection.close() connection.close()
@@ -60,7 +79,7 @@ def anonymous(method, path):
try: try:
connection.request(method, path) connection.request(method, path)
response = connection.getresponse() response = connection.getresponse()
return response.status, response.read() return response.status, read_response(response)
finally: finally:
connection.close() connection.close()
@@ -75,7 +94,7 @@ if len(sys.argv) > 2 and sys.argv[2] == "acl":
path = f"/{bucket}/cluster-test/acl-multipart" path = f"/{bucket}/cluster-test/acl-multipart"
status, content, _ = request("POST", path + "?uploads", extra={"x-amz-acl": "public-read"}) status, content, _ = request("POST", path + "?uploads", extra={"x-amz-acl": "public-read"})
assert status == 200, (status, content) assert status == 200, (status, content)
upload_id = ET.fromstring(content).findtext("UploadId") upload_id = parse_xml(content).findtext("UploadId")
status, _, headers = request("PUT", path + f"?partNumber=1&uploadId={upload_id}", b"public part") status, _, headers = request("PUT", path + f"?partNumber=1&uploadId={upload_id}", b"public part")
assert status == 200, status assert status == 200, status
completion = ("<CompleteMultipartUpload><Part><PartNumber>1</PartNumber><ETag>" + completion = ("<CompleteMultipartUpload><Part><PartNumber>1</PartNumber><ETag>" +
@@ -102,9 +121,9 @@ if len(sys.argv) > 4 and sys.argv[2] == "status":
if len(sys.argv) > 2 and sys.argv[2] == "version-survivor": if len(sys.argv) > 2 and sys.argv[2] == "version-survivor":
status, listing, _ = request("GET", "/version-bucket?versions") status, listing, _ = request("GET", "/version-bucket?versions")
assert status == 200, status assert status == 200, status
root = ET.fromstring(listing) root = parse_xml(listing)
namespace = {"s3": "http://s3.amazonaws.com/doc/2006-03-01/"} namespace = {"s3": "http://s3.amazonaws.com/doc/2006-03-01/"}
expected_etag = '"' + hashlib.md5(b"older cluster version").hexdigest() + '"' expected_etag = '"' + hashlib.md5(b"older cluster version", usedforsecurity=False).hexdigest() + '"'
versions = [version for version in root.findall("s3:Version", namespace) versions = [version for version in root.findall("s3:Version", namespace)
if version.findtext("s3:ETag", namespaces=namespace) == expected_etag] if version.findtext("s3:ETag", namespaces=namespace) == expected_etag]
assert len(versions) == 1, listing assert len(versions) == 1, listing
@@ -166,7 +185,7 @@ copy_source = f"/{bucket}/cluster-test/copy-source.txt"
copy_target = f"/{bucket}/cluster-test/copied.txt" copy_target = f"/{bucket}/cluster-test/copied.txt"
body = b"cluster copy and checksum test" body = b"cluster copy and checksum test"
crc32 = base64.b64encode(zlib.crc32(body).to_bytes(4, "big")).decode() crc32 = base64.b64encode(zlib.crc32(body).to_bytes(4, "big")).decode()
md5 = base64.b64encode(hashlib.md5(body).digest()).decode() md5 = base64.b64encode(hashlib.md5(body, usedforsecurity=False).digest()).decode()
status, _, headers = request("PUT", copy_source, body, status, _, headers = request("PUT", copy_source, body,
{"content-type": "text/plain", "content-md5": md5, {"content-type": "text/plain", "content-md5": md5,
"x-amz-checksum-crc32": crc32, "x-amz-checksum-crc32": crc32,
@@ -257,7 +276,7 @@ assert status == 200 and b"<DeleteMarker>" in content and old_version.encode() i
versioned_multipart = "/version-bucket/multipart.txt" versioned_multipart = "/version-bucket/multipart.txt"
status, content, _ = request("POST", versioned_multipart + "?uploads") status, content, _ = request("POST", versioned_multipart + "?uploads")
assert status == 200, (status, content) assert status == 200, (status, content)
versioned_upload = ET.fromstring(content).findtext("UploadId") versioned_upload = parse_xml(content).findtext("UploadId")
assert versioned_upload, content assert versioned_upload, content
versioned_part = b"retained multipart version" versioned_part = b"retained multipart version"
status, _, headers = request("PUT", versioned_multipart + status, _, headers = request("PUT", versioned_multipart +
+10 -1
View File
@@ -8,6 +8,8 @@ compose() { docker compose --env-file "$env_file" -f compose.cluster.yaml "$@";
restore() { restore() {
compose stop maintenance >/dev/null 2>&1 || true compose stop maintenance >/dev/null 2>&1 || true
compose start metadata node-a node-b node-c >/dev/null 2>&1 || true compose start metadata node-a node-b node-c >/dev/null 2>&1 || true
compose exec -T metadata psql -U objectstore -d postgres -c \
'ALTER DATABASE objectstore RESET default_transaction_read_only' >/dev/null 2>&1 || true
if [ -n "$backup_dir" ]; then rm -rf "$backup_dir"; fi if [ -n "$backup_dir" ]; then rm -rf "$backup_dir"; fi
} }
trap restore EXIT trap restore EXIT
@@ -59,6 +61,13 @@ expected=$(compose exec -T metadata psql -U objectstore -d objectstore -At -c \
"SELECT encode(s.sha256,'hex') FROM cluster_segments s JOIN cluster_objects o ON o.generation=s.generation WHERE o.object_key='cluster-test/survivor' LIMIT 1") "SELECT encode(s.sha256,'hex') FROM cluster_segments s JOIN cluster_objects o ON o.generation=s.generation WHERE o.object_key='cluster-test/survivor' LIMIT 1")
actual=$(compose exec -T node-a sha256sum "/data/segments/$shard/$segment_id" | cut -d' ' -f1) actual=$(compose exec -T node-a sha256sum "/data/segments/$shard/$segment_id" | cut -d' ' -f1)
[ "$expected" = "$actual" ] [ "$expected" = "$actual" ]
compose exec -T metadata psql -U objectstore -d postgres -c \
'ALTER DATABASE objectstore SET default_transaction_read_only=on' >/dev/null
status=$(curl -sS -o /dev/null -w '%{http_code}' "http://127.0.0.1:$host_port/ready")
[ "$status" = 503 ]
compose exec -T metadata psql -U objectstore -d postgres -c \
'ALTER DATABASE objectstore RESET default_transaction_read_only' >/dev/null
wait_ready
compose stop metadata compose stop metadata
status=$(curl -sS -o /dev/null -w '%{http_code}' "http://127.0.0.1:$host_port/ready") status=$(curl -sS -o /dev/null -w '%{http_code}' "http://127.0.0.1:$host_port/ready")
[ "$status" = 503 ] [ "$status" = 503 ]
@@ -124,7 +133,7 @@ compose exec -T metadata-recovery pg_restore -U objectstore -d objectstore --no-
< "$backup_dir/metadata.dump" < "$backup_dir/metadata.dump"
compose stop metadata compose stop metadata
compose run --rm -T --no-deps \ compose run --rm -T --no-deps \
-e 'POSTGRES_JDBC_URL=jdbc:postgresql://metadata-recovery:5432/objectstore?connectTimeout=3&socketTimeout=10' \ -e 'POSTGRES_JDBC_URL=jdbc:postgresql://metadata:5432,metadata-recovery:5432/objectstore?connectTimeout=3&socketTimeout=10&targetServerType=primary&hostRecheckSeconds=0' \
--entrypoint java gateway --add-modules jdk.httpserver,java.net.http \ --entrypoint java gateway --add-modules jdk.httpserver,java.net.http \
-cp /app:/app/postgresql.jar:/app/hash4j.jar cloud.lunarsky.store.ClusterIntegrationTest recovered -cp /app:/app/postgresql.jar:/app/hash4j.jar cloud.lunarsky.store.ClusterIntegrationTest recovered
compose start metadata compose start metadata
+7
View File
@@ -0,0 +1,7 @@
#!/bin/sh
set -eu
cd "$(dirname "$0")/.."
project="objectstore-metadata-test-$$"
compose() { docker compose -p "$project" -f tests/metadata/compose.yaml "$@"; }
trap 'compose down -v --remove-orphans >/dev/null 2>&1 || true' EXIT
compose up --build --abort-on-container-exit --exit-code-from check check
+2
View File
@@ -8,6 +8,8 @@ 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 -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.HttpTest
java --add-modules jdk.httpserver,java.net.http -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.ClientLimitsTest 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.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 java -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.CliTest
bash client/scripts/test.sh bash client/scripts/test.sh
@@ -105,42 +105,51 @@ final class AwsChunkedInputStream extends FilterInputStream {
catch (NumberFormatException error) { throw invalid("Invalid signed chunk size"); } catch (NumberFormatException error) { throw invalid("Invalid signed chunk size"); }
if (chunkLeft > decodedLength - decoded) throw invalid("Signed chunks exceed decoded length"); if (chunkLeft > decodedLength - decoded) throw invalid("Signed chunks exceed decoded length");
chunkHash.reset(); chunkHash.reset();
if (chunkLeft == 0) { if (chunkLeft == 0) finishPayload();
finishChunk(); }
if (decoded != decodedLength) throw invalid("Decoded length mismatch");
if (trailerName == null) { private void finishPayload() throws IOException {
if (!line().isEmpty()) throw invalid("Invalid signed chunk ending"); finishChunk();
} else { if (decoded != decodedLength) throw invalid("Decoded length mismatch");
String trailer = line(); if (trailerName == null) {
if (!trailer.startsWith(trailerName + ":")) throw invalid("Missing signed checksum trailer"); if (!line().isEmpty()) throw invalid("Invalid signed chunk ending");
trailerValue = trailer.substring(trailerName.length() + 1); } else {
byte[] actual; verifyTrailer();
if (trailerCrc != null) {
long value = trailerCrc.getValue();
actual = new byte[trailerName.equals("x-amz-checksum-crc64nvme") ? 8 : 4];
for (int i = actual.length - 1; i >= 0; i--) {
actual[i] = (byte) value;
value >>>= 8;
}
} else actual = trailerXxhash != null ? trailerXxhash.digest() : trailerHash.digest();
if (!Base64.getEncoder().encodeToString(actual).equals(trailerValue))
throw new StoreException(400, "BadDigest", "Checksum trailer mismatch");
String signature = line();
if (!signature.matches("x-amz-trailer-signature=[0-9a-f]{64}"))
throw invalid("Missing trailer signature");
String toSign = "AWS4-HMAC-SHA256-TRAILER\n" + authorization.date() + "\n" +
authorization.scope() + "\n" + previousSignature + "\n" +
SigV4.hex(SigV4.hash((trailerName + ":" + trailerValue + "\n")
.getBytes(StandardCharsets.UTF_8)));
String expected = SigV4.hex(SigV4.hmac(authorization.signingKey(), toSign));
if (!MessageDigest.isEqual(expected.getBytes(StandardCharsets.US_ASCII),
signature.substring(24).getBytes(StandardCharsets.US_ASCII)))
throw invalid("Trailer signature mismatch");
if (!line().isEmpty()) throw invalid("Invalid trailer ending");
}
if (in.read() != -1) throw invalid("Extra bytes after signed payload");
finished = true;
} }
if (in.read() != -1) throw invalid("Extra bytes after signed payload");
finished = true;
}
private void verifyTrailer() throws IOException {
String trailer = line();
if (!trailer.startsWith(trailerName + ":")) throw invalid("Missing signed checksum trailer");
trailerValue = trailer.substring(trailerName.length() + 1);
if (!Base64.getEncoder().encodeToString(trailerChecksum()).equals(trailerValue))
throw new StoreException(400, "BadDigest", "Checksum trailer mismatch");
String signature = line();
if (!signature.matches("x-amz-trailer-signature=[0-9a-f]{64}"))
throw invalid("Missing trailer signature");
String toSign = "AWS4-HMAC-SHA256-TRAILER\n" + authorization.date() + "\n" +
authorization.scope() + "\n" + previousSignature + "\n" +
SigV4.hex(SigV4.hash((trailerName + ":" + trailerValue + "\n")
.getBytes(StandardCharsets.UTF_8)));
String expected = SigV4.hex(SigV4.hmac(authorization.signingKey(), toSign));
if (!MessageDigest.isEqual(expected.getBytes(StandardCharsets.US_ASCII),
signature.substring(24).getBytes(StandardCharsets.US_ASCII)))
throw invalid("Trailer signature mismatch");
if (!line().isEmpty()) throw invalid("Invalid trailer ending");
}
private byte[] trailerChecksum() {
if (trailerCrc == null)
return trailerXxhash != null ? trailerXxhash.digest() : trailerHash.digest();
long value = trailerCrc.getValue();
byte[] actual = new byte[trailerName.equals("x-amz-checksum-crc64nvme") ? 8 : 4];
for (int i = actual.length - 1; i >= 0; i--) {
actual[i] = (byte) value;
value >>>= 8;
}
return actual;
} }
private void finishChunk() throws IOException { private void finishChunk() throws IOException {
+25 -23
View File
@@ -114,29 +114,7 @@ final class ClientLimits {
if (exchange.getRemoteAddress().getAddress().isLoopbackAddress() && if (exchange.getRemoteAddress().getAddress().isLoopbackAddress() &&
exchange.getRequestHeaders().get("X-Real-IP") == null && exchange.getRequestHeaders().get("X-Real-IP") == null &&
path.equals("/health")) return null; path.equals("/health")) return null;
String address = address(exchange); Client client = admit(address(exchange));
Client client;
synchronized (this) {
long now = System.nanoTime();
if (++admissions % 1024 == 0 || clients.size() >= MAX_CLIENTS)
clients.entrySet().removeIf(entry -> entry.getValue().inFlight == 0 &&
now - entry.getValue().lastSeen > IDLE_NANOS);
client = clients.get(address);
if (client == null) {
if (clients.size() >= MAX_CLIENTS)
throw new StoreException(503, "SlowDown", "Client limit table is full");
client = new Client(now, requestBurst, byteBurst);
clients.put(address, client);
}
refill(client, now);
client.lastSeen = now;
if (maxInFlight > 0 && client.inFlight >= maxInFlight)
throw new StoreException(503, "SlowDown", "Too many concurrent requests from this client");
if (requestsPerSecond > 0 && client.requestTokens < 1)
throw new StoreException(503, "SlowDown", "Client request rate exceeded");
if (requestsPerSecond > 0) client.requestTokens--;
client.inFlight++;
}
if (bytesPerSecond > 0) { if (bytesPerSecond > 0) {
try { try {
exchange.setStreams(new LimitedInput(exchange.getRequestBody(), client), exchange.setStreams(new LimitedInput(exchange.getRequestBody(), client),
@@ -149,6 +127,30 @@ final class ClientLimits {
return client; return client;
} }
private synchronized Client admit(String address) {
Client client;
long now = System.nanoTime();
if (++admissions % 1024 == 0 || clients.size() >= MAX_CLIENTS)
clients.entrySet().removeIf(entry -> entry.getValue().inFlight == 0 &&
now - entry.getValue().lastSeen > IDLE_NANOS);
client = clients.get(address);
if (client == null) {
if (clients.size() >= MAX_CLIENTS)
throw new StoreException(503, "SlowDown", "Client limit table is full");
client = new Client(now, requestBurst, byteBurst);
clients.put(address, client);
}
refill(client, now);
client.lastSeen = now;
if (maxInFlight > 0 && client.inFlight >= maxInFlight)
throw new StoreException(503, "SlowDown", "Too many concurrent requests from this client");
if (requestsPerSecond > 0 && client.requestTokens < 1)
throw new StoreException(503, "SlowDown", "Client request rate exceeded");
if (requestsPerSecond > 0) client.requestTokens--;
client.inFlight++;
return client;
}
synchronized void leave(Client client) { synchronized void leave(Client client) {
if (client != null) { if (client != null) {
client.inFlight--; client.inFlight--;
+3 -1
View File
@@ -15,6 +15,7 @@ public final class ClusterGc {
if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV"))) 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"); throw new IllegalArgumentException("Cluster garbage collection is only enabled in local cluster mode");
boolean testDomains = "true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS")); 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")); boolean disposableTest = testDomains && "true".equals(env.get("CLUSTER_GC_TEST_MODE"));
long age = Long.parseLong(env.getOrDefault("CLUSTER_GC_MIN_AGE_SECONDS", "1209600")); long age = Long.parseLong(env.getOrDefault("CLUSTER_GC_MIN_AGE_SECONDS", "1209600"));
long backupRetention = Long.parseLong(env.getOrDefault("CLUSTER_BACKUP_RETENTION_SECONDS", "0")); 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"), try (ClusterStore store = new ClusterStore(env.get("POSTGRES_JDBC_URL"), env.get("POSTGRES_USER"),
env.get("POSTGRES_PASSWORD"), env.get("S3_BUCKET"), env.get("POSTGRES_PASSWORD"), env.get("S3_BUCKET"),
Arrays.stream(env.get("CLUSTER_NODES").split(",")).map(URI::create).toList(), 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)) { testDomains)) {
if (apply) { if (apply) {
var repair = store.repairOnce(); var repair = store.repairOnce();
@@ -12,6 +12,8 @@ public final class ClusterJoin {
if (args.length != 2) if (args.length != 2)
throw new IllegalArgumentException("Usage: objectstore cluster-join node-url expected-host-uuid"); throw new IllegalArgumentException("Usage: objectstore cluster-join node-url expected-host-uuid");
Map<String, String> env = System.getenv(); Map<String, String> 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"))) 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"); throw new IllegalArgumentException("Node registration is only enabled in local cluster mode");
URI url = URI.create(args[0]); URI url = URI.create(args[0]);
@@ -24,6 +24,8 @@ public final class ClusterMigrate {
(args.length != 2 || !args[0].equals("--apply"))) (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"); throw new IllegalArgumentException("Usage: objectstore cluster-migrate --check | --apply node-id-0,node-id-1,node-id-2");
Map<String, String> env = System.getenv(); Map<String, String> 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"))) 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"); throw new IllegalArgumentException("Migration is enabled only in local cluster mode");
List<URI> urls = Arrays.stream(env.get("CLUSTER_NODES").split(",", -1)).map(URI::create).toList(); List<URI> urls = Arrays.stream(env.get("CLUSTER_NODES").split(",", -1)).map(URI::create).toList();
+5 -3
View File
@@ -1,7 +1,6 @@
package cloud.lunarsky.store; package cloud.lunarsky.store;
import com.sun.net.httpserver.HttpExchange; import com.sun.net.httpserver.HttpExchange;
import com.sun.net.httpserver.HttpServer;
import java.io.IOException; import java.io.IOException;
import java.io.InputStream; import java.io.InputStream;
import java.io.OutputStream; import java.io.OutputStream;
@@ -33,6 +32,7 @@ public final class ClusterNode implements AutoCloseable {
private final FileLock lock; private final FileLock lock;
ClusterNode(Path root, String token, String repairToken, UUID hostId) throws IOException { 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 (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)) if (repairToken == null || repairToken.length() < 32 || repairToken.equals(token))
throw new IllegalArgumentException("A separate repair token of at least 32 characters is required"); throw new IllegalArgumentException("A separate repair token of at least 32 characters is required");
@@ -325,7 +325,8 @@ public final class ClusterNode implements AutoCloseable {
env.get("CLUSTER_REPAIR_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"),
UUID.fromString(env.get("CLUSTER_HOST_ID"))); UUID.fromString(env.get("CLUSTER_HOST_ID")));
int port = Integer.parseInt(env.getOrDefault("NODE_PORT", "9100")); int port = Integer.parseInt(env.getOrDefault("NODE_PORT", "9100"));
var server = HttpServer.create(new InetSocketAddress(env.getOrDefault("NODE_BIND", "127.0.0.1"), port), 64); var server = ClusterTls.nodeServer(
new InetSocketAddress(env.getOrDefault("NODE_BIND", "127.0.0.1"), port), env);
var executor = Executors.newVirtualThreadPerTaskExecutor(); var executor = Executors.newVirtualThreadPerTaskExecutor();
server.setExecutor(executor); server.setExecutor(executor);
server.createContext("/", node::handle); server.createContext("/", node::handle);
@@ -336,6 +337,7 @@ public final class ClusterNode implements AutoCloseable {
catch (IOException error) { System.err.println("Node close failed: " + error); } catch (IOException error) { System.err.println("Node close failed: " + error); }
})); }));
server.start(); server.start();
System.out.println("ObjectStore cluster node listening on :" + port); System.out.println("ObjectStore cluster node listening on :" + port +
(server instanceof com.sun.net.httpserver.HttpsServer ? " (TLS)" : " (HTTP)"));
} }
} }
+3 -1
View File
@@ -12,6 +12,7 @@ public final class ClusterRepair {
if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV"))) 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"); throw new IllegalArgumentException("Cluster repair is only enabled in local cluster mode");
boolean loop = args.length == 1; boolean loop = args.length == 1;
ObjectStorage.Limits limits = StorageLimits.fromEnvironment(env);
long seconds = Long.parseLong(env.getOrDefault("CLUSTER_MAINTENANCE_INTERVAL_SECONDS", "60")); long seconds = Long.parseLong(env.getOrDefault("CLUSTER_MAINTENANCE_INTERVAL_SECONDS", "60"));
if (seconds < 1 || seconds > 3600) throw new IllegalArgumentException("Invalid maintenance interval"); if (seconds < 1 || seconds > 3600) throw new IllegalArgumentException("Invalid maintenance interval");
boolean gcEnabled = loop && "true".equals(env.get("CLUSTER_GC_ENABLED")); 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"), try (ClusterStore store = new ClusterStore(env.get("POSTGRES_JDBC_URL"), env.get("POSTGRES_USER"),
env.get("POSTGRES_PASSWORD"), env.get("S3_BUCKET"), env.get("POSTGRES_PASSWORD"), env.get("S3_BUCKET"),
Arrays.stream(env.get("CLUSTER_NODES").split(",")).map(URI::create).toList(), 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")))) { "true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS")))) {
var report = store.repairOnce(); var report = store.repairOnce();
System.out.println("segments_scanned=" + report.scanned()); System.out.println("segments_scanned=" + report.scanned());
+5 -2
View File
@@ -43,6 +43,8 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
ClusterStore(String jdbcUrl, String user, String password, String bucket, ClusterStore(String jdbcUrl, String user, String password, String bucket,
List<URI> nodeUrls, String token, String repairToken, long maxObject, long maxTotal, List<URI> nodeUrls, String token, String repairToken, long maxObject, long maxTotal,
boolean testNodeDomains) throws IOException { 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) if (jdbcUrl == null || !jdbcUrl.startsWith("jdbc:postgresql://") || user == null || password == null)
throw new IllegalArgumentException("Invalid metadata database configuration"); throw new IllegalArgumentException("Invalid metadata database configuration");
this.jdbcUrl = jdbcUrl; this.jdbcUrl = jdbcUrl;
@@ -1347,8 +1349,9 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
@Override public boolean ready() { @Override public boolean ready() {
if (!nodes.availableHostsAtLeast(2, testNodeDomains)) return false; if (!nodes.availableHostsAtLeast(2, testNodeDomains)) return false;
try (Connection connection = connect(); var statement = connection.createStatement(); try (Connection connection = connect(); var statement = connection.createStatement();
ResultSet result = statement.executeQuery("SELECT 1")) { ResultSet result = statement.executeQuery(
return result.next() && result.getInt(1) == 1; "SELECT NOT pg_is_in_recovery() AND current_setting('transaction_read_only') = 'off'")) {
return result.next() && result.getBoolean(1);
} catch (SQLException error) { return false; } } catch (SQLException error) { return false; }
} }
RepairReport repairOnce() throws IOException { RepairReport repairOnce() throws IOException {
+99
View File
@@ -0,0 +1,99 @@
package cloud.lunarsky.store;
import com.sun.net.httpserver.HttpServer;
import com.sun.net.httpserver.HttpsConfigurator;
import com.sun.net.httpserver.HttpsServer;
import java.io.IOException;
import java.io.InputStream;
import java.net.InetSocketAddress;
import java.net.http.HttpClient;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.security.GeneralSecurityException;
import java.security.KeyStore;
import java.time.Duration;
import java.util.Arrays;
import java.util.Map;
import javax.net.ssl.KeyManagerFactory;
import javax.net.ssl.SSLContext;
import javax.net.ssl.TrustManagerFactory;
final class ClusterTls {
private ClusterTls() { }
static HttpServer nodeServer(InetSocketAddress address, Map<String, String> env) throws IOException {
String keyStore = env.get("NODE_TLS_KEYSTORE");
String passwordFile = env.get("NODE_TLS_PASSWORD_FILE");
if (missing(keyStore) && missing(passwordFile)) return HttpServer.create(address, 64);
if (missing(keyStore) || missing(passwordFile))
throw new IOException("Node TLS requires both NODE_TLS_KEYSTORE and NODE_TLS_PASSWORD_FILE");
SSLContext context = serverContext(Path.of(keyStore), Path.of(passwordFile));
HttpsServer server = HttpsServer.create(address, 64);
server.setHttpsConfigurator(new HttpsConfigurator(context));
return server;
}
static HttpClient client(Map<String, String> env, Duration timeout) throws IOException {
String trustStore = env.get("CLUSTER_TLS_TRUSTSTORE");
String passwordFile = env.get("CLUSTER_TLS_PASSWORD_FILE");
if (missing(trustStore) != missing(passwordFile))
throw new IOException("Cluster TLS requires both CLUSTER_TLS_TRUSTSTORE and CLUSTER_TLS_PASSWORD_FILE");
HttpClient.Builder builder = HttpClient.newBuilder().connectTimeout(timeout);
if (!missing(trustStore))
builder.sslContext(clientContext(Path.of(trustStore), Path.of(passwordFile)));
return builder.build();
}
private static SSLContext serverContext(Path keyStore, Path passwordFile) throws IOException {
char[] password = password(passwordFile);
try {
KeyStore keys = load(keyStore, password);
KeyManagerFactory managers = KeyManagerFactory.getInstance(KeyManagerFactory.getDefaultAlgorithm());
managers.init(keys, password);
SSLContext context = SSLContext.getInstance("TLS");
context.init(managers.getKeyManagers(), null, null);
return context;
} catch (GeneralSecurityException error) {
throw new IOException("Could not configure node TLS", error);
} finally {
Arrays.fill(password, '\0');
}
}
private static SSLContext clientContext(Path trustStore, Path passwordFile) throws IOException {
char[] password = password(passwordFile);
try {
KeyStore trust = load(trustStore, password);
TrustManagerFactory managers = TrustManagerFactory.getInstance(TrustManagerFactory.getDefaultAlgorithm());
managers.init(trust);
SSLContext context = SSLContext.getInstance("TLS");
context.init(null, managers.getTrustManagers(), null);
return context;
} catch (GeneralSecurityException error) {
throw new IOException("Could not configure cluster TLS trust", error);
} finally {
Arrays.fill(password, '\0');
}
}
private static KeyStore load(Path path, char[] password) throws IOException, GeneralSecurityException {
KeyStore store = KeyStore.getInstance("PKCS12");
try (InputStream input = Files.newInputStream(path)) {
store.load(input, password);
}
return store;
}
private static char[] password(Path file) throws IOException {
String value = Files.readString(file, StandardCharsets.UTF_8);
if (value.endsWith("\n")) value = value.substring(0, value.length() - 1);
if (value.endsWith("\r")) value = value.substring(0, value.length() - 1);
if (value.isEmpty()) throw new IOException("Cluster TLS password file is empty");
return value.toCharArray();
}
private static boolean missing(String value) {
return value == null || value.isBlank();
}
}
+1
View File
@@ -48,6 +48,7 @@ final class DiskStore implements ObjectStorage {
private record VersionRecord(String id, String storageId, boolean marker, long modified) {} private record VersionRecord(String id, String storageId, boolean marker, long modified) {}
DiskStore(Path root, long maxObject, long maxTotal) throws IOException { DiskStore(Path root, long maxObject, long maxTotal) throws IOException {
EncryptedVolume.requireConfigured(System.getenv(), root);
this.root = root; this.root = root;
objects = root.resolve("objects"); objects = root.resolve("objects");
temporary = root.resolve("pending"); temporary = root.resolve("pending");
@@ -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<String, String> environment) throws IOException {
requireConfigured(environment, null);
}
static void requireConfigured(Map<String, String> 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");
}
}
+117 -55
View File
@@ -24,6 +24,11 @@ import java.util.concurrent.Executors;
import java.util.concurrent.Semaphore; import java.util.concurrent.Semaphore;
public final class Main { public final class Main {
//Todo: Main is too complex. Chop Chop.
private static final String CAPABILITIES_PATH = "/_objectstore/capabilities"; private static final String CAPABILITIES_PATH = "/_objectstore/capabilities";
private final ObjectStorage store; private final ObjectStorage store;
private final SigV4 authentication; private final SigV4 authentication;
@@ -31,7 +36,8 @@ public final class Main {
private final String bucket; private final String bucket;
private final MultipartStorage multipart; private final MultipartStorage multipart;
private final ClientLimits clientLimits; 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 { Main(DiskStore store, SigV4 authentication, String bucket) throws IOException {
this(store, new MultipartStore(store), authentication, bucket); this(store, new MultipartStore(store), authentication, bucket);
@@ -43,15 +49,26 @@ public final class Main {
Main(ObjectStorage store, MultipartStorage multipart, SigV4 authentication, String bucket, Main(ObjectStorage store, MultipartStorage multipart, SigV4 authentication, String bucket,
ClientLimits clientLimits) throws IOException { 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.store = store;
this.multipart = multipart; this.multipart = multipart;
this.authentication = authentication; this.authentication = authentication;
this.owner = authentication.root(); this.owner = authentication.root();
this.bucket = bucket; this.bucket = bucket;
this.clientLimits = clientLimits; this.clientLimits = clientLimits;
this.slots = new Semaphore(maxInFlight);
this.virtualHostSuffix = validateVirtualHostSuffix(virtualHostSuffix);
store.ensureBucket(bucket); store.ensureBucket(bucket);
} }
//Todo: These methods dont need to be in Main.java.
void handle(HttpExchange exchange) throws IOException { void handle(HttpExchange exchange) throws IOException {
boolean admitted = false; boolean admitted = false;
ClientLimits.Client client = null; ClientLimits.Client client = null;
@@ -63,55 +80,10 @@ public final class Main {
admitted = slots.tryAcquire(); admitted = slots.tryAcquire();
if (!admitted) throw new StoreException(503, "SlowDown", "Too many concurrent requests"); if (!admitted) throw new StoreException(503, "SlowDown", "Too many concurrent requests");
if (handleStatus(exchange)) return; if (handleStatus(exchange)) return;
SigV4.Verified verified = anonymousRead(exchange) dispatch(exchange);
? new SigV4.Verified("UNSIGNED-PAYLOAD", exchange.getRequestURI().getRawQuery(),
null, null, null, null, null)
: authentication.verifyRequest(exchange.getRequestMethod(),
exchange.getRequestURI(), exchange.getRequestHeaders());
String hash = verified.payload();
String principal = verified.principal();
String path = SigV4.decode(exchange.getRequestURI().getRawPath());
Map<String, String> query = query(verified.applicationQuery());
if (path.equals(CAPABILITIES_PATH)) {
requireOwner(principal);
if (!exchange.getRequestMethod().equals("GET")) unsupported("Capability operation");
if (!query.isEmpty())
throw new StoreException(400, "InvalidArgument", "Capability request has unsupported query parameters");
requireEmptyBody(exchange, hash);
capabilities(exchange);
} else if (path.equals("/")) {
requireOwner(principal);
if (!exchange.getRequestMethod().equals("GET") ||
!(query.isEmpty() || query.size() == 1 && "ListBuckets".equals(query.get("x-id"))))
unsupported("Service operation");
requireEmptyBody(exchange, hash);
listBuckets(exchange);
} else {
int slash = path.indexOf('/', 1);
String requestedBucket = slash < 0 ? path.substring(1) : path.substring(1, slash);
if (requestedBucket.isEmpty()) throw new StoreException(404, "NoSuchBucket", "Bucket not found");
if (slash < 0 || slash == path.length() - 1) {
handleBucket(exchange, query, hash, requestedBucket, principal);
} else {
store.bucket(requestedBucket);
handleObject(exchange, path.substring(slash + 1), query, verified, requestedBucket);
}
}
} catch (StoreException error) { } catch (StoreException error) {
if (error.status == 503 && error.code.equals("SlowDown")) sendStoreError(exchange, error, requestId);
exchange.getResponseHeaders().set("Retry-After", "1"); } catch (Exception error) {
if (anonymousRead(exchange) && error.status == 404)
error = new StoreException(403, "AccessDenied", "Access denied");
if (error.deleteMarker) {
exchange.getResponseHeaders().set("x-amz-delete-marker", "true");
exchange.getResponseHeaders().set("x-amz-version-id", error.versionId);
if (error.modified >= 0) exchange.getResponseHeaders().set("Last-Modified",
DateTimeFormatter.RFC_1123_DATE_TIME.withZone(ZoneOffset.UTC)
.format(Instant.ofEpochMilli(error.modified)));
}
sendError(exchange, error.status, error.code, error.getMessage(), requestId);
}
catch (Exception error) {
System.err.println("ObjectStore request failed: " + requestId + " " + error.getClass().getSimpleName()); System.err.println("ObjectStore request failed: " + requestId + " " + error.getClass().getSimpleName());
sendError(exchange, 500, "InternalError", "Storage operation failed", requestId); sendError(exchange, 500, "InternalError", "Storage operation failed", requestId);
} finally { } finally {
@@ -123,6 +95,84 @@ public final class Main {
} }
} }
private void dispatch(HttpExchange exchange) throws IOException {
SigV4.Verified verified = anonymousRead(exchange)
? new SigV4.Verified("UNSIGNED-PAYLOAD", exchange.getRequestURI().getRawQuery(),
null, null, null, null, null)
: authentication.verifyRequest(exchange.getRequestMethod(),
exchange.getRequestURI(), exchange.getRequestHeaders());
String path = SigV4.decode(exchange.getRequestURI().getRawPath());
Map<String, String> query = query(verified.applicationQuery());
if (path.equals(CAPABILITIES_PATH)) {
requireOwner(verified.principal());
if (!exchange.getRequestMethod().equals("GET")) unsupported("Capability operation");
if (!query.isEmpty())
throw new StoreException(400, "InvalidArgument", "Capability request has unsupported query parameters");
requireEmptyBody(exchange, verified.payload());
capabilities(exchange);
return;
}
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"))))
unsupported("Service operation");
requireEmptyBody(exchange, verified.payload());
listBuckets(exchange);
return;
}
int slash = path.indexOf('/', 1);
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 (hostBucket != null ? path.equals("/") : slash < 0 || slash == path.length() - 1) {
handleBucket(exchange, query, verified.payload(), requestedBucket, verified.principal());
} else {
store.bucket(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");
if (anonymousRead(exchange) && error.status == 404)
error = new StoreException(403, "AccessDenied", "Access denied");
if (error.deleteMarker) {
exchange.getResponseHeaders().set("x-amz-delete-marker", "true");
exchange.getResponseHeaders().set("x-amz-version-id", error.versionId);
if (error.modified >= 0) exchange.getResponseHeaders().set("Last-Modified",
DateTimeFormatter.RFC_1123_DATE_TIME.withZone(ZoneOffset.UTC)
.format(Instant.ofEpochMilli(error.modified)));
}
sendError(exchange, error.status, error.code, error.getMessage(), requestId);
}
private static boolean anonymousRead(HttpExchange exchange) { private static boolean anonymousRead(HttpExchange exchange) {
if (!exchange.getRequestMethod().equals("GET") && if (!exchange.getRequestMethod().equals("GET") &&
!exchange.getRequestMethod().equals("HEAD")) return false; !exchange.getRequestMethod().equals("HEAD")) return false;
@@ -1183,10 +1233,8 @@ public final class Main {
if (!access.matches("[A-Za-z0-9]{16,128}") || secret.length() < 32 || 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]")) !bucket.matches("[a-z0-9][a-z0-9-]{1,61}[a-z0-9]"))
throw new IllegalArgumentException("Invalid storage credentials/bucket configuration"); throw new IllegalArgumentException("Invalid storage credentials/bucket configuration");
long maxObject = Long.parseLong(env.getOrDefault("MAX_OBJECT_BYTES", "134217728")); ObjectStorage.Limits limits = StorageLimits.fromEnvironment(env);
long maxTotal = Long.parseLong(env.getOrDefault("MAX_TOTAL_BYTES", "2147483648")); long maxObject = limits.maxObjectBytes(), maxTotal = limits.maxTotalBytes();
if (maxObject < 1 || maxObject > 1073741824L || maxTotal < maxObject)
throw new IllegalArgumentException("Invalid size limits");
String mode = env.getOrDefault("STORE_MODE", "disk"); String mode = env.getOrDefault("STORE_MODE", "disk");
ObjectStorage store; ObjectStorage store;
MultipartStorage multipart; MultipartStorage multipart;
@@ -1206,11 +1254,13 @@ public final class Main {
multipart = new MultipartStore(disk); multipart = new MultipartStore(disk);
} else throw new IllegalArgumentException("Invalid STORE_MODE"); } else throw new IllegalArgumentException("Invalid STORE_MODE");
var app = new Main(store, multipart, new SigV4(identities(env, access, secret), 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")); int port = Integer.parseInt(env.getOrDefault("PORT", "9000"));
var server = HttpServer.create(mode.equals("cluster") var server = HttpServer.create(mode.equals("cluster")
? new InetSocketAddress(env.getOrDefault("BIND_ADDRESS", "127.0.0.1"), port) ? 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(); var executor = Executors.newVirtualThreadPerTaskExecutor();
server.setExecutor(executor); server.setExecutor(executor);
server.createContext("/", app::handle); server.createContext("/", app::handle);
@@ -1277,4 +1327,16 @@ public final class Main {
if (value == null || value.isBlank()) throw new IllegalArgumentException("Missing " + key); if (value == null || value.isBlank()) throw new IllegalArgumentException("Missing " + key);
return value; return value;
} }
private static int positiveInt(Map<String, String> 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);
}
}
} }
+35 -8
View File
@@ -21,18 +21,21 @@ final class NodeClient {
record Node(UUID id, UUID hostId, URI url) {} record Node(UUID id, UUID hostId, URI url) {}
record StoredSegment(UUID id, long modified) {} record StoredSegment(UUID id, long modified) {}
private static final HttpClient IDENTITY_HTTP = HttpClient.newBuilder() private static HttpClient identityHttp;
.connectTimeout(Duration.ofSeconds(2)).build();
private final List<Node> nodes; private final List<Node> nodes;
private final String token; private final String token;
private final String repairToken; private final String repairToken;
private final HttpClient http = HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(3)).build(); private final HttpClient http;
private final ConcurrentHashMap<UUID, Long> unreadableUntil = new ConcurrentHashMap<>(); private final ConcurrentHashMap<UUID, Long> unreadableUntil = new ConcurrentHashMap<>();
private final ConcurrentHashMap<UUID, Long> healthyUntil = new ConcurrentHashMap<>(); private final ConcurrentHashMap<UUID, Long> healthyUntil = new ConcurrentHashMap<>();
private static final long READ_RETRY_NANOS = TimeUnit.SECONDS.toNanos(5); private static final long READ_RETRY_NANOS = TimeUnit.SECONDS.toNanos(5);
private static final long HEALTH_FRESH_NANOS = TimeUnit.SECONDS.toNanos(3); private static final long HEALTH_FRESH_NANOS = TimeUnit.SECONDS.toNanos(3);
NodeClient(List<Node> nodes, String token, String repairToken) { NodeClient(List<Node> nodes, String token, String repairToken) throws IOException {
this(nodes, token, repairToken, ClusterTls.client(System.getenv(), Duration.ofSeconds(3)));
}
NodeClient(List<Node> nodes, String token, String repairToken, HttpClient http) {
if (nodes.isEmpty() || nodes.stream().map(Node::id).distinct().count() != nodes.size() || if (nodes.isEmpty() || nodes.stream().map(Node::id).distinct().count() != nodes.size() ||
nodes.stream().map(Node::url).distinct().count() != nodes.size()) nodes.stream().map(Node::url).distinct().count() != nodes.size())
throw new IllegalArgumentException("Cluster node IDs and URLs must be unique"); throw new IllegalArgumentException("Cluster node IDs and URLs must be unique");
@@ -45,6 +48,7 @@ final class NodeClient {
this.nodes = List.copyOf(nodes); this.nodes = List.copyOf(nodes);
this.token = token; this.token = token;
this.repairToken = repairToken; this.repairToken = repairToken;
this.http = java.util.Objects.requireNonNull(http);
} }
int count() { return nodes.size(); } int count() { return nodes.size(); }
@@ -63,12 +67,16 @@ final class NodeClient {
} }
static NodeIdentity probe(URI url, String token) throws IOException { static NodeIdentity probe(URI url, String token) throws IOException {
return probe(url, token, identityHttp());
}
static NodeIdentity probe(URI url, String token, HttpClient http) throws IOException {
validateUrl(url); validateUrl(url);
if (token == null || token.length() < 32) throw new IllegalArgumentException("Invalid cluster token"); if (token == null || token.length() < 32) throw new IllegalArgumentException("Invalid cluster token");
HttpRequest request = HttpRequest.newBuilder(url.resolve("/identity")) HttpRequest request = HttpRequest.newBuilder(url.resolve("/identity"))
.timeout(Duration.ofSeconds(2)).header("X-Cluster-Token", token).GET().build(); .timeout(Duration.ofSeconds(2)).header("X-Cluster-Token", token).GET().build();
try { try {
HttpResponse<InputStream> response = IDENTITY_HTTP.send(request, HttpResponse.BodyHandlers.ofInputStream()); HttpResponse<InputStream> response = http.send(request, HttpResponse.BodyHandlers.ofInputStream());
try (InputStream body = response.body()) { try (InputStream body = response.body()) {
if (response.statusCode() != 200) throw new IOException("Node identity request failed: " + response.statusCode()); if (response.statusCode() != 200) throw new IOException("Node identity request failed: " + response.statusCode());
byte[] bytes = body.readNBytes(128); byte[] bytes = body.readNBytes(128);
@@ -90,19 +98,38 @@ final class NodeClient {
catch (IOException offline) { return null; } catch (IOException offline) { return null; }
} }
private static NodeIdentity probeIfAvailable(URI url, String token, HttpClient http) {
try { return probe(url, token, http); }
catch (IOException offline) { return null; }
}
private static synchronized HttpClient identityHttp() throws IOException {
if (identityHttp == null)
identityHttp = ClusterTls.client(System.getenv(), Duration.ofSeconds(2));
return identityHttp;
}
static void validateUrl(URI url) { static void validateUrl(URI url) {
if (url == null || !"http".equals(url.getScheme()) || url.getHost() == null || String trustStore = System.getenv("CLUSTER_TLS_TRUSTSTORE");
validateUrl(url, trustStore != null && !trustStore.isBlank());
}
static void validateUrl(URI url, boolean requireHttps) {
if (url == null || !("http".equals(url.getScheme()) || "https".equals(url.getScheme())) ||
url.getHost() == null ||
url.getPort() < 1 || url.getRawUserInfo() != null || url.getPort() < 1 || url.getRawUserInfo() != null ||
(url.getRawPath() != null && !url.getRawPath().isEmpty()) || (url.getRawPath() != null && !url.getRawPath().isEmpty()) ||
url.getRawQuery() != null || url.getRawFragment() != null) url.getRawQuery() != null || url.getRawFragment() != null)
throw new IllegalArgumentException("Invalid private storage node URL"); throw new IllegalArgumentException("Invalid private storage node URL");
if (requireHttps && !"https".equals(url.getScheme()))
throw new IllegalArgumentException("Cluster TLS truststore requires HTTPS node URLs");
} }
boolean availableHostsAtLeast(int required, boolean testNodeDomains) { boolean availableHostsAtLeast(int required, boolean testNodeDomains) {
Set<UUID> healthy = new HashSet<>(); Set<UUID> healthy = new HashSet<>();
for (int i = 0; i < nodes.size(); i++) { for (int i = 0; i < nodes.size(); i++) {
Node node = nodes.get(i); Node node = nodes.get(i);
NodeIdentity actual = probeIfAvailable(node.url(), token); NodeIdentity actual = probeIfAvailable(node.url(), token, http);
if (actual == null || !actual.nodeId().equals(node.id()) || !actual.hostId().equals(node.hostId())) { if (actual == null || !actual.nodeId().equals(node.id()) || !actual.hostId().equals(node.hostId())) {
markUnreadable(node); markUnreadable(node);
continue; continue;
@@ -159,7 +186,7 @@ final class NodeClient {
Node node = nodes.get(index); Node node = nodes.get(index);
if (unreadable(node)) throw new IOException("Storage node is temporarily unreachable"); if (unreadable(node)) throw new IOException("Storage node is temporarily unreachable");
if (!recentlyHealthy(node)) { if (!recentlyHealthy(node)) {
NodeIdentity actual = probeIfAvailable(node.url(), token); NodeIdentity actual = probeIfAvailable(node.url(), token, http);
if (actual == null || !actual.nodeId().equals(node.id()) || !actual.hostId().equals(node.hostId())) { if (actual == null || !actual.nodeId().equals(node.id()) || !actual.hostId().equals(node.hostId())) {
markUnreadable(node); markUnreadable(node);
throw new IOException("Storage node is temporarily unreachable"); throw new IOException("Storage node is temporarily unreachable");
+40 -24
View File
@@ -74,10 +74,36 @@ final class SigV4 {
private Verified verifyPresigned(String method, URI uri, Headers headers) { private Verified verifyPresigned(String method, URI uri, Headers headers) {
if (headers.containsKey("authorization")) denied("Use one authentication method"); if (headers.containsKey("authorization")) denied("Use one authentication method");
PresignedQuery query = presignedQuery(uri.getRawQuery());
Map<String, String> fields = query.fields();
if (!fields.keySet().equals(Set.of("X-Amz-Algorithm", "X-Amz-Credential", "X-Amz-Date",
"X-Amz-Expires", "X-Amz-SignedHeaders", "X-Amz-Signature")) ||
!"AWS4-HMAC-SHA256".equals(fields.get("X-Amz-Algorithm")))
denied("Invalid presigned parameters");
String[] credential = credentialScope(fields.get("X-Amz-Credential"));
String date = fields.get("X-Amz-Date");
validatePresignedTime(date, credential[1], fields.get("X-Amz-Expires"));
String signedHeaders = fields.get("X-Amz-SignedHeaders");
String canonicalHeaders = canonicalHeaders(headers, signedHeaders, Set.of("host"));
String scope = String.join("/", Arrays.copyOfRange(credential, 1, 5));
String canonical = method + "\n" + encode(decode(uri.getRawPath()), true) + "\n"
+ canonicalQuery(query.signed()) + "\n" + canonicalHeaders + "\n"
+ signedHeaders + "\nUNSIGNED-PAYLOAD";
String toSign = "AWS4-HMAC-SHA256\n" + date + "\n" + scope + "\n"
+ hex(hash(canonical.getBytes(StandardCharsets.UTF_8)));
String signature = fields.get("X-Amz-Signature");
byte[] key = signingKey(secret(credential[0]), credential[1], region);
if (!HEX.matcher(signature).matches() ||
!MessageDigest.isEqual(hmac(key, toSign), HexFormat.of().parseHex(signature)))
denied("Signature mismatch");
return new Verified("UNSIGNED-PAYLOAD", query.application(), key, date, scope, signature, credential[0]);
}
private static PresignedQuery presignedQuery(String rawQuery) {
Map<String, String> fields = new TreeMap<>(); Map<String, String> fields = new TreeMap<>();
StringBuilder application = new StringBuilder(); StringBuilder application = new StringBuilder();
StringBuilder signed = new StringBuilder(); StringBuilder signed = new StringBuilder();
for (String part : uri.getRawQuery().split("&", -1)) { for (String part : rawQuery.split("&", -1)) {
String[] pair = part.split("=", 2); String[] pair = part.split("=", 2);
String name = decode(pair[0]); String name = decode(pair[0]);
String value = decode(pair.length == 2 ? pair[1] : ""); String value = decode(pair.length == 2 ? pair[1] : "");
@@ -89,17 +115,19 @@ final class SigV4 {
appendQuery(signed, part); appendQuery(signed, part);
} }
} }
if (!fields.keySet().equals(Set.of("X-Amz-Algorithm", "X-Amz-Credential", "X-Amz-Date", return new PresignedQuery(fields, application.toString(), signed.toString());
"X-Amz-Expires", "X-Amz-SignedHeaders", "X-Amz-Signature")) || }
!"AWS4-HMAC-SHA256".equals(fields.get("X-Amz-Algorithm")))
denied("Invalid presigned parameters"); private void validatePresignedTime(String date, String scopeDate, String rawExpires) {
String[] credential = credentialScope(fields.get("X-Amz-Credential")); if (!date.matches("[0-9]{8}T[0-9]{6}Z") || !date.startsWith(scopeDate))
String date = fields.get("X-Amz-Date");
if (!date.matches("[0-9]{8}T[0-9]{6}Z") || !date.startsWith(credential[1]))
denied("Invalid signing date"); denied("Invalid signing date");
long expires; long expires;
try { expires = Long.parseLong(fields.get("X-Amz-Expires")); } try {
catch (NumberFormatException error) { denied("Invalid presigned expiry"); return null; } expires = Long.parseLong(rawExpires);
} catch (NumberFormatException error) {
denied("Invalid presigned expiry");
return;
}
if (expires < 1 || expires > 604800) denied("Invalid presigned expiry"); if (expires < 1 || expires > 604800) denied("Invalid presigned expiry");
try { try {
Instant start = Instant.from(DATE.parse(date)); Instant start = Instant.from(DATE.parse(date));
@@ -107,22 +135,10 @@ final class SigV4 {
if (now.isBefore(start.minus(Duration.ofMinutes(5))) || now.isAfter(start.plusSeconds(expires))) if (now.isBefore(start.minus(Duration.ofMinutes(5))) || now.isAfter(start.plusSeconds(expires)))
denied("Presigned URL has expired or is not yet valid"); denied("Presigned URL has expired or is not yet valid");
} catch (java.time.DateTimeException error) { denied("Invalid signing date"); } } catch (java.time.DateTimeException error) { denied("Invalid signing date"); }
String signedHeaders = fields.get("X-Amz-SignedHeaders");
String canonicalHeaders = canonicalHeaders(headers, signedHeaders, Set.of("host"));
String scope = String.join("/", Arrays.copyOfRange(credential, 1, 5));
String canonical = method + "\n" + encode(decode(uri.getRawPath()), true) + "\n"
+ canonicalQuery(signed.toString()) + "\n" + canonicalHeaders + "\n"
+ signedHeaders + "\nUNSIGNED-PAYLOAD";
String toSign = "AWS4-HMAC-SHA256\n" + date + "\n" + scope + "\n"
+ hex(hash(canonical.getBytes(StandardCharsets.UTF_8)));
String signature = fields.get("X-Amz-Signature");
byte[] key = signingKey(secret(credential[0]), credential[1], region);
if (!HEX.matcher(signature).matches() ||
!MessageDigest.isEqual(hmac(key, toSign), HexFormat.of().parseHex(signature)))
denied("Signature mismatch");
return new Verified("UNSIGNED-PAYLOAD", application.toString(), key, date, scope, signature, credential[0]);
} }
private record PresignedQuery(Map<String, String> fields, String application, String signed) { }
private static boolean hasPresignedQuery(String raw) { private static boolean hasPresignedQuery(String raw) {
return raw != null && (raw.startsWith("X-Amz-Algorithm=") || raw.contains("&X-Amz-Algorithm=")); return raw != null && (raw.startsWith("X-Amz-Algorithm=") || raw.contains("&X-Amz-Algorithm="));
} }
@@ -0,0 +1,19 @@
package cloud.lunarsky.store;
import java.util.Map;
final class StorageLimits {
private StorageLimits() {}
static ObjectStorage.Limits fromEnvironment(Map<String, String> 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);
}
}
}
@@ -0,0 +1,63 @@
package cloud.lunarsky.store;
import com.sun.net.httpserver.HttpServer;
import java.net.InetSocketAddress;
import java.net.URI;
import java.nio.file.Files;
import java.nio.file.Path;
import java.sql.DriverManager;
import java.util.List;
import java.util.Map;
import java.util.UUID;
public final class ClusterMetadataTest {
public static void main(String[] args) throws Exception {
Map<String, String> env = System.getenv();
String token = "metadata-test-cluster-token-0123456789";
String repairToken = "metadata-test-repair-token-0123456789";
Path root = Files.createTempDirectory("objectstore-metadata-test-");
try (ClusterNode first = new ClusterNode(root.resolve("first"), token, repairToken, UUID.randomUUID());
ClusterNode second = new ClusterNode(root.resolve("second"), token, repairToken, UUID.randomUUID())) {
HttpServer firstServer = server(first);
HttpServer secondServer = server(second);
try {
List<URI> nodes = List.of(url(firstServer), url(secondServer));
try (ClusterStore store = new ClusterStore(env.get("POSTGRES_JDBC_URL"),
env.get("POSTGRES_USER"), env.get("POSTGRES_PASSWORD"), "objects",
nodes, token, repairToken, 1024 * 1024, 16 * 1024 * 1024, false)) {
require(store.ready(), "Writable primary was not selected from the multi-host URL");
try (var admin = DriverManager.getConnection(env.get("POSTGRES_ADMIN_JDBC_URL"),
env.get("POSTGRES_USER"), env.get("POSTGRES_PASSWORD"));
var statement = admin.createStatement()) {
statement.execute("ALTER DATABASE objectstore_test SET default_transaction_read_only=on");
try {
require(!store.ready(), "Read-only metadata was reported ready");
} finally {
statement.execute("ALTER DATABASE objectstore_test RESET default_transaction_read_only");
}
}
require(store.ready(), "Writable metadata did not recover after read-only mode ended");
}
} finally {
firstServer.stop(0);
secondServer.stop(0);
}
}
System.out.println("Metadata routing tests passed: second JDBC host and writable readiness");
}
private static HttpServer server(ClusterNode node) throws Exception {
HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
server.createContext("/", node::handle);
server.start();
return server;
}
private static URI url(HttpServer server) {
return URI.create("http://127.0.0.1:" + server.getAddress().getPort());
}
private static void require(boolean condition, String message) {
if (!condition) throw new AssertionError(message);
}
}
@@ -0,0 +1,107 @@
package cloud.lunarsky.store;
import com.sun.net.httpserver.HttpServer;
import com.sun.net.httpserver.HttpsServer;
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.Map;
import java.util.UUID;
public final class ClusterTlsTest {
private static final String PASSWORD = "local-test-password-0123456789";
public static void main(String[] args) throws Exception {
Path directory = Files.createTempDirectory("objectstore-tls-");
Path keyStore = directory.resolve("node.p12");
Path trustStore = directory.resolve("trust.p12");
Path certificate = directory.resolve("node.crt");
Path passwordFile = directory.resolve("password");
Files.writeString(passwordFile, PASSWORD + "\n");
String keytool = Path.of(System.getProperty("java.home"), "bin", "keytool").toString();
run(keytool, "-genkeypair", "-alias", "node", "-keyalg", "RSA", "-keysize", "2048",
"-validity", "2", "-dname", "CN=localhost", "-ext", "SAN=DNS:localhost",
"-storetype", "PKCS12", "-keystore", keyStore.toString(), "-storepass", PASSWORD,
"-keypass", PASSWORD, "-noprompt");
run(keytool, "-exportcert", "-alias", "node", "-keystore", keyStore.toString(),
"-storepass", PASSWORD, "-file", certificate.toString());
run(keytool, "-importcert", "-alias", "node", "-file", certificate.toString(),
"-keystore", trustStore.toString(), "-storetype", "PKCS12", "-storepass", PASSWORD,
"-noprompt");
Map<String, String> serverConfig = Map.of(
"NODE_TLS_KEYSTORE", keyStore.toString(), "NODE_TLS_PASSWORD_FILE", passwordFile.toString());
Map<String, String> clientConfig = Map.of(
"CLUSTER_TLS_TRUSTSTORE", trustStore.toString(),
"CLUSTER_TLS_PASSWORD_FILE", passwordFile.toString());
String token = "tls-test-cluster-token-0123456789";
String repairToken = "tls-test-repair-token-0123456789";
UUID hostId = UUID.randomUUID();
byte[] data = "encrypted transport".getBytes(java.nio.charset.StandardCharsets.UTF_8);
try (ClusterNode node = new ClusterNode(directory.resolve("data"), token, repairToken, hostId)) {
HttpServer server = ClusterTls.nodeServer(new InetSocketAddress("127.0.0.1", 0), serverConfig);
require(server instanceof HttpsServer, "Node did not enable HTTPS");
server.createContext("/", node::handle);
server.start();
try {
URI url = URI.create("https://localhost:" + server.getAddress().getPort());
HttpClient trusted = ClusterTls.client(clientConfig, Duration.ofSeconds(3));
NodeIdentity identity = NodeClient.probe(url, token, trusted);
require(identity.hostId().equals(hostId), "TLS probe returned wrong node identity");
NodeClient client = new NodeClient(List.of(
new NodeClient.Node(identity.nodeId(), hostId, url)), token, repairToken, trusted);
UUID segment = UUID.randomUUID();
client.put(0, segment, data, SigV4.hash(data));
require(java.util.Arrays.equals(data, client.get(0, segment, data.length, SigV4.hash(data))),
"TLS segment roundtrip failed");
HttpRequest request = HttpRequest.newBuilder(url.resolve("/identity"))
.header("X-Cluster-Token", token).GET().build();
try {
ClusterTls.client(Map.of(), Duration.ofSeconds(3))
.send(request, HttpResponse.BodyHandlers.discarding());
throw new AssertionError("Untrusted certificate was accepted");
} catch (IOException expected) { }
URI wrongHost = URI.create("https://127.0.0.1:" + server.getAddress().getPort());
try {
NodeClient.probe(wrongHost, token, trusted);
throw new AssertionError("Wrong certificate hostname was accepted");
} catch (IOException expected) { }
} finally {
server.stop(0);
}
}
try {
ClusterTls.nodeServer(new InetSocketAddress("127.0.0.1", 0),
Map.of("NODE_TLS_KEYSTORE", keyStore.toString()));
throw new AssertionError("Incomplete TLS configuration was accepted");
} catch (IOException expected) { }
try {
ClusterTls.client(Map.of("CLUSTER_TLS_TRUSTSTORE", trustStore.toString()),
Duration.ofSeconds(3));
throw new AssertionError("Incomplete cluster trust configuration was accepted");
} catch (IOException expected) { }
try {
NodeClient.validateUrl(URI.create("http://localhost:9100"), true);
throw new AssertionError("HTTP node URL was accepted with cluster TLS enabled");
} catch (IllegalArgumentException expected) { }
System.out.println("Cluster TLS tests passed: trusted roundtrip, untrusted and hostname rejection, no HTTP downgrade");
}
private static void run(String... command) throws Exception {
Process process = new ProcessBuilder(command).redirectErrorStream(true).start();
String output = new String(process.getInputStream().readAllBytes(),
java.nio.charset.StandardCharsets.UTF_8);
if (process.waitFor() != 0) throw new AssertionError("keytool failed: " + output);
}
private static void require(boolean condition, String message) {
if (!condition) throw new AssertionError(message);
}
}
@@ -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<String, String> 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<String, String> 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");
}
}
+68 -2
View File
@@ -806,15 +806,70 @@ public final class HttpTest {
throw new AssertionError("Completed multipart version was not retained"); 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 { public static void main(String[] args) throws Exception {
Path root = Files.createTempDirectory("store-http-test-"); Path root = Files.createTempDirectory("store-http-test-");
var executor = Executors.newVirtualThreadPerTaskExecutor(); var executor = Executors.newVirtualThreadPerTaskExecutor();
HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 16); HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 16);
DiskStore store = new DiskStore(root, 1024, 4096); DiskStore store = new DiskStore(root, 1024, 4096);
try { try {
var app = new Main(store, var app = new Main(store, new MultipartStore(store),
new SigV4(Map.of(ACCESS, SECRET, SECONDARY, SECONDARY_SECRET), 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.setExecutor(executor);
server.createContext("/", app::handle); server.createContext("/", app::handle);
server.start(); server.start();
@@ -834,6 +889,17 @@ public final class HttpTest {
testBuckets(client, base); testBuckets(client, base);
testVersioning(client, base); testVersioning(client, base);
testAcl(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"); System.out.println("HTTP tests passed: capabilities, objects, copy, checksums, listing, multipart, attributes, buckets, versioning, ACLs");
} finally { } finally {
server.stop(0); server.stop(0);
+30
View File
@@ -0,0 +1,30 @@
services:
metadata:
image: postgres:17-alpine
environment:
POSTGRES_DB: objectstore_test
POSTGRES_USER: objectstore_test
POSTGRES_PASSWORD: local-metadata-test-only
tmpfs:
- /var/lib/postgresql/data:size=268435456
healthcheck:
test: ["CMD-SHELL", "pg_isready -U objectstore_test -d objectstore_test"]
interval: 2s
timeout: 2s
retries: 20
mem_limit: 256m
check:
build:
context: ../..
image: lunarsky-objectstore:metadata-test
entrypoint: ["java", "--add-modules", "jdk.httpserver,java.net.http", "-cp", "/app:/app/postgresql.jar:/app/hash4j.jar", "cloud.lunarsky.store.ClusterMetadataTest"]
environment:
POSTGRES_JDBC_URL: jdbc:postgresql://127.0.0.2:5432,metadata:5432/objectstore_test?connectTimeout=1&socketTimeout=10&targetServerType=primary&hostRecheckSeconds=0
POSTGRES_ADMIN_JDBC_URL: jdbc:postgresql://metadata:5432/postgres?connectTimeout=3&socketTimeout=10
POSTGRES_USER: objectstore_test
POSTGRES_PASSWORD: local-metadata-test-only
depends_on:
metadata:
condition: service_healthy
mem_limit: 384m
+1 -1
View File
@@ -1,6 +1,6 @@
# Two-machine durability drill # Two-machine durability drill
Run this disposable test on two machines. Machine A runs the gateway, PostgreSQL, and one storage node; machine B runs a second storage node. Keep the node connection on a private network: the node protocol uses bearer tokens over HTTP. Run this disposable test on two machines. Machine A runs the gateway, PostgreSQL, and one storage node; machine B runs a second storage node. Keep the node connection on a private network: this drill uses bearer tokens over HTTP and does not enable the optional node TLS setup.
Both machines need Docker. Machine A also needs Docker Compose and Python 3. The example uses loopback port 9003 on machine A and private-network port 9103 on machine B; change them if needed. Use separate test volumes, and do not point this drill at an existing ObjectStore cluster. Both machines need Docker. Machine A also needs Docker Compose and Python 3. The example uses loopback port 9003 on machine A and private-network port 9103 on machine B; change them if needed. Use separate test volumes, and do not point this drill at an existing ObjectStore cluster.