Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
510e519bd1 | ||
|
|
43987bbadb | ||
|
|
a71c4fa5f1 | ||
|
|
1d1f329220 | ||
|
|
91f6825726 | ||
|
|
885239be1c |
No files matched your search
+11
-3
@@ -1,7 +1,15 @@
|
||||
S3_ACCESS_KEY=
|
||||
S3_SECRET_KEY=
|
||||
S3_ACCESS_KEY=[Replace_with_your_S3_Access_Key]
|
||||
S3_SECRET_KEY=[Replace_with_your_S3_Secret_Key]
|
||||
CLUSTER_TOKEN=
|
||||
CLUSTER_REPAIR_TOKEN=
|
||||
POSTGRES_PASSWORD=
|
||||
POSTGRES_PASSWORD=[PG_Pass]
|
||||
S3_BUCKET=objects
|
||||
CLUSTER_HOST_PORT=9001
|
||||
MAX_OBJECT_BYTES=134217728
|
||||
MAX_TOTAL_BYTES=2147483648
|
||||
MAX_IN_FLIGHT_REQUESTS=16
|
||||
HTTP_BACKLOG=64
|
||||
S3_VIRTUAL_HOST_SUFFIX=
|
||||
# For compose.cluster.encrypted.yaml only; provision and mount LUKS before setting these.
|
||||
ENCRYPTED_STORAGE_ROOT=
|
||||
ENCRYPTED_VOLUME_ID=
|
||||
@@ -13,3 +13,9 @@ PUBLIC_BYTES_PER_SECOND=0
|
||||
PUBLIC_BYTE_BURST=
|
||||
PUBLIC_MAX_IN_FLIGHT_PER_IP=
|
||||
PUBLIC_TRUSTED_PROXY_IPS=
|
||||
MAX_IN_FLIGHT_REQUESTS=16
|
||||
HTTP_BACKLOG=64
|
||||
S3_VIRTUAL_HOST_SUFFIX=
|
||||
# For compose.encrypted.yaml only; provision and mount LUKS before setting these.
|
||||
ENCRYPTED_STORAGE_ROOT=
|
||||
ENCRYPTED_VOLUME_ID=
|
||||
@@ -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,java.net.http -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.HttpTest
|
||||
RUN java --add-modules jdk.httpserver,java.net.http -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.ClientLimitsTest
|
||||
RUN java -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.EncryptedVolumeTest
|
||||
RUN java --add-modules jdk.httpserver,java.net.http -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.ClusterNodeTest
|
||||
RUN java --add-modules jdk.httpserver,java.net.http -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.ClusterTlsTest
|
||||
RUN java -cp /out:/tmp/hash4j.jar cloud.lunarsky.store.CliTest
|
||||
|
||||
FROM eclipse-temurin:21-jre-alpine
|
||||
|
||||
@@ -16,11 +16,14 @@ Source: [GitHub](https://github.com/LunarSkyOSS/ObjectStore) · [Gitea mirror](h
|
||||
- [S3 API support checklist](#s3-api-support-checklist)
|
||||
- [Tests](TESTS.md)
|
||||
- [Single-node setup](#single-node-setup)
|
||||
- [Encrypted storage](#encrypted-storage)
|
||||
- [Access keys and ACLs](#access-keys-and-acls)
|
||||
- [CLI and tests](#cli-and-tests)
|
||||
- [Capability discovery](#capability-discovery)
|
||||
- [Java client](#java-client)
|
||||
- [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)
|
||||
- [Adding a cluster node](#adding-a-cluster-node)
|
||||
- [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.
|
||||
|
||||
- ✅ Persistent single-node storage with checksum verification
|
||||
- ✅ Optional dm-crypt-backed data mounts with startup guards
|
||||
- ✅ Configurable per-object and total logical size limits
|
||||
- ✅ CLI status, version, and full payload verification
|
||||
- ✅ Local cluster prototype with stable node IDs and host-aware placement code
|
||||
- ✅ Optional HTTPS between cluster processes and storage nodes
|
||||
- ✅ Opt-in automatic repair, rebalance, and guarded garbage collection in the local cluster
|
||||
- ✅ Metadata backup and tested restore to a separate local PostgreSQL instance
|
||||
- ✅ 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
|
||||
|
||||
- ✅ Header-based and presigned-query AWS Signature Version 4 authentication
|
||||
- ✅ Path-style and configurable virtual-hosted-style bucket addresses
|
||||
- ✅ `PutObject`, `GetObject`, `HeadObject`, and `DeleteObject` in both modes
|
||||
- ✅ Single-range GET and `ListObjectsV2` in both modes
|
||||
- ✅ SHA-256 payload verification and `x-amz-checksum-sha256` in both modes
|
||||
@@ -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.
|
||||
|
||||
## Encrypted storage
|
||||
|
||||
At-rest encryption is an **opt-in deployment setting** backed by a host-mounted dm-crypt/LUKS filesystem. ObjectStore does not encrypt bytes itself or store an encryption key in its environment. Prepare and unlock the LUKS filesystem outside Docker, then set `ENCRYPTED_STORAGE_ROOT` to its exact mount point and `ENCRYPTED_VOLUME_ID` to a stable 32-character lowercase hex identifier in your private environment file. Generate the identifier with `openssl rand -hex 16`; it is a mount guard, **not** a cryptographic key. Keep the LUKS recovery key separately from the data and backups.
|
||||
|
||||
With the encrypted filesystem mounted, prepare empty directories using `sh scripts/prepare-encrypted-storage.sh single /path/to/objectstore.env` or `sh scripts/prepare-encrypted-storage.sh cluster /path/to/cluster.env`. The script reads only the two encryption settings from that file; it does not execute it. It requires an active dm-crypt-backed mount at `ENCRYPTED_STORAGE_ROOT` and refuses to mark a nonempty directory. Run it with permission to create the directories and set container ownership. Enable the matching Compose override:
|
||||
|
||||
```sh
|
||||
docker compose --env-file /path/to/objectstore.env -f compose.yaml -f compose.encrypted.yaml up -d --build
|
||||
docker compose --env-file /path/to/cluster.env -f compose.cluster.yaml -f compose.cluster.encrypted.yaml up -d --build
|
||||
```
|
||||
|
||||
Use the first command for single-node mode or the second for the **local development cluster**, not both. The overrides bind encrypted directories for object data, cluster nodes, PostgreSQL, and cluster staging files. Each process checks a marker stored on its mounted directory before opening data; PostgreSQL checks before starting. If a mount is missing after reboot, the service refuses to start instead of silently creating a new plaintext store. Provision an unlock-and-mount service before starting Compose; a container restart cannot unlock LUKS. The preparation script checks the host mount, while the runtime marker is only a guard against an absent or wrong mount.
|
||||
|
||||
**Existing Docker volumes are not migrated automatically.** Stop the stack, take and verify a backup, prepare the empty encrypted directories, then copy the stopped volumes into their corresponding directories before starting with the override. For PostgreSQL, copy the existing data directory into `cluster/metadata/pgdata` because the encrypted override changes `PGDATA`. Verify the restored objects and database before retiring the original volumes. The cluster overlay places all local test containers on one encrypted host; a future multi-host deployment needs an encrypted data mount and key recovery plan on each host.
|
||||
|
||||
This covers the specified data and staging mounts, but not Docker logs, host swap, other temporary files, external backups, or a compromised running host. Encrypt those separately as appropriate, use TLS for network paths, restrict access, and test backup recovery. Encryption at rest is one security measure; enabling it does not by itself establish GDPR compliance.
|
||||
|
||||
## Access keys and ACLs
|
||||
|
||||
`S3_ACCESS_KEY` is the owner identity. Its secret is `S3_SECRET_KEY`. Additional keys are optional: place one `ACCESSKEY:secret` pair per line in a file mounted read-only inside the container, and set `S3_CREDENTIALS_FILE=/run/secrets/s3-credentials` in the environment file. Use a Compose override to mount an absolute host path at `/run/secrets/s3-credentials` for the `objectstore` service (or `gateway` in cluster mode). Each access key needs 16–128 alphanumeric characters and each secret at least 32 characters. A restart loads changes to that file. Keep it outside Git and protect it as a secret. Additional identities have no access until the owner grants it.
|
||||
@@ -117,7 +140,7 @@ The [JDK-only Java client](client/README.md) works with ObjectStore and other S3
|
||||
|
||||
## 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
|
||||
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 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
|
||||
|
||||
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
|
||||
|
||||
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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
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.
|
||||
|
||||
|
||||
@@ -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. |
|
||||
| `ConcurrencyTest` | Atomic local overwrites and consistent reads, listings, and deletes during concurrent access. |
|
||||
| `HttpTest` | Signed capability discovery, presigned URLs, streaming uploads and trailers, object and bucket operations, ACL grants with a second access key and public reads, copies, checksum persistence and rejection, ranges, listing, metadata, tags, multipart uploads, and versioning in single-node mode. |
|
||||
| `HttpTest` | Signed capability discovery, presigned URLs, streaming uploads and trailers, object and bucket operations, ACL grants with a second access key and public reads, copies, checksum persistence and rejection, ranges, listing, metadata, tags, multipart uploads, versioning, and signed virtual-hosted requests in single-node mode. |
|
||||
| `ClientLimitsTest` | Disabled defaults, trusted-proxy address validation, ignored untrusted headers, per-IP request refusal, and paced response bytes. |
|
||||
| `EncryptedVolumeTest` | Marker acceptance and refusal when the configured encrypted directory is missing, mismatched, or a symlink. |
|
||||
| `ClusterNodeTest` | Node identity and locking, authenticated segment transfers, checksum rejection, repair authorization, inventory and guarded deletion, and restart cleanup. |
|
||||
| `ClusterTlsTest` | HTTPS node identity and segment roundtrip with a trusted certificate, plus rejection of untrusted and wrong-host certificates. |
|
||||
| `CliTest` | Version, status, verification, and a nonzero result for corrupt data. |
|
||||
| `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
|
||||
|
||||
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
|
||||
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.
|
||||
|
||||
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).
|
||||
|
||||
|
||||
@@ -371,7 +371,9 @@ public final class ImageManager extends Frame {
|
||||
Object value = event.getTransferable().getTransferData(DataFlavor.javaFileListFlavor);
|
||||
List<?> dropped = (List<?>) value;
|
||||
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);
|
||||
EventQueue.invokeLater(() -> uploadImages(files));
|
||||
} catch (Exception error) {
|
||||
@@ -391,7 +393,10 @@ public final class ImageManager extends Frame {
|
||||
Button cancel = new Button("Cancel");
|
||||
Button proceed = new Button("Continue");
|
||||
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(proceed);
|
||||
dialog.add(buttons, BorderLayout.SOUTH);
|
||||
|
||||
@@ -124,7 +124,10 @@ final class Json {
|
||||
return result.toString();
|
||||
}
|
||||
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");
|
||||
char escaped = source.charAt(index++);
|
||||
switch (escaped) {
|
||||
@@ -159,31 +162,42 @@ final class Json {
|
||||
private BigDecimal number() throws ProtocolException {
|
||||
int start = index;
|
||||
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 (index < source.length() && Character.isDigit(source.charAt(index)))
|
||||
throw new ProtocolException("Invalid JSON number");
|
||||
} else {
|
||||
if (index >= source.length() || source.charAt(index) < '1' || source.charAt(index) > '9')
|
||||
throw new ProtocolException("Invalid JSON number");
|
||||
while (index < source.length() && source.charAt(index) >= '0' && source.charAt(index) <= '9') index++;
|
||||
scanDigits();
|
||||
}
|
||||
if (take('.')) {
|
||||
}
|
||||
|
||||
private void requireDigits(String message) throws ProtocolException {
|
||||
int first = index;
|
||||
while (index < source.length() && source.charAt(index) >= '0' && source.charAt(index) <= '9') index++;
|
||||
if (first == index) throw new ProtocolException("Invalid JSON fraction");
|
||||
scanDigits();
|
||||
if (first == index) throw new ProtocolException(message);
|
||||
}
|
||||
if (take('e') || take('E')) {
|
||||
if (!take('+')) take('-');
|
||||
int first = index;
|
||||
|
||||
private void scanDigits() {
|
||||
while (index < source.length() && source.charAt(index) >= '0' && source.charAt(index) <= '9') index++;
|
||||
if (first == index) throw new ProtocolException("Invalid JSON exponent");
|
||||
}
|
||||
try { return new BigDecimal(source.substring(start, index)); }
|
||||
catch (NumberFormatException e) { throw new ProtocolException("Invalid JSON number", e); }
|
||||
}
|
||||
|
||||
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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
@@ -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
@@ -15,6 +15,9 @@ services:
|
||||
S3_CREDENTIALS_FILE: ${S3_CREDENTIALS_FILE:-}
|
||||
MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728}
|
||||
MAX_TOTAL_BYTES: ${MAX_TOTAL_BYTES:-2147483648}
|
||||
MAX_IN_FLIGHT_REQUESTS: ${MAX_IN_FLIGHT_REQUESTS:-16}
|
||||
HTTP_BACKLOG: ${HTTP_BACKLOG:-64}
|
||||
S3_VIRTUAL_HOST_SUFFIX: ${S3_VIRTUAL_HOST_SUFFIX:-}
|
||||
PUBLIC_REQUESTS_PER_SECOND: ${PUBLIC_REQUESTS_PER_SECOND:-0}
|
||||
PUBLIC_REQUEST_BURST: ${PUBLIC_REQUEST_BURST:-}
|
||||
PUBLIC_BYTES_PER_SECOND: ${PUBLIC_BYTES_PER_SECOND:-0}
|
||||
@@ -23,7 +26,7 @@ services:
|
||||
PUBLIC_TRUSTED_PROXY_IPS: ${PUBLIC_TRUSTED_PROXY_IPS:-}
|
||||
CLUSTER_TOKEN: ${CLUSTER_TOKEN:?Set CLUSTER_TOKEN}
|
||||
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_PASSWORD: ${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD}
|
||||
ports:
|
||||
@@ -96,9 +99,11 @@ services:
|
||||
CLUSTER_TOKEN: ${CLUSTER_TOKEN:?Set CLUSTER_TOKEN}
|
||||
CLUSTER_REPAIR_TOKEN: ${CLUSTER_REPAIR_TOKEN:?Set CLUSTER_REPAIR_TOKEN}
|
||||
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_PASSWORD: ${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD}
|
||||
MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728}
|
||||
MAX_TOTAL_BYTES: ${MAX_TOTAL_BYTES:-2147483648}
|
||||
CLUSTER_GC_MIN_AGE_SECONDS: ${CLUSTER_GC_MIN_AGE_SECONDS:-1209600}
|
||||
CLUSTER_BACKUP_RETENTION_SECONDS: ${CLUSTER_BACKUP_RETENTION_SECONDS:-0}
|
||||
CLUSTER_GC_TEST_MODE: ${CLUSTER_GC_TEST_MODE:-false}
|
||||
@@ -131,9 +136,11 @@ services:
|
||||
CLUSTER_BACKUP_RETENTION_SECONDS: ${CLUSTER_BACKUP_RETENTION_SECONDS:-0}
|
||||
CLUSTER_GC_TEST_MODE: ${CLUSTER_GC_TEST_MODE:-false}
|
||||
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_PASSWORD: ${POSTGRES_PASSWORD:?Set POSTGRES_PASSWORD}
|
||||
MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728}
|
||||
MAX_TOTAL_BYTES: ${MAX_TOTAL_BYTES:-2147483648}
|
||||
depends_on:
|
||||
metadata:
|
||||
condition: service_healthy
|
||||
|
||||
@@ -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
|
||||
@@ -11,6 +11,9 @@ services:
|
||||
S3_CREDENTIALS_FILE: ${S3_CREDENTIALS_FILE:-}
|
||||
MAX_OBJECT_BYTES: ${MAX_OBJECT_BYTES:-134217728}
|
||||
MAX_TOTAL_BYTES: ${MAX_TOTAL_BYTES:-2147483648}
|
||||
MAX_IN_FLIGHT_REQUESTS: ${MAX_IN_FLIGHT_REQUESTS:-16}
|
||||
HTTP_BACKLOG: ${HTTP_BACKLOG:-64}
|
||||
S3_VIRTUAL_HOST_SUFFIX: ${S3_VIRTUAL_HOST_SUFFIX:-}
|
||||
PUBLIC_REQUESTS_PER_SECOND: ${PUBLIC_REQUESTS_PER_SECOND:-0}
|
||||
PUBLIC_REQUEST_BURST: ${PUBLIC_REQUEST_BURST:-}
|
||||
PUBLIC_BYTES_PER_SECOND: ${PUBLIC_BYTES_PER_SECOND:-0}
|
||||
|
||||
@@ -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"
|
||||
@@ -22,12 +22,31 @@ port = int(values.get("CLUSTER_HOST_PORT", "9001"))
|
||||
if not 1 <= port <= 65535:
|
||||
raise ValueError("CLUSTER_HOST_PORT must be between 1 and 65535")
|
||||
host = f"127.0.0.1:{port}"
|
||||
MAX_RESPONSE_BYTES = 1024 * 1024
|
||||
|
||||
|
||||
def sign(key, message):
|
||||
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):
|
||||
extra = extra or {}
|
||||
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:
|
||||
connection.request(method, path, body=body if method in ("PUT", "POST") else None, headers=headers)
|
||||
response = connection.getresponse()
|
||||
return response.status, response.read(), response.headers
|
||||
return response.status, read_response(response), response.headers
|
||||
finally:
|
||||
connection.close()
|
||||
|
||||
@@ -60,7 +79,7 @@ def anonymous(method, path):
|
||||
try:
|
||||
connection.request(method, path)
|
||||
response = connection.getresponse()
|
||||
return response.status, response.read()
|
||||
return response.status, read_response(response)
|
||||
finally:
|
||||
connection.close()
|
||||
|
||||
@@ -75,7 +94,7 @@ if len(sys.argv) > 2 and sys.argv[2] == "acl":
|
||||
path = f"/{bucket}/cluster-test/acl-multipart"
|
||||
status, content, _ = request("POST", path + "?uploads", extra={"x-amz-acl": "public-read"})
|
||||
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")
|
||||
assert status == 200, status
|
||||
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":
|
||||
status, listing, _ = request("GET", "/version-bucket?versions")
|
||||
assert status == 200, status
|
||||
root = ET.fromstring(listing)
|
||||
root = parse_xml(listing)
|
||||
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)
|
||||
if version.findtext("s3:ETag", namespaces=namespace) == expected_etag]
|
||||
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"
|
||||
body = b"cluster copy and checksum test"
|
||||
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,
|
||||
{"content-type": "text/plain", "content-md5": md5,
|
||||
"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"
|
||||
status, content, _ = request("POST", versioned_multipart + "?uploads")
|
||||
assert status == 200, (status, content)
|
||||
versioned_upload = ET.fromstring(content).findtext("UploadId")
|
||||
versioned_upload = parse_xml(content).findtext("UploadId")
|
||||
assert versioned_upload, content
|
||||
versioned_part = b"retained multipart version"
|
||||
status, _, headers = request("PUT", versioned_multipart +
|
||||
|
||||
+10
-1
@@ -8,6 +8,8 @@ compose() { docker compose --env-file "$env_file" -f compose.cluster.yaml "$@";
|
||||
restore() {
|
||||
compose stop maintenance >/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
|
||||
}
|
||||
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")
|
||||
actual=$(compose exec -T node-a sha256sum "/data/segments/$shard/$segment_id" | cut -d' ' -f1)
|
||||
[ "$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
|
||||
status=$(curl -sS -o /dev/null -w '%{http_code}' "http://127.0.0.1:$host_port/ready")
|
||||
[ "$status" = 503 ]
|
||||
@@ -124,7 +133,7 @@ compose exec -T metadata-recovery pg_restore -U objectstore -d objectstore --no-
|
||||
< "$backup_dir/metadata.dump"
|
||||
compose stop metadata
|
||||
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 \
|
||||
-cp /app:/app/postgresql.jar:/app/hash4j.jar cloud.lunarsky.store.ClusterIntegrationTest recovered
|
||||
compose start metadata
|
||||
|
||||
@@ -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
|
||||
@@ -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,java.net.http -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.HttpTest
|
||||
java --add-modules jdk.httpserver,java.net.http -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.ClientLimitsTest
|
||||
java -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.EncryptedVolumeTest
|
||||
java --add-modules jdk.httpserver,java.net.http -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.ClusterNodeTest
|
||||
java --add-modules jdk.httpserver,java.net.http -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.ClusterTlsTest
|
||||
java -cp out/classes:lib/hash4j-0.30.0.jar cloud.lunarsky.store.CliTest
|
||||
bash client/scripts/test.sh
|
||||
@@ -105,25 +105,26 @@ final class AwsChunkedInputStream extends FilterInputStream {
|
||||
catch (NumberFormatException error) { throw invalid("Invalid signed chunk size"); }
|
||||
if (chunkLeft > decodedLength - decoded) throw invalid("Signed chunks exceed decoded length");
|
||||
chunkHash.reset();
|
||||
if (chunkLeft == 0) {
|
||||
if (chunkLeft == 0) finishPayload();
|
||||
}
|
||||
|
||||
private void finishPayload() throws IOException {
|
||||
finishChunk();
|
||||
if (decoded != decodedLength) throw invalid("Decoded length mismatch");
|
||||
if (trailerName == null) {
|
||||
if (!line().isEmpty()) throw invalid("Invalid signed chunk ending");
|
||||
} else {
|
||||
verifyTrailer();
|
||||
}
|
||||
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);
|
||||
byte[] actual;
|
||||
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))
|
||||
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}"))
|
||||
@@ -138,9 +139,17 @@ final class AwsChunkedInputStream extends FilterInputStream {
|
||||
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;
|
||||
|
||||
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 {
|
||||
|
||||
@@ -114,9 +114,21 @@ final class ClientLimits {
|
||||
if (exchange.getRemoteAddress().getAddress().isLoopbackAddress() &&
|
||||
exchange.getRequestHeaders().get("X-Real-IP") == null &&
|
||||
path.equals("/health")) return null;
|
||||
String address = address(exchange);
|
||||
Client client = admit(address(exchange));
|
||||
if (bytesPerSecond > 0) {
|
||||
try {
|
||||
exchange.setStreams(new LimitedInput(exchange.getRequestBody(), client),
|
||||
new LimitedOutput(exchange.getResponseBody(), client));
|
||||
} catch (RuntimeException error) {
|
||||
leave(client);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
return client;
|
||||
}
|
||||
|
||||
private synchronized Client admit(String address) {
|
||||
Client client;
|
||||
synchronized (this) {
|
||||
long now = System.nanoTime();
|
||||
if (++admissions % 1024 == 0 || clients.size() >= MAX_CLIENTS)
|
||||
clients.entrySet().removeIf(entry -> entry.getValue().inFlight == 0 &&
|
||||
@@ -136,16 +148,6 @@ final class ClientLimits {
|
||||
throw new StoreException(503, "SlowDown", "Client request rate exceeded");
|
||||
if (requestsPerSecond > 0) client.requestTokens--;
|
||||
client.inFlight++;
|
||||
}
|
||||
if (bytesPerSecond > 0) {
|
||||
try {
|
||||
exchange.setStreams(new LimitedInput(exchange.getRequestBody(), client),
|
||||
new LimitedOutput(exchange.getResponseBody(), client));
|
||||
} catch (RuntimeException error) {
|
||||
leave(client);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
return client;
|
||||
}
|
||||
|
||||
|
||||
@@ -15,6 +15,7 @@ public final class ClusterGc {
|
||||
if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV")))
|
||||
throw new IllegalArgumentException("Cluster garbage collection is only enabled in local cluster mode");
|
||||
boolean testDomains = "true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS"));
|
||||
ObjectStorage.Limits limits = StorageLimits.fromEnvironment(env);
|
||||
boolean disposableTest = testDomains && "true".equals(env.get("CLUSTER_GC_TEST_MODE"));
|
||||
long age = Long.parseLong(env.getOrDefault("CLUSTER_GC_MIN_AGE_SECONDS", "1209600"));
|
||||
long backupRetention = Long.parseLong(env.getOrDefault("CLUSTER_BACKUP_RETENTION_SECONDS", "0"));
|
||||
@@ -25,7 +26,8 @@ public final class ClusterGc {
|
||||
try (ClusterStore store = new ClusterStore(env.get("POSTGRES_JDBC_URL"), env.get("POSTGRES_USER"),
|
||||
env.get("POSTGRES_PASSWORD"), env.get("S3_BUCKET"),
|
||||
Arrays.stream(env.get("CLUSTER_NODES").split(",")).map(URI::create).toList(),
|
||||
env.get("CLUSTER_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"), 134217728, 2147483648L,
|
||||
env.get("CLUSTER_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"),
|
||||
limits.maxObjectBytes(), limits.maxTotalBytes(),
|
||||
testDomains)) {
|
||||
if (apply) {
|
||||
var repair = store.repairOnce();
|
||||
|
||||
@@ -12,6 +12,8 @@ public final class ClusterJoin {
|
||||
if (args.length != 2)
|
||||
throw new IllegalArgumentException("Usage: objectstore cluster-join node-url expected-host-uuid");
|
||||
Map<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")))
|
||||
throw new IllegalArgumentException("Node registration is only enabled in local cluster mode");
|
||||
URI url = URI.create(args[0]);
|
||||
|
||||
@@ -24,6 +24,8 @@ public final class ClusterMigrate {
|
||||
(args.length != 2 || !args[0].equals("--apply")))
|
||||
throw new IllegalArgumentException("Usage: objectstore cluster-migrate --check | --apply node-id-0,node-id-1,node-id-2");
|
||||
Map<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")))
|
||||
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();
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
package cloud.lunarsky.store;
|
||||
|
||||
import com.sun.net.httpserver.HttpExchange;
|
||||
import com.sun.net.httpserver.HttpServer;
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.io.OutputStream;
|
||||
@@ -33,6 +32,7 @@ public final class ClusterNode implements AutoCloseable {
|
||||
private final FileLock lock;
|
||||
|
||||
ClusterNode(Path root, String token, String repairToken, UUID hostId) throws IOException {
|
||||
EncryptedVolume.requireConfigured(System.getenv(), root);
|
||||
if (token == null || token.length() < 32) throw new IllegalArgumentException("Cluster token must have at least 32 characters");
|
||||
if (repairToken == null || repairToken.length() < 32 || repairToken.equals(token))
|
||||
throw new IllegalArgumentException("A separate repair token of at least 32 characters is required");
|
||||
@@ -325,7 +325,8 @@ public final class ClusterNode implements AutoCloseable {
|
||||
env.get("CLUSTER_REPAIR_TOKEN"),
|
||||
UUID.fromString(env.get("CLUSTER_HOST_ID")));
|
||||
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();
|
||||
server.setExecutor(executor);
|
||||
server.createContext("/", node::handle);
|
||||
@@ -336,6 +337,7 @@ public final class ClusterNode implements AutoCloseable {
|
||||
catch (IOException error) { System.err.println("Node close failed: " + error); }
|
||||
}));
|
||||
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)"));
|
||||
}
|
||||
}
|
||||
@@ -12,6 +12,7 @@ public final class ClusterRepair {
|
||||
if (!"cluster".equals(env.get("STORE_MODE")) || !"true".equals(env.get("CLUSTER_LOCAL_DEV")))
|
||||
throw new IllegalArgumentException("Cluster repair is only enabled in local cluster mode");
|
||||
boolean loop = args.length == 1;
|
||||
ObjectStorage.Limits limits = StorageLimits.fromEnvironment(env);
|
||||
long seconds = Long.parseLong(env.getOrDefault("CLUSTER_MAINTENANCE_INTERVAL_SECONDS", "60"));
|
||||
if (seconds < 1 || seconds > 3600) throw new IllegalArgumentException("Invalid maintenance interval");
|
||||
boolean gcEnabled = loop && "true".equals(env.get("CLUSTER_GC_ENABLED"));
|
||||
@@ -23,7 +24,8 @@ public final class ClusterRepair {
|
||||
try (ClusterStore store = new ClusterStore(env.get("POSTGRES_JDBC_URL"), env.get("POSTGRES_USER"),
|
||||
env.get("POSTGRES_PASSWORD"), env.get("S3_BUCKET"),
|
||||
Arrays.stream(env.get("CLUSTER_NODES").split(",")).map(URI::create).toList(),
|
||||
env.get("CLUSTER_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"), 134217728, 2147483648L,
|
||||
env.get("CLUSTER_TOKEN"), env.get("CLUSTER_REPAIR_TOKEN"),
|
||||
limits.maxObjectBytes(), limits.maxTotalBytes(),
|
||||
"true".equals(env.get("CLUSTER_TEST_NODE_DOMAINS")))) {
|
||||
var report = store.repairOnce();
|
||||
System.out.println("segments_scanned=" + report.scanned());
|
||||
|
||||
@@ -43,6 +43,8 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
|
||||
ClusterStore(String jdbcUrl, String user, String password, String bucket,
|
||||
List<URI> nodeUrls, String token, String repairToken, long maxObject, long maxTotal,
|
||||
boolean testNodeDomains) throws IOException {
|
||||
EncryptedVolume.requireConfigured(System.getenv(),
|
||||
java.nio.file.Path.of(System.getProperty("java.io.tmpdir")));
|
||||
if (jdbcUrl == null || !jdbcUrl.startsWith("jdbc:postgresql://") || user == null || password == null)
|
||||
throw new IllegalArgumentException("Invalid metadata database configuration");
|
||||
this.jdbcUrl = jdbcUrl;
|
||||
@@ -1347,8 +1349,9 @@ final class ClusterStore implements ObjectStorage, MultipartStorage {
|
||||
@Override public boolean ready() {
|
||||
if (!nodes.availableHostsAtLeast(2, testNodeDomains)) return false;
|
||||
try (Connection connection = connect(); var statement = connection.createStatement();
|
||||
ResultSet result = statement.executeQuery("SELECT 1")) {
|
||||
return result.next() && result.getInt(1) == 1;
|
||||
ResultSet result = statement.executeQuery(
|
||||
"SELECT NOT pg_is_in_recovery() AND current_setting('transaction_read_only') = 'off'")) {
|
||||
return result.next() && result.getBoolean(1);
|
||||
} catch (SQLException error) { return false; }
|
||||
}
|
||||
RepairReport repairOnce() throws IOException {
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
}
|
||||
@@ -48,6 +48,7 @@ final class DiskStore implements ObjectStorage {
|
||||
private record VersionRecord(String id, String storageId, boolean marker, long modified) {}
|
||||
|
||||
DiskStore(Path root, long maxObject, long maxTotal) throws IOException {
|
||||
EncryptedVolume.requireConfigured(System.getenv(), root);
|
||||
this.root = root;
|
||||
objects = root.resolve("objects");
|
||||
temporary = root.resolve("pending");
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
}
|
||||
@@ -24,6 +24,11 @@ import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.Semaphore;
|
||||
|
||||
public final class Main {
|
||||
|
||||
|
||||
//Todo: Main is too complex. Chop Chop.
|
||||
|
||||
|
||||
private static final String CAPABILITIES_PATH = "/_objectstore/capabilities";
|
||||
private final ObjectStorage store;
|
||||
private final SigV4 authentication;
|
||||
@@ -31,7 +36,8 @@ public final class Main {
|
||||
private final String bucket;
|
||||
private final MultipartStorage multipart;
|
||||
private final ClientLimits clientLimits;
|
||||
private final Semaphore slots = new Semaphore(16);
|
||||
private final Semaphore slots;
|
||||
private final String virtualHostSuffix;
|
||||
|
||||
Main(DiskStore store, SigV4 authentication, String bucket) throws IOException {
|
||||
this(store, new MultipartStore(store), authentication, bucket);
|
||||
@@ -43,15 +49,26 @@ public final class Main {
|
||||
|
||||
Main(ObjectStorage store, MultipartStorage multipart, SigV4 authentication, String bucket,
|
||||
ClientLimits clientLimits) throws IOException {
|
||||
this(store, multipart, authentication, bucket, clientLimits, 16, "");
|
||||
}
|
||||
|
||||
Main(ObjectStorage store, MultipartStorage multipart, SigV4 authentication, String bucket,
|
||||
ClientLimits clientLimits, int maxInFlight, String virtualHostSuffix) throws IOException {
|
||||
if (maxInFlight < 1 || maxInFlight > 1024)
|
||||
throw new IllegalArgumentException("MAX_IN_FLIGHT_REQUESTS must be between 1 and 1024");
|
||||
this.store = store;
|
||||
this.multipart = multipart;
|
||||
this.authentication = authentication;
|
||||
this.owner = authentication.root();
|
||||
this.bucket = bucket;
|
||||
this.clientLimits = clientLimits;
|
||||
this.slots = new Semaphore(maxInFlight);
|
||||
this.virtualHostSuffix = validateVirtualHostSuffix(virtualHostSuffix);
|
||||
store.ensureBucket(bucket);
|
||||
}
|
||||
|
||||
//Todo: These methods dont need to be in Main.java.
|
||||
|
||||
void handle(HttpExchange exchange) throws IOException {
|
||||
boolean admitted = false;
|
||||
ClientLimits.Client client = null;
|
||||
@@ -63,41 +80,85 @@ public final class Main {
|
||||
admitted = slots.tryAcquire();
|
||||
if (!admitted) throw new StoreException(503, "SlowDown", "Too many concurrent requests");
|
||||
if (handleStatus(exchange)) return;
|
||||
dispatch(exchange);
|
||||
} catch (StoreException error) {
|
||||
sendStoreError(exchange, error, requestId);
|
||||
} catch (Exception error) {
|
||||
System.err.println("ObjectStore request failed: " + requestId + " " + error.getClass().getSimpleName());
|
||||
sendError(exchange, 500, "InternalError", "Storage operation failed", requestId);
|
||||
} finally {
|
||||
try { exchange.close(); }
|
||||
finally {
|
||||
if (admitted) slots.release();
|
||||
clientLimits.leave(client);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
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 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);
|
||||
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, hash);
|
||||
requireEmptyBody(exchange, verified.payload());
|
||||
capabilities(exchange);
|
||||
} else if (path.equals("/")) {
|
||||
requireOwner(principal);
|
||||
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, hash);
|
||||
requireEmptyBody(exchange, verified.payload());
|
||||
listBuckets(exchange);
|
||||
} else {
|
||||
return;
|
||||
}
|
||||
int slash = path.indexOf('/', 1);
|
||||
String requestedBucket = slash < 0 ? path.substring(1) : path.substring(1, slash);
|
||||
String requestedBucket = hostBucket == null ?
|
||||
(slash < 0 ? path.substring(1) : path.substring(1, slash)) : hostBucket;
|
||||
if (requestedBucket.isEmpty()) throw new StoreException(404, "NoSuchBucket", "Bucket not found");
|
||||
if (slash < 0 || slash == path.length() - 1) {
|
||||
handleBucket(exchange, query, hash, requestedBucket, principal);
|
||||
if (hostBucket != null ? path.equals("/") : slash < 0 || slash == path.length() - 1) {
|
||||
handleBucket(exchange, query, verified.payload(), requestedBucket, verified.principal());
|
||||
} else {
|
||||
store.bucket(requestedBucket);
|
||||
handleObject(exchange, path.substring(slash + 1), query, verified, requestedBucket);
|
||||
handleObject(exchange, hostBucket == null ? path.substring(slash + 1) : path.substring(1),
|
||||
query, verified, requestedBucket);
|
||||
}
|
||||
}
|
||||
} catch (StoreException error) {
|
||||
|
||||
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)
|
||||
@@ -111,17 +172,6 @@ public final class Main {
|
||||
}
|
||||
sendError(exchange, error.status, error.code, error.getMessage(), requestId);
|
||||
}
|
||||
catch (Exception error) {
|
||||
System.err.println("ObjectStore request failed: " + requestId + " " + error.getClass().getSimpleName());
|
||||
sendError(exchange, 500, "InternalError", "Storage operation failed", requestId);
|
||||
} finally {
|
||||
try { exchange.close(); }
|
||||
finally {
|
||||
if (admitted) slots.release();
|
||||
clientLimits.leave(client);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static boolean anonymousRead(HttpExchange exchange) {
|
||||
if (!exchange.getRequestMethod().equals("GET") &&
|
||||
@@ -1183,10 +1233,8 @@ public final class Main {
|
||||
if (!access.matches("[A-Za-z0-9]{16,128}") || secret.length() < 32 ||
|
||||
!bucket.matches("[a-z0-9][a-z0-9-]{1,61}[a-z0-9]"))
|
||||
throw new IllegalArgumentException("Invalid storage credentials/bucket configuration");
|
||||
long maxObject = Long.parseLong(env.getOrDefault("MAX_OBJECT_BYTES", "134217728"));
|
||||
long maxTotal = Long.parseLong(env.getOrDefault("MAX_TOTAL_BYTES", "2147483648"));
|
||||
if (maxObject < 1 || maxObject > 1073741824L || maxTotal < maxObject)
|
||||
throw new IllegalArgumentException("Invalid size limits");
|
||||
ObjectStorage.Limits limits = StorageLimits.fromEnvironment(env);
|
||||
long maxObject = limits.maxObjectBytes(), maxTotal = limits.maxTotalBytes();
|
||||
String mode = env.getOrDefault("STORE_MODE", "disk");
|
||||
ObjectStorage store;
|
||||
MultipartStorage multipart;
|
||||
@@ -1206,11 +1254,13 @@ public final class Main {
|
||||
multipart = new MultipartStore(disk);
|
||||
} else throw new IllegalArgumentException("Invalid STORE_MODE");
|
||||
var app = new Main(store, multipart, new SigV4(identities(env, access, secret),
|
||||
access, region, Clock.systemUTC()), bucket, ClientLimits.fromEnvironment(env));
|
||||
access, region, Clock.systemUTC()), bucket, ClientLimits.fromEnvironment(env),
|
||||
positiveInt(env, "MAX_IN_FLIGHT_REQUESTS", 16, 1024),
|
||||
env.getOrDefault("S3_VIRTUAL_HOST_SUFFIX", ""));
|
||||
int port = Integer.parseInt(env.getOrDefault("PORT", "9000"));
|
||||
var server = HttpServer.create(mode.equals("cluster")
|
||||
? new InetSocketAddress(env.getOrDefault("BIND_ADDRESS", "127.0.0.1"), port)
|
||||
: new InetSocketAddress(port), 64);
|
||||
: new InetSocketAddress(port), positiveInt(env, "HTTP_BACKLOG", 64, 4096));
|
||||
var executor = Executors.newVirtualThreadPerTaskExecutor();
|
||||
server.setExecutor(executor);
|
||||
server.createContext("/", app::handle);
|
||||
@@ -1277,4 +1327,16 @@ public final class Main {
|
||||
if (value == null || value.isBlank()) throw new IllegalArgumentException("Missing " + key);
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -21,18 +21,21 @@ final class NodeClient {
|
||||
record Node(UUID id, UUID hostId, URI url) {}
|
||||
record StoredSegment(UUID id, long modified) {}
|
||||
|
||||
private static final HttpClient IDENTITY_HTTP = HttpClient.newBuilder()
|
||||
.connectTimeout(Duration.ofSeconds(2)).build();
|
||||
private static HttpClient identityHttp;
|
||||
private final List<Node> nodes;
|
||||
private final String token;
|
||||
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> healthyUntil = new ConcurrentHashMap<>();
|
||||
private static final long READ_RETRY_NANOS = TimeUnit.SECONDS.toNanos(5);
|
||||
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() ||
|
||||
nodes.stream().map(Node::url).distinct().count() != nodes.size())
|
||||
throw new IllegalArgumentException("Cluster node IDs and URLs must be unique");
|
||||
@@ -45,6 +48,7 @@ final class NodeClient {
|
||||
this.nodes = List.copyOf(nodes);
|
||||
this.token = token;
|
||||
this.repairToken = repairToken;
|
||||
this.http = java.util.Objects.requireNonNull(http);
|
||||
}
|
||||
|
||||
int count() { return nodes.size(); }
|
||||
@@ -63,12 +67,16 @@ final class NodeClient {
|
||||
}
|
||||
|
||||
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);
|
||||
if (token == null || token.length() < 32) throw new IllegalArgumentException("Invalid cluster token");
|
||||
HttpRequest request = HttpRequest.newBuilder(url.resolve("/identity"))
|
||||
.timeout(Duration.ofSeconds(2)).header("X-Cluster-Token", token).GET().build();
|
||||
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()) {
|
||||
if (response.statusCode() != 200) throw new IOException("Node identity request failed: " + response.statusCode());
|
||||
byte[] bytes = body.readNBytes(128);
|
||||
@@ -90,19 +98,38 @@ final class NodeClient {
|
||||
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) {
|
||||
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.getRawPath() != null && !url.getRawPath().isEmpty()) ||
|
||||
url.getRawQuery() != null || url.getRawFragment() != null)
|
||||
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) {
|
||||
Set<UUID> healthy = new HashSet<>();
|
||||
for (int i = 0; i < nodes.size(); 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())) {
|
||||
markUnreadable(node);
|
||||
continue;
|
||||
@@ -159,7 +186,7 @@ final class NodeClient {
|
||||
Node node = nodes.get(index);
|
||||
if (unreadable(node)) throw new IOException("Storage node is temporarily unreachable");
|
||||
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())) {
|
||||
markUnreadable(node);
|
||||
throw new IOException("Storage node is temporarily unreachable");
|
||||
|
||||
@@ -74,10 +74,36 @@ final class SigV4 {
|
||||
|
||||
private Verified verifyPresigned(String method, URI uri, Headers headers) {
|
||||
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<>();
|
||||
StringBuilder application = 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 name = decode(pair[0]);
|
||||
String value = decode(pair.length == 2 ? pair[1] : "");
|
||||
@@ -89,17 +115,19 @@ final class SigV4 {
|
||||
appendQuery(signed, part);
|
||||
}
|
||||
}
|
||||
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");
|
||||
if (!date.matches("[0-9]{8}T[0-9]{6}Z") || !date.startsWith(credential[1]))
|
||||
return new PresignedQuery(fields, application.toString(), signed.toString());
|
||||
}
|
||||
|
||||
private void validatePresignedTime(String date, String scopeDate, String rawExpires) {
|
||||
if (!date.matches("[0-9]{8}T[0-9]{6}Z") || !date.startsWith(scopeDate))
|
||||
denied("Invalid signing date");
|
||||
long expires;
|
||||
try { expires = Long.parseLong(fields.get("X-Amz-Expires")); }
|
||||
catch (NumberFormatException error) { denied("Invalid presigned expiry"); return null; }
|
||||
try {
|
||||
expires = Long.parseLong(rawExpires);
|
||||
} catch (NumberFormatException error) {
|
||||
denied("Invalid presigned expiry");
|
||||
return;
|
||||
}
|
||||
if (expires < 1 || expires > 604800) denied("Invalid presigned expiry");
|
||||
try {
|
||||
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)))
|
||||
denied("Presigned URL has expired or is not yet valid");
|
||||
} 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) {
|
||||
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");
|
||||
}
|
||||
}
|
||||
@@ -806,15 +806,70 @@ public final class HttpTest {
|
||||
throw new AssertionError("Completed multipart version was not retained");
|
||||
}
|
||||
|
||||
private static int virtualHostPut(int port, String host, String signedHost, byte[] body) throws Exception {
|
||||
URI signed = URI.create("http://" + signedHost + ":" + port + "/photos/cat.jpg");
|
||||
HttpRequest request = signedUri(signed, "PUT", body, Map.of());
|
||||
String headers = "PUT /photos/cat.jpg HTTP/1.1\r\nHost: " + host + ":" + port +
|
||||
"\r\nAuthorization: " + request.headers().firstValue("authorization").orElseThrow() +
|
||||
"\r\nx-amz-date: " + request.headers().firstValue("x-amz-date").orElseThrow() +
|
||||
"\r\nx-amz-content-sha256: " + request.headers().firstValue("x-amz-content-sha256").orElseThrow() +
|
||||
"\r\nContent-Length: " + body.length + "\r\nConnection: close\r\n\r\n";
|
||||
try (var socket = new java.net.Socket("127.0.0.1", port)) {
|
||||
socket.setSoTimeout(5000);
|
||||
socket.getOutputStream().write(headers.getBytes(StandardCharsets.US_ASCII));
|
||||
socket.getOutputStream().write(body);
|
||||
String response = new String(socket.getInputStream().readAllBytes(), StandardCharsets.ISO_8859_1);
|
||||
return Integer.parseInt(response.split(" ", 3)[1]);
|
||||
}
|
||||
}
|
||||
|
||||
private static void testGatewayConcurrency(DiskStore store, HttpClient client,
|
||||
java.util.concurrent.Executor executor) throws Exception {
|
||||
HttpServer limited = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 16);
|
||||
var app = new Main(store, new MultipartStore(store),
|
||||
new SigV4(ACCESS, SECRET, REGION, Clock.systemUTC()), "objects",
|
||||
ClientLimits.disabled(), 1, "");
|
||||
limited.setExecutor(executor);
|
||||
limited.createContext("/", app::handle);
|
||||
limited.start();
|
||||
int port = limited.getAddress().getPort();
|
||||
String base = "http://127.0.0.1:" + port;
|
||||
byte[] body = {42};
|
||||
HttpRequest request = signedUri(URI.create(base + "/objects/held"), "PUT", body, Map.of());
|
||||
try (var socket = new java.net.Socket("127.0.0.1", port)) {
|
||||
socket.setSoTimeout(5000);
|
||||
String headers = "PUT /objects/held HTTP/1.1\r\nHost: 127.0.0.1:" + port +
|
||||
"\r\nAuthorization: " + request.headers().firstValue("authorization").orElseThrow() +
|
||||
"\r\nx-amz-date: " + request.headers().firstValue("x-amz-date").orElseThrow() +
|
||||
"\r\nx-amz-content-sha256: " + request.headers().firstValue("x-amz-content-sha256").orElseThrow() +
|
||||
"\r\nContent-Length: 1\r\nConnection: close\r\n\r\n";
|
||||
socket.getOutputStream().write(headers.getBytes(StandardCharsets.US_ASCII));
|
||||
boolean refused = false;
|
||||
for (int attempt = 0; attempt < 50 && !refused; attempt++) {
|
||||
int status = client.send(HttpRequest.newBuilder(URI.create(base + "/health")).GET().build(),
|
||||
HttpResponse.BodyHandlers.discarding()).statusCode();
|
||||
refused = status == 503;
|
||||
if (!refused) Thread.sleep(10);
|
||||
}
|
||||
if (!refused) throw new AssertionError("Configured gateway concurrency cap was not enforced");
|
||||
socket.getOutputStream().write(body);
|
||||
String finished = new String(socket.getInputStream().readAllBytes(), StandardCharsets.ISO_8859_1);
|
||||
if (!finished.startsWith("HTTP/1.1 200 "))
|
||||
throw new AssertionError("Held request did not complete after releasing its body");
|
||||
} finally {
|
||||
limited.stop(0);
|
||||
}
|
||||
}
|
||||
|
||||
public static void main(String[] args) throws Exception {
|
||||
Path root = Files.createTempDirectory("store-http-test-");
|
||||
var executor = Executors.newVirtualThreadPerTaskExecutor();
|
||||
HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 16);
|
||||
DiskStore store = new DiskStore(root, 1024, 4096);
|
||||
try {
|
||||
var app = new Main(store,
|
||||
var app = new Main(store, new MultipartStore(store),
|
||||
new SigV4(Map.of(ACCESS, SECRET, SECONDARY, SECONDARY_SECRET),
|
||||
ACCESS, REGION, Clock.systemUTC()), "objects");
|
||||
ACCESS, REGION, Clock.systemUTC()), "objects", ClientLimits.disabled(), 16, "s3.test");
|
||||
server.setExecutor(executor);
|
||||
server.createContext("/", app::handle);
|
||||
server.start();
|
||||
@@ -834,6 +889,17 @@ public final class HttpTest {
|
||||
testBuckets(client, base);
|
||||
testVersioning(client, base);
|
||||
testAcl(client, base);
|
||||
byte[] hostedBody = "virtual host".getBytes(StandardCharsets.UTF_8);
|
||||
int port = server.getAddress().getPort();
|
||||
if (virtualHostPut(port, "objects.s3.test", "objects.s3.test", hostedBody) != 200)
|
||||
throw new AssertionError("Signed virtual-hosted upload failed");
|
||||
try (var opened = store.open("objects", "photos/cat.jpg")) {
|
||||
if (!java.util.Arrays.equals(hostedBody, opened.stream().readAllBytes()))
|
||||
throw new AssertionError("Virtual-hosted bucket or key was parsed incorrectly");
|
||||
}
|
||||
if (virtualHostPut(port, "objects.s3.test", "other.s3.test", hostedBody) != 403)
|
||||
throw new AssertionError("Changing a signed virtual hostname was accepted");
|
||||
testGatewayConcurrency(store, client, executor);
|
||||
System.out.println("HTTP tests passed: capabilities, objects, copy, checksums, listing, multipart, attributes, buckets, versioning, ACLs");
|
||||
} finally {
|
||||
server.stop(0);
|
||||
|
||||
@@ -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,6 +1,6 @@
|
||||
# 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.
|
||||
|
||||
|
||||
Reference in new issue
Block a user