diff --git a/README.md b/README.md index fc338c9..90e260e 100644 --- a/README.md +++ b/README.md @@ -1,13 +1,22 @@ # cephrgw-operator A Kubernetes operator that provisions Ceph RGW (S3) **buckets** and **access -keys** declaratively, driving the Ceph **manager dashboard REST API**. You +keys** declaratively, talking **directly to radosgw** the way the CLI does. You describe a bucket, its owner, and who may read or write it as custom resources; the operator creates the RGW users and bucket, delivers the access/secret keys into Kubernetes Secrets, and maintains the bucket's S3 policy. -It talks only to the dashboard API (e.g. `https://dashboard.ceph.unkin.net`) — -no RADOS access, no admin socket, no in-cluster Ceph required. +It uses native Go libraries against radosgw (e.g. +`https://radosgw.service.consul:443`) — no manager dashboard, no RADOS access, +no admin socket, no in-cluster Ceph required: + +- **[go-ceph](https://github.com/ceph/go-ceph) `rgw/admin`** drives the RGW + **Admin Ops API** (`/admin/...`) for users, keys, quotas and bucket + info/removal — pure Go, no cgo. +- **[aws-sdk-go-v2](https://github.com/aws/aws-sdk-go-v2)** drives the **S3 API** + for bucket creation, versioning, policy, tagging and object lock — the + operations the Admin Ops API does not expose. These are signed as the bucket + **owner**, so the owner owns the bucket directly. ## Custom resources @@ -37,10 +46,10 @@ order objects are created or deleted. It watches `BucketAccess` and `ObjectStoreUser`, re-reconciling the bucket whenever a grant or user changes. ``` -ObjectStoreUser ──create user──▶ dashboard /api/rgw/user ──▶ Secret (AK/SK) -Bucket ──create bucket─▶ dashboard /api/rgw/bucket ─▶ owns S3 policy -BucketAccess ──ensure user───▶ dashboard /api/rgw/user ──▶ Secret (AK/SK, RW or RO) - └────── enqueues Bucket ──▶ PUT bucket_policy (aggregate) +ObjectStoreUser ──admin PUT /admin/user──────▶ Secret (AK/SK) +Bucket ──S3 CreateBucket (as owner)─▶ owns S3 policy +BucketAccess ──admin PUT /admin/user──────▶ Secret (AK/SK, RW or RO) + └────── enqueues Bucket ──▶ S3 PutBucketPolicy (aggregate) ``` ## Credential Secrets @@ -50,7 +59,7 @@ into a workload: - `AWS_ACCESS_KEY_ID`, `AWS_SECRET_ACCESS_KEY` - `RGW_UID` -- `S3_ENDPOINT`, `BUCKET_HOST` (when `CEPH_RGW_ENDPOINT` is configured) +- `S3_ENDPOINT`, `BUCKET_HOST` (when `CEPH_RGW_ENDPOINT` is set) - `BUCKET_NAME` (on `BucketAccess` Secrets) Secrets are owner-referenced by the resource that produced them, so they are @@ -58,10 +67,11 @@ garbage-collected when the resource is deleted. ## Prerequisites -The operator needs a dashboard login with the `rgw-manager` role, a dashboard -that is wired to RGW, and (for `read-only`/non-owner `read-write` grants) Ceph -**Reef 18.2+ / Squid**. See **[docs/ceph-setup.md](docs/ceph-setup.md)** for the -exact commands and the `cephrgw-credentials` Secret schema. +The operator needs an RGW user with admin caps (`users=*;buckets=*`) and its +access/secret key, the radosgw endpoint, and (for `read-only`/non-owner +`read-write` grants) Ceph **Reef 18.2+ / Squid**. See +**[docs/ceph-setup.md](docs/ceph-setup.md)** for the exact commands and the +`cephrgw-credentials` Secret schema. ## Quickstart @@ -104,12 +114,11 @@ with `make patch|minor|major`. ## Notes & caveats -- **Policy clearing.** Removing the last `BucketAccess` asks the dashboard to - clear the bucket policy. Not every release honours an empty policy string; if - a stale policy lingers, clear it once by hand. Adding/replacing grants always - works. +- **Policy clearing.** Removing the last `BucketAccess` issues an S3 + `DeleteBucketPolicy`. A `NoSuchBucketPolicy` response is treated as already + clear. Adding/replacing grants always works. - **Per-bucket quota.** `Bucket.spec.quota` is applied as the owner's default - bucket quota via the dashboard, which is per-owner rather than strictly + bucket quota via the Admin Ops API, which is per-owner rather than strictly per-bucket. Use distinct owners if you need independent bucket quotas. - **Immutability.** `bucketName`, an `ObjectStoreUser`'s `uid`, and object lock are fixed at creation; changing them on an existing object has no effect. diff --git a/api/v1alpha1/objectstoreuser_types.go b/api/v1alpha1/objectstoreuser_types.go index 5439623..6831c3d 100644 --- a/api/v1alpha1/objectstoreuser_types.go +++ b/api/v1alpha1/objectstoreuser_types.go @@ -5,7 +5,7 @@ import ( ) // ObjectStoreUserSpec defines a Ceph RGW (S3) user. The operator creates the -// user through the Ceph dashboard API and writes its generated access/secret +// user through the radosgw Admin Ops API and writes its generated access/secret // key pair into a Kubernetes Secret. The key material is never stored on the // resource itself. type ObjectStoreUserSpec struct { diff --git a/cmd/operator/main.go b/cmd/operator/main.go index 4ca6757..7d9e0f0 100644 --- a/cmd/operator/main.go +++ b/cmd/operator/main.go @@ -40,19 +40,19 @@ func main() { cephCfg, endpoint, err := cephConfigFromEnv() if err != nil { - logger.Error(err, "invalid Ceph dashboard configuration") + logger.Error(err, "invalid radosgw configuration") os.Exit(1) } cephClient, err := ceph.NewClient(cephCfg) if err != nil { - logger.Error(err, "unable to build Ceph dashboard client") + logger.Error(err, "unable to build radosgw client") os.Exit(1) } // Fail fast on obviously-broken credentials, but do not block startup on a - // transiently unreachable dashboard. + // transiently unreachable radosgw. pingCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second) if err := cephClient.Ping(pingCtx); err != nil { - logger.Error(err, "initial dashboard authentication failed; continuing and will retry per-reconcile") + logger.Error(err, "initial radosgw authentication failed; continuing and will retry per-reconcile") } cancel() @@ -89,24 +89,34 @@ func main() { } } -// cephConfigFromEnv reads dashboard connection settings from the environment, +// cephConfigFromEnv reads radosgw connection settings from the environment, // which the deployment sources from the cephrgw-credentials Secret. +// +// The operator talks to the radosgw Admin Ops and S3 APIs at +// CEPH_RGW_ADMIN_ENDPOINT (falling back to CEPH_RGW_ENDPOINT). The returned +// endpoint string is the S3 endpoint written into consumer credential Secrets, +// which may differ (e.g. a public S3 name) from the API endpoint. func cephConfigFromEnv() (ceph.Config, string, error) { - cfg := ceph.Config{ - BaseURL: os.Getenv("CEPH_DASHBOARD_URL"), - Username: os.Getenv("CEPH_DASHBOARD_USERNAME"), - Password: os.Getenv("CEPH_DASHBOARD_PASSWORD"), - Insecure: os.Getenv("CEPH_DASHBOARD_INSECURE") == "true", + consumerEndpoint := os.Getenv("CEPH_RGW_ENDPOINT") + apiEndpoint := os.Getenv("CEPH_RGW_ADMIN_ENDPOINT") + if apiEndpoint == "" { + apiEndpoint = consumerEndpoint } - if f := os.Getenv("CEPH_DASHBOARD_CA_FILE"); f != "" { + cfg := ceph.Config{ + Endpoint: apiEndpoint, + AccessKey: os.Getenv("CEPH_RGW_ACCESS_KEY"), + SecretKey: os.Getenv("CEPH_RGW_SECRET_KEY"), + Region: os.Getenv("CEPH_RGW_REGION"), + Insecure: os.Getenv("CEPH_RGW_INSECURE") == "true", + } + if f := os.Getenv("CEPH_RGW_CA_FILE"); f != "" { b, err := os.ReadFile(f) if err != nil { return cfg, "", err } cfg.CACert = b - } else if inline := os.Getenv("CEPH_DASHBOARD_CA"); inline != "" { + } else if inline := os.Getenv("CEPH_RGW_CA"); inline != "" { cfg.CACert = []byte(inline) } - endpoint := os.Getenv("CEPH_RGW_ENDPOINT") - return cfg, endpoint, nil + return cfg, consumerEndpoint, nil } diff --git a/docs/ceph-setup.md b/docs/ceph-setup.md index 12a84d0..caa2a4c 100644 --- a/docs/ceph-setup.md +++ b/docs/ceph-setup.md @@ -1,102 +1,73 @@ # Ceph setup: credentials and permissions the operator needs -`cephrgw-operator` never talks to RADOS or the RGW admin socket directly. It -drives the **Ceph manager dashboard REST API** (the same API the web dashboard -uses) at `https://dashboard.ceph.unkin.net`. Everything below is about giving -the operator a dashboard login with enough RGW authority, and making sure the -dashboard itself is wired to your RGW. +`cephrgw-operator` talks **directly to radosgw**, the same way the `radosgw-admin` +CLI and S3 clients do — no manager dashboard involved. It uses two native Go +libraries against the RGW endpoint (e.g. `https://radosgw.service.consul:443`): -There are **two** credentials involved. Don't confuse them: +- **go-ceph `rgw/admin`** → the RGW **Admin Ops API** (`/admin/user`, + `/admin/bucket`), signed with the operator's access/secret key, for users, + keys, quotas and bucket info/removal. +- **aws-sdk-go-v2** → the **S3 API**, signed as each bucket's **owner**, for + bucket creation, versioning, policy, tagging and object lock. -| # | Credential | Who uses it | What it is | -|---|------------|-------------|------------| -| 1 | Dashboard login (username + password) | the operator → `POST /api/auth` | a **dashboard account** with the `rgw-manager` role | -| 2 | RGW admin connection | the dashboard → RGW | a **radosgw system user** (access/secret key) the dashboard is configured with | - -The operator only holds #1. #2 is what actually lets the dashboard create RGW -users, buckets and bucket policies on the operator's behalf, so it must exist -and be privileged. +So there is exactly **one** credential to provision: a radosgw user with admin +caps, plus its access/secret key. --- -## 1. Create the dashboard login for the operator +## 1. Create the operator's RGW admin user -Create a dedicated dashboard user with the built-in **`rgw-manager`** role. That -role grants full create/read/update/delete on the dashboard's `rgw` scope -(users, buckets, policies) and nothing else — least privilege for this operator. - -```bash -# Put the password in a file so it never lands in shell history. -printf '%s' 'REPLACE-WITH-A-STRONG-PASSWORD' > /tmp/cephrgw.pw - -ceph dashboard ac-user-create k8s-cephrgw-operator -i /tmp/cephrgw.pw rgw-manager - -rm -f /tmp/cephrgw.pw -``` - -If your Ceph version wants the arguments in a different order, check -`ceph dashboard ac-user-create -h`. To confirm the role exists and what it -grants: - -```bash -ceph dashboard ac-role-show rgw-manager -``` - -> Prefer `rgw-manager` over `administrator`. The operator only needs RGW -> authority; giving it full dashboard admin is unnecessary blast radius. - -## 2. Make sure the dashboard can manage RGW - -The dashboard performs RGW operations through a **radosgw system user**. On -recent Ceph (Pacific and later) the mgr/dashboard module usually auto-discovers -and configures this. Verify it first: - -```bash -ceph dashboard get-rgw-api-access-key # should print a key, not empty -``` - -If it is empty, create a system user and point the dashboard at it: +Create a dedicated radosgw user and give it the admin caps the operator needs. +Only `users` and `buckets` caps are required (the operator never reads usage or +metadata endpoints): ```bash radosgw-admin user create \ - --uid=dashboard \ - --display-name="Ceph Dashboard" \ - --system + --uid=cephrgw-operator \ + --display-name="cephrgw-operator" \ + --caps="users=*;buckets=*" -# Feed the returned keys to the dashboard. -radosgw-admin user info --uid=dashboard \ - | jq -r '.keys[0].access_key' > /tmp/ak -radosgw-admin user info --uid=dashboard \ - | jq -r '.keys[0].secret_key' > /tmp/sk - -ceph dashboard set-rgw-api-access-key -i /tmp/ak -ceph dashboard set-rgw-api-secret-key -i /tmp/sk -rm -f /tmp/ak /tmp/sk +# Grab its keys (these become CEPH_RGW_ACCESS_KEY / CEPH_RGW_SECRET_KEY): +radosgw-admin user info --uid=cephrgw-operator \ + | jq -r '.keys[0] | .access_key, .secret_key' ``` -A `--system` user has the admin caps the dashboard needs to create/delete RGW -users and buckets and to set bucket policies on any bucket. If you would rather -not use `--system`, grant an equivalent admin cap set instead: +If the user already exists, add the caps instead: ```bash -radosgw-admin caps add --uid=dashboard \ - --caps="users=*;buckets=*;metadata=*;usage=read;zone=read" +radosgw-admin caps add --uid=cephrgw-operator --caps="users=*;buckets=*" ``` -If the dashboard reaches RGW over TLS with a private CA, you may also need: +> `users=*;buckets=*` lets the operator create/read/delete RGW users and read/ +> remove buckets through the Admin Ops API. Bucket **creation** and all bucket +> sub-resources (versioning, policy, tagging, object lock) go over the S3 API +> signed as the bucket owner, so they need no extra admin cap — every RGW user +> can manage its own buckets. The `--system` flag is **not** required. + +## 2. Admin Ops API must be enabled on radosgw + +The Admin Ops API is served by radosgw at the `admin` resource and is enabled by +default. If your deployment has trimmed `rgw_enable_apis`, make sure it includes +both `s3` and `admin`: + +``` +rgw_enable_apis = s3, admin +``` + +Quick check from your workstation (a `403`/`AccessDenied` still proves the +endpoint is reachable and the API is on; a connection error means it is not): ```bash -ceph dashboard set-rgw-api-ssl-verify true # keep verification on in prod +curl -sk "https://radosgw.service.consul:443/admin/user?format=json" ``` ## 3. Bucket policy support (read-only / non-owner read-write) -The operator enforces `read-only` and non-owner `read-write` grants by writing -an **S3 bucket policy** through the dashboard's bucket API (the `bucket_policy` -field on `PUT /api/rgw/bucket/{name}`). That field is available on **Ceph Reef -18.2+ / Squid**. On older releases bucket creation and owner (`full`) access -still work, but policy-based grants will fail — upgrade the cluster, or only use -owner credentials, if you are pre-Reef. +The operator enforces `read-only` and non-owner `read-write` grants by writing an +**S3 bucket policy** (`PutBucketPolicy`). Bucket-policy support is available on +**Ceph Reef 18.2+ / Squid**. On older releases bucket creation and owner +(`full`) access still work, but policy-based grants will fail — upgrade the +cluster, or only use owner credentials, if you are pre-Reef. Check your version: @@ -108,27 +79,31 @@ ceph versions | jq -r '.mon | keys[]' The operator can stamp the S3 endpoint into every credential Secret it writes (`S3_ENDPOINT` and `BUCKET_HOST`) so applications don't have to hard-code it. -This is the RGW/S3 endpoint your clients use — **not** the dashboard URL. Provide -it via `CEPH_RGW_ENDPOINT` (see below); if unset, those keys are simply omitted. +Provide it via `CEPH_RGW_ENDPOINT` (see below); if unset, those keys are simply +omitted. This is also the default endpoint for the Admin Ops and S3 API calls +when `CEPH_RGW_ADMIN_ENDPOINT` is not set separately. --- ## 5. Give the operator its credentials (the `cephrgw-credentials` Secret) The operator reads its configuration from environment variables, which the -deployment sources (via `envFrom`) from a Secret named **`cephrgw-credentials`** -in its namespace (`cephrgw-system`). The Secret data keys map 1:1 to the env -vars: +deployment sources from a Secret named **`cephrgw-credentials`** in its namespace +(`cephrgw-system`). The Secret data keys map 1:1 to the env vars: | Secret key | Required | Meaning | |------------|----------|---------| -| `CEPH_DASHBOARD_URL` | yes | dashboard base URL, e.g. `https://dashboard.ceph.unkin.net` | -| `CEPH_DASHBOARD_USERNAME` | yes | the `rgw-manager` account from step 1 | -| `CEPH_DASHBOARD_PASSWORD` | yes | its password | -| `CEPH_RGW_ENDPOINT` | no | S3 endpoint written into consumer Secrets | -| `CEPH_DASHBOARD_CA` | no | PEM CA bundle to verify the dashboard TLS cert (inline) | -| `CEPH_DASHBOARD_CA_FILE` | no | path to a mounted CA file (alternative to the above) | -| `CEPH_DASHBOARD_INSECURE` | no | `"true"` to skip TLS verification (dev only) | +| `CEPH_RGW_ACCESS_KEY` | yes | access key of the RGW admin user from step 1 | +| `CEPH_RGW_SECRET_KEY` | yes | its secret key | +| `CEPH_RGW_ENDPOINT` | see note | S3 endpoint; written into consumer Secrets and used for API calls unless `CEPH_RGW_ADMIN_ENDPOINT` is set | +| `CEPH_RGW_ADMIN_ENDPOINT` | no | radosgw endpoint for the Admin Ops + S3 API calls, if it differs from the public `CEPH_RGW_ENDPOINT` | +| `CEPH_RGW_REGION` | no | SigV4 credential-scope region for S3 requests (default `default`) | +| `CEPH_RGW_CA` | no | PEM CA bundle to verify the radosgw TLS cert (inline) | +| `CEPH_RGW_CA_FILE` | no | path to a mounted CA file (alternative to the above) | +| `CEPH_RGW_INSECURE` | no | `"true"` to skip TLS verification (dev only) | + +> At least one of `CEPH_RGW_ENDPOINT` or `CEPH_RGW_ADMIN_ENDPOINT` must be set — +> the API endpoint falls back to `CEPH_RGW_ENDPOINT` when the admin one is unset. ### Primary method: Vault + VSO @@ -143,10 +118,10 @@ Vault role or policy is required** — you only seed the values: ```bash vault kv put kv/kubernetes/namespace/cephrgw-system/default/cephrgw-credentials \ - CEPH_DASHBOARD_URL=https://dashboard.ceph.unkin.net \ - CEPH_DASHBOARD_USERNAME=k8s-cephrgw-operator \ - CEPH_DASHBOARD_PASSWORD='REPLACE-WITH-A-STRONG-PASSWORD' \ - CEPH_RGW_ENDPOINT=https://s3.ceph.unkin.net + CEPH_RGW_ENDPOINT=https://s3.ceph.unkin.net \ + CEPH_RGW_ADMIN_ENDPOINT=https://radosgw.service.consul:443 \ + CEPH_RGW_ACCESS_KEY='REPLACE-WITH-ACCESS-KEY' \ + CEPH_RGW_SECRET_KEY='REPLACE-WITH-SECRET-KEY' ``` The keys under that KV path are copied verbatim into the Secret, so they must be @@ -162,10 +137,10 @@ directly instead of using Vault: ```bash kubectl -n cephrgw-system create secret generic cephrgw-credentials \ - --from-literal=CEPH_DASHBOARD_URL=https://dashboard.ceph.unkin.net \ - --from-literal=CEPH_DASHBOARD_USERNAME=k8s-cephrgw-operator \ - --from-literal=CEPH_DASHBOARD_PASSWORD='REPLACE-WITH-A-STRONG-PASSWORD' \ - --from-literal=CEPH_RGW_ENDPOINT=https://s3.ceph.unkin.net + --from-literal=CEPH_RGW_ENDPOINT=https://s3.ceph.unkin.net \ + --from-literal=CEPH_RGW_ADMIN_ENDPOINT=https://radosgw.service.consul:443 \ + --from-literal=CEPH_RGW_ACCESS_KEY='REPLACE-WITH-ACCESS-KEY' \ + --from-literal=CEPH_RGW_SECRET_KEY='REPLACE-WITH-SECRET-KEY' ``` The operator does not care where the Secret comes from, only that those keys @@ -175,22 +150,19 @@ exist. ## Quick verification -Once the Secret and dashboard account exist, a smoke test from your workstation: +Once the Secret and RGW admin user exist, a smoke test from your workstation +using the operator's keys (this is the same Admin Ops call the operator's +readiness `Ping` makes): ```bash -# 1. Log in and capture a token. -TOKEN=$(curl -sk -X POST https://dashboard.ceph.unkin.net/api/auth \ - -H 'Accept: application/vnd.ceph.api.v1.0+json' \ - -H 'Content-Type: application/json' \ - -d '{"username":"k8s-cephrgw-operator","password":"REPLACE-WITH-A-STRONG-PASSWORD"}' \ - | jq -r .token) - -# 2. List RGW users — a 200 with a JSON array means the role + RGW wiring work. -curl -sk https://dashboard.ceph.unkin.net/api/rgw/user \ - -H 'Accept: application/vnd.ceph.api.v1.0+json' \ - -H "Authorization: Bearer $TOKEN" +# Signing an Admin Ops request by hand is fiddly; the simplest proof is to use +# the AWS CLI configured with the operator's keys against the S3 endpoint: +AWS_ACCESS_KEY_ID=REPLACE-WITH-ACCESS-KEY \ +AWS_SECRET_ACCESS_KEY=REPLACE-WITH-SECRET-KEY \ +aws --endpoint-url https://s3.ceph.unkin.net s3 ls ``` -If step 1 fails the login/role is wrong (step 1–2 above); if step 1 works but -step 2 returns 500/empty, the dashboard→RGW connection is not configured -(step 2). +A successful (even empty) listing proves the keys and endpoint work. If the +operator logs `initial radosgw authentication failed`, the keys are wrong or the +`admin` API is disabled (steps 1–2); if users are created but bucket policy +grants fail, the cluster is likely pre-Reef (step 3). diff --git a/go.mod b/go.mod index 86953b1..0a7e64d 100644 --- a/go.mod +++ b/go.mod @@ -1,8 +1,13 @@ module git.unkin.net/unkin/cephrgw-operator -go 1.25 +go 1.25.0 require ( + github.com/aws/aws-sdk-go-v2 v1.43.0 + github.com/aws/aws-sdk-go-v2/credentials v1.19.30 + github.com/aws/aws-sdk-go-v2/service/s3 v1.106.0 + github.com/aws/smithy-go v1.27.4 + github.com/ceph/go-ceph v0.40.0 k8s.io/api v0.34.4 k8s.io/apimachinery v0.34.4 k8s.io/client-go v0.34.4 @@ -10,6 +15,14 @@ require ( ) require ( + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.14 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.31 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.31 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.32 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.13 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.24 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.31 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.32 // indirect github.com/beorn7/perks v1.0.1 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/davecgh/go-spew v1.1.1 // indirect @@ -48,7 +61,7 @@ require ( golang.org/x/net v0.38.0 // indirect golang.org/x/oauth2 v0.27.0 // indirect golang.org/x/sync v0.12.0 // indirect - golang.org/x/sys v0.31.0 // indirect + golang.org/x/sys v0.45.0 // indirect golang.org/x/term v0.30.0 // indirect golang.org/x/text v0.23.0 // indirect golang.org/x/time v0.9.0 // indirect diff --git a/go.sum b/go.sum index 2261246..bb98464 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,43 @@ +github.com/aws/aws-sdk-go-v2 v1.43.0 h1:fharf/WhbRAVZ1du0QL7roNFxZ6T/sWr+4Ni617bwSI= +github.com/aws/aws-sdk-go-v2 v1.43.0/go.mod h1:5pKeft2eJj+gElQ38Jqg4ibCqh+/AK33/0X3hip7IjM= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.14 h1:3IZY0XAJquT3aHzbkHfPzy4ACPcEjVG0x87KOwtpqGY= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.14/go.mod h1:zwM6veDkhGgQFqkBy+uT28AAYpLu+uFMlPl+rCg/73E= +github.com/aws/aws-sdk-go-v2/config v1.32.22 h1:Vfvp7+fYKsVCADcWOEllqEV47aIBXhNchvyDFu1B5fY= +github.com/aws/aws-sdk-go-v2/config v1.32.22/go.mod h1:0+H+0nPKbvWltf5vSIGkApv+hGbaQ4FfwTjGIYQREcw= +github.com/aws/aws-sdk-go-v2/credentials v1.19.30 h1:TTCvvzFU6gXa4iJecNG/0F/B0oYTiazoRECr2XyLHrY= +github.com/aws/aws-sdk-go-v2/credentials v1.19.30/go.mod h1:jKxAp2AEncnliinzpgOSZDFv6+VjvWhjw/AtbfsWT9U= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.31 h1:kfVL5wAunCJycL6MOQ6aNh6PlAYEymflcjuKmrWUA0o= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.31/go.mod h1:nWfRNDAppujCQgOUd43lKT4yeLv9z3nJ3bw1G3BgQKo= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.31 h1:Z8F3hfCY33IGpJjFAnv0wvtv1FIKj1GHmRDEYqy64tw= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.31/go.mod h1:aVyUoytEyOViR6jhq6jula0xkc5NfBE2hgeF6BvOrao= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.31 h1:hyOxUyXdh3AyjE93gBgsfziJag9ACwcs+ZpDBLzi8mw= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.31/go.mod h1:OERqI9k0draSLB8O8woxY3q25ZWTELRK4RRoLMuMZFo= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.32 h1:0MrUL35H/Y4kdFfItoR5jCgtDQ4Z/8LudAoIHRfA4hE= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.32/go.mod h1:2tNZkuWz54arj8mHVf+8Y7cKkcD8Wr/fBpENgEXpjLc= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.13 h1:mbRIur/BiHK6SKPjoBIXSE/hJ6g6JGRLuxQy1jGjlN4= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.13/go.mod h1:ITg9em2KbJx1s0y4aqRX5OYWG6HBZ5TVR//OdpEZ2CQ= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.24 h1:mdPwDQPqxlw9Sc62Nt15yjEcARaDbPXkjRYtXsUripo= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.24/go.mod h1:ls5ytnwLTcQaUu32fMYXFI3MjpKuTwL840PAm9iqyEg= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.31 h1:w2SIhW92DZPFrSL4ksVCr8IYff5OZwIcxg8+95tzvAI= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.31/go.mod h1:wAhpCQbkov+IcvjozJbd2xRCoZybUEHNkcFunssNACg= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.32 h1:jWXtZdCnhXa9sGFixRaU2AxT4DIVse9HS4E2f+/KwV0= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.32/go.mod h1:9JS1UpfVvyD/ZPX8GsKb/Pq8scEM+7GP5fqh9SwH7po= +github.com/aws/aws-sdk-go-v2/service/s3 v1.106.0 h1:7QZWVJZWzHivHWIa+5TELLaBBkbuoj0GPwQtMlJ0sqk= +github.com/aws/aws-sdk-go-v2/service/s3 v1.106.0/go.mod h1:fcvq5L7dK+5cQFicEJwpI6e6Wn8NY2i6yT5wRLYVc7s= +github.com/aws/aws-sdk-go-v2/service/signin v1.5.0 h1:OHH5iTQvVGmfHjX/5Q+vFuA/Rf2x6/95aJ/75QCQSm4= +github.com/aws/aws-sdk-go-v2/service/signin v1.5.0/go.mod h1:mCF3AK9PpL49oOrhniUXWAfhVBVQ/XbytoE5eccZUIs= +github.com/aws/aws-sdk-go-v2/service/sso v1.33.0 h1:CaJyYhxBE0M/HJX/YvSaSmQlsI91VHB0lKU8LtLxL3A= +github.com/aws/aws-sdk-go-v2/service/sso v1.33.0/go.mod h1:+e6BMRMPjBQoCw/WovYR9GLy2IU0z4Q77smOB1DraSg= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.0 h1:tC323YV77QdafeBr6LUhLDTsboyuyHLNRwAyCP44kGU= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.0/go.mod h1:SfLK1sgviHmbI+MozR9iDwDjL4cdCVZtahsjoR+z7wg= +github.com/aws/aws-sdk-go-v2/service/sts v1.45.0 h1:Pd6PNlp4t8PTXxqzstICl52Wsy78vpjFZ7PRUj44mJc= +github.com/aws/aws-sdk-go-v2/service/sts v1.45.0/go.mod h1:rmQ0TnHzuLPmabgjPcsywhsSOmaBDgzR4zvDxSPsGdg= +github.com/aws/smithy-go v1.27.4 h1:JQcphmBN4f0q/sPqXqROIItRNV/hy10cgu7CsFy616M= +github.com/aws/smithy-go v1.27.4/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= +github.com/ceph/go-ceph v0.40.0 h1:Wz9WOX6i73Hz74mpwhTO9S6IyX3eFPv88VUc7FRMPRk= +github.com/ceph/go-ceph v0.40.0/go.mod h1:1oFtT/x/4y+teLsNiogdd/Kj81Gmrw6JM5w1bWnZoc4= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= @@ -101,8 +139,8 @@ github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UV github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= -github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= -github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= @@ -138,8 +176,8 @@ golang.org/x/sync v0.12.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200930185726-fdedc70b468f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.31.0 h1:ioabZlmFYtWhL+TRYpcnNlLwhyxaM9kWTDEmfnprqik= -golang.org/x/sys v0.31.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= +golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY= +golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/term v0.30.0 h1:PQ39fJZ+mfadBm0y5WlL4vlM7Sx1Hgf13sMIY2+QS9Y= golang.org/x/term v0.30.0/go.mod h1:NYYFdzHoI5wRh/h5tDMdMqCqPJZEuNqVR5xJLd/n67g= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= diff --git a/hack/kind/manifests/operator-deployment.yaml b/hack/kind/manifests/operator-deployment.yaml index 95d47ff..15a3507 100644 --- a/hack/kind/manifests/operator-deployment.yaml +++ b/hack/kind/manifests/operator-deployment.yaml @@ -1,6 +1,8 @@ --- -# Dashboard credentials for local testing. Replace the values, or create the +# radosgw admin credentials for local testing. Replace the values, or create the # Secret out-of-band, before applying. Keys map 1:1 to the operator env vars. +# The access/secret key belong to an RGW user with admin caps +# (users=*, buckets=*, metadata=read). apiVersion: v1 kind: Secret metadata: @@ -8,13 +10,17 @@ metadata: namespace: cephrgw-system type: Opaque stringData: - CEPH_DASHBOARD_URL: "https://dashboard.ceph.unkin.net" - CEPH_DASHBOARD_USERNAME: "k8s-cephrgw-operator" - CEPH_DASHBOARD_PASSWORD: "change-me" - # Optional: the S3 endpoint written into credential Secrets for consumers. + # radosgw endpoint the operator talks to (Admin Ops + S3 APIs). + CEPH_RGW_ADMIN_ENDPOINT: "https://radosgw.service.consul:443" + CEPH_RGW_ACCESS_KEY: "change-me" + CEPH_RGW_SECRET_KEY: "change-me" + # The S3 endpoint written into credential Secrets for consumers (may be a + # public name that differs from the API endpoint above). CEPH_RGW_ENDPOINT: "https://s3.ceph.unkin.net" + # Optional: SigV4 credential-scope region (defaults to "default"). + # CEPH_RGW_REGION: "default" # Optional: set to "true" to skip TLS verification (dev only). - # CEPH_DASHBOARD_INSECURE: "true" + # CEPH_RGW_INSECURE: "true" --- apiVersion: apps/v1 kind: Deployment diff --git a/internal/ceph/buckets.go b/internal/ceph/buckets.go index 2759d8d..458e381 100644 --- a/internal/ceph/buckets.go +++ b/internal/ceph/buckets.go @@ -2,22 +2,25 @@ package ceph import ( "context" - "net/http" - "net/url" - "strconv" + "encoding/json" + "fmt" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/service/s3" + s3types "github.com/aws/aws-sdk-go-v2/service/s3/types" + "github.com/ceph/go-ceph/rgw/admin" ) -// BucketInfo is the subset of an RGW bucket record the operator consumes. -// Different Ceph releases name the id/name fields slightly differently, so the -// struct captures the known variants and Name/ID normalise them. +// BucketInfo is the subset of an RGW bucket record the operator consumes, as +// returned by the Admin Ops API (GET /admin/bucket). type BucketInfo struct { - Bucket string `json:"bucket"` - Bid string `json:"bid"` - ID string `json:"id"` - Owner string `json:"owner"` + Bucket string + Bid string + ID string + Owner string } -// Name returns the bucket name regardless of the field the dashboard used. +// Name returns the bucket name regardless of the field radosgw used. func (b *BucketInfo) Name() string { if b.Bucket != "" { return b.Bucket @@ -40,90 +43,160 @@ type CreateBucketSpec struct { LockYears *int32 } -type createBucketRequest struct { - Bucket string `json:"bucket"` - UID string `json:"uid"` - Zonegroup string `json:"zonegroup,omitempty"` - PlacementTarget string `json:"placement_target,omitempty"` - LockEnabled string `json:"lock_enabled"` - LockMode string `json:"lock_mode,omitempty"` - LockDays string `json:"lock_retention_period_days,omitempty"` - LockYears string `json:"lock_retention_period_years,omitempty"` -} - -// GetBucket fetches a bucket by name, returning an *APIError with status 404 -// (see IsNotFound) when it does not exist. +// GetBucket fetches a bucket by name via the Admin Ops API, returning an error +// classified by IsNotFound (admin.ErrNoSuchBucket) when it does not exist. func (c *Client) GetBucket(ctx context.Context, name string) (*BucketInfo, error) { - var b BucketInfo - if err := c.do(ctx, http.MethodGet, "/api/rgw/bucket/"+url.PathEscape(name), nil, &b, ""); err != nil { + b, err := c.admin.GetBucketInfo(ctx, admin.Bucket{Bucket: name}) + if err != nil { return nil, err } - return &b, nil + return &BucketInfo{Bucket: b.Bucket, ID: b.ID, Owner: b.Owner}, nil } -// CreateBucket provisions a bucket owned by spec.OwnerUID. +// CreateBucket provisions a bucket owned by spec.OwnerUID. The Admin Ops API +// cannot create buckets, so the operator issues an S3 CreateBucket signed as the +// owner (which makes the owner the bucket owner directly). func (c *Client) CreateBucket(ctx context.Context, spec CreateBucketSpec) (*BucketInfo, error) { - req := createBucketRequest{ - Bucket: spec.Bucket, - UID: spec.OwnerUID, - Zonegroup: spec.Zonegroup, - PlacementTarget: spec.PlacementTarget, - LockEnabled: strconv.FormatBool(spec.LockEnabled), - LockMode: spec.LockMode, - } - if spec.LockDays != nil { - req.LockDays = strconv.Itoa(int(*spec.LockDays)) - } - if spec.LockYears != nil { - req.LockYears = strconv.Itoa(int(*spec.LockYears)) - } - var b BucketInfo - if err := c.do(ctx, http.MethodPost, "/api/rgw/bucket", req, &b, ""); err != nil { + owner, err := c.asOwner(ctx, spec.OwnerUID) + if err != nil { return nil, err } - return &b, nil -} -type setBucketRequest struct { - BucketID string `json:"bucket_id"` - UID string `json:"uid"` - VersioningState *string `json:"versioning_state,omitempty"` - BucketPolicy *string `json:"bucket_policy,omitempty"` - Tags *string `json:"tags,omitempty"` -} - -// SetBucketVersioning enables or suspends S3 versioning on a bucket. -func (c *Client) SetBucketVersioning(ctx context.Context, name, bucketID, ownerUID string, enabled bool) error { - state := "Suspended" - if enabled { - state = "Enabled" + input := &s3.CreateBucketInput{Bucket: aws.String(spec.Bucket)} + if loc := locationConstraint(spec.Zonegroup, spec.PlacementTarget); loc != "" { + input.CreateBucketConfiguration = &s3types.CreateBucketConfiguration{ + LocationConstraint: s3types.BucketLocationConstraint(loc), + } } - req := setBucketRequest{BucketID: bucketID, UID: ownerUID, VersioningState: &state} - return c.do(ctx, http.MethodPut, "/api/rgw/bucket/"+url.PathEscape(name), req, nil, "") + if spec.LockEnabled { + input.ObjectLockEnabledForBucket = aws.Bool(true) + } + + if _, err := c.s3.CreateBucket(ctx, input, owner); err != nil && !IsConflict(err) { + return nil, err + } + + // Apply a default object-lock retention when requested. + if spec.LockEnabled && spec.LockMode != "" && (spec.LockDays != nil || spec.LockYears != nil) { + if err := c.setObjectLockDefault(ctx, owner, spec); err != nil { + return nil, err + } + } + + return c.GetBucket(ctx, spec.Bucket) } -// SetBucketPolicy replaces the S3 bucket policy. An empty policy string asks the -// dashboard to clear it; not every release honours clearing, so callers should -// treat a clear as best-effort. +// SetBucketVersioning enables or suspends S3 versioning on a bucket. bucketID is +// unused (kept for call-site stability). +func (c *Client) SetBucketVersioning(ctx context.Context, name, bucketID, ownerUID string, enabled bool) error { + owner, err := c.asOwner(ctx, ownerUID) + if err != nil { + return err + } + status := s3types.BucketVersioningStatusSuspended + if enabled { + status = s3types.BucketVersioningStatusEnabled + } + _, err = c.s3.PutBucketVersioning(ctx, &s3.PutBucketVersioningInput{ + Bucket: aws.String(name), + VersioningConfiguration: &s3types.VersioningConfiguration{Status: status}, + }, owner) + return err +} + +// SetBucketPolicy replaces the S3 bucket policy. An empty policy clears it. +// bucketID is unused (kept for call-site stability). func (c *Client) SetBucketPolicy(ctx context.Context, name, bucketID, ownerUID, policy string) error { - req := setBucketRequest{BucketID: bucketID, UID: ownerUID, BucketPolicy: &policy} - return c.do(ctx, http.MethodPut, "/api/rgw/bucket/"+url.PathEscape(name), req, nil, "") + owner, err := c.asOwner(ctx, ownerUID) + if err != nil { + return err + } + if policy == "" { + _, err := c.s3.DeleteBucketPolicy(ctx, &s3.DeleteBucketPolicyInput{Bucket: aws.String(name)}, owner) + if IsNotFound(err) { + return nil + } + return err + } + _, err = c.s3.PutBucketPolicy(ctx, &s3.PutBucketPolicyInput{ + Bucket: aws.String(name), + Policy: aws.String(policy), + }, owner) + return err } -// SetBucketTags replaces the bucket tag set. tagsJSON is the RGW/S3 tag JSON -// (a list of {"Key","Value"} objects). +// SetBucketTags replaces the bucket tag set. tagsJSON is the JSON produced by +// BuildTagJSON (a list of {"Key","Value"} objects). bucketID is unused (kept for +// call-site stability). func (c *Client) SetBucketTags(ctx context.Context, name, bucketID, ownerUID, tagsJSON string) error { - req := setBucketRequest{BucketID: bucketID, UID: ownerUID, Tags: &tagsJSON} - return c.do(ctx, http.MethodPut, "/api/rgw/bucket/"+url.PathEscape(name), req, nil, "") + owner, err := c.asOwner(ctx, ownerUID) + if err != nil { + return err + } + + type tag struct { + Key string `json:"Key"` + Value string `json:"Value"` + } + var tags []tag + if tagsJSON != "" { + if err := json.Unmarshal([]byte(tagsJSON), &tags); err != nil { + return fmt.Errorf("ceph: parse bucket tags: %w", err) + } + } + if len(tags) == 0 { + _, err := c.s3.DeleteBucketTagging(ctx, &s3.DeleteBucketTaggingInput{Bucket: aws.String(name)}, owner) + if IsNotFound(err) { + return nil + } + return err + } + + tagSet := make([]s3types.Tag, 0, len(tags)) + for _, t := range tags { + tagSet = append(tagSet, s3types.Tag{Key: aws.String(t.Key), Value: aws.String(t.Value)}) + } + _, err = c.s3.PutBucketTagging(ctx, &s3.PutBucketTaggingInput{ + Bucket: aws.String(name), + Tagging: &s3types.Tagging{TagSet: tagSet}, + }, owner) + return err } -// DeleteBucket removes a bucket. When purge is true its objects are deleted too; -// otherwise deletion of a non-empty bucket fails. A 404 is treated as success. +// DeleteBucket removes a bucket via the Admin Ops API. When purge is true its +// objects are deleted too. A NoSuchBucket response is treated as success. func (c *Client) DeleteBucket(ctx context.Context, name string, purge bool) error { - path := "/api/rgw/bucket/" + url.PathEscape(name) + "?purge_objects=" + strconv.FormatBool(purge) - err := c.do(ctx, http.MethodDelete, path, nil, nil, "") + err := c.admin.RemoveBucket(ctx, admin.Bucket{Bucket: name, PurgeObject: &purge}) if IsNotFound(err) { return nil } return err } + +// setObjectLockDefault sets the bucket's default object-lock retention. +func (c *Client) setObjectLockDefault(ctx context.Context, owner func(*s3.Options), spec CreateBucketSpec) error { + _, err := c.s3.PutObjectLockConfiguration(ctx, &s3.PutObjectLockConfigurationInput{ + Bucket: aws.String(spec.Bucket), + ObjectLockConfiguration: &s3types.ObjectLockConfiguration{ + ObjectLockEnabled: s3types.ObjectLockEnabledEnabled, + Rule: &s3types.ObjectLockRule{ + DefaultRetention: &s3types.DefaultRetention{ + Mode: s3types.ObjectLockRetentionMode(spec.LockMode), + Days: spec.LockDays, + Years: spec.LockYears, + }, + }, + }, + }, owner) + return err +} + +// locationConstraint renders the RGW LocationConstraint from a zonegroup and +// placement target (":"), or "" for default placement. +func locationConstraint(zonegroup, placement string) string { + loc := zonegroup + if placement != "" { + loc = zonegroup + ":" + placement + } + return loc +} diff --git a/internal/ceph/client.go b/internal/ceph/client.go index 8203d06..0bd985c 100644 --- a/internal/ceph/client.go +++ b/internal/ceph/client.go @@ -1,39 +1,50 @@ -// Package ceph is a small client for the Ceph manager dashboard REST API, -// scoped to the RGW (S3) user and bucket endpoints the operator needs. +// Package ceph is a small client for the Ceph RGW (radosgw) admin and S3 APIs, +// scoped to the user, bucket and policy operations the operator needs. // -// The dashboard authenticates with a username/password to POST /api/auth, which -// returns a bearer (JWT) token. The client caches that token and transparently -// re-authenticates when the server returns 401 (expired/invalid token). +// It talks directly to radosgw (e.g. https://radosgw.service.consul:443) rather +// than the manager dashboard, via two native Go libraries: +// +// - github.com/ceph/go-ceph/rgw/admin drives the RGW Admin Ops API +// (/admin/...), signed with the operator's admin access/secret key, to +// manage users, keys, quotas and bucket info/removal. +// - github.com/aws/aws-sdk-go-v2/service/s3 drives the S3 API (/), signed as +// the bucket's owner, to create buckets and set versioning, tagging, policy +// and object lock — operations the Admin Ops API does not expose. package ceph import ( - "bytes" "context" "crypto/tls" "crypto/x509" - "encoding/json" "errors" "fmt" - "io" "net/http" "strings" "sync" "time" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/credentials" + "github.com/aws/aws-sdk-go-v2/service/s3" + smithy "github.com/aws/smithy-go" + awshttp "github.com/aws/smithy-go/transport/http" + "github.com/ceph/go-ceph/rgw/admin" ) -// defaultAccept is the versioned media type the Ceph dashboard requires on its -// RGW endpoints. The dashboard rejects requests without a matching version. -const defaultAccept = "application/vnd.ceph.api.v1.0+json" - -// Config configures a dashboard Client. +// Config configures a Client. type Config struct { - // BaseURL is the dashboard root, e.g. https://dashboard.ceph.unkin.net. - BaseURL string - // Username / Password authenticate to POST /api/auth. The account needs the - // rgw-manager role (or admin) on the dashboard. - Username string - Password string - // CACert is an optional PEM bundle used to verify the dashboard TLS cert. + // Endpoint is the radosgw root, e.g. https://radosgw.service.consul:443. + Endpoint string + // AccessKey / SecretKey are the S3 credentials of an RGW user holding the + // admin caps the operator needs (users=*, buckets=*). + AccessKey string + SecretKey string + // Region is the SigV4 credential-scope region used for S3 requests. radosgw + // verifies the signature against whatever region the client used, so any + // consistent value works; defaults to "default". (The go-ceph admin client + // always signs its own requests with region "default".) + Region string + // CACert is an optional PEM bundle used to verify the radosgw TLS cert. CACert []byte // Insecure disables TLS verification (not recommended). Insecure bool @@ -41,49 +52,80 @@ type Config struct { Timeout time.Duration } -// Client talks to the Ceph dashboard API. It is safe for concurrent use. +// Client talks to radosgw. It is safe for concurrent use. type Client struct { - base string - user string - pass string - http *http.Client + admin *admin.API + s3 *s3.Client + region string - mu sync.Mutex - token string + // keyCache memoises owner uid -> S3 credentials (via the Admin Ops API) so + // per-owner S3 calls do not re-fetch keys on every reconcile. + mu sync.Mutex + keyCache map[string]aws.CredentialsProvider } -// APIError is returned for any non-2xx dashboard response. -type APIError struct { - Status int - Method string - Path string - Body string +// notFoundCodes and conflictCodes classify RGW/S3 error codes that surface only +// as a generic smithy.APIError (i.e. not a modeled S3 error type). +var notFoundCodes = map[string]bool{ + "NoSuchUser": true, "NoSuchBucket": true, "NoSuchKey": true, + "NoSuchBucketPolicy": true, "NoSuchTagSet": true, + "NoSuchTagSetError": true, "NotFound": true, } -func (e *APIError) Error() string { - return fmt.Sprintf("ceph dashboard %s %s: status %d: %s", e.Method, e.Path, e.Status, e.Body) +var conflictCodes = map[string]bool{ + "BucketAlreadyExists": true, "BucketAlreadyOwnedByYou": true, "UserAlreadyExists": true, } -// IsNotFound reports whether err is a 404 from the dashboard. +// IsNotFound reports whether err represents a missing user, bucket, key, policy +// or tag set, on either the admin or the S3 path. func IsNotFound(err error) bool { - var a *APIError - return errors.As(err, &a) && a.Status == http.StatusNotFound + if err == nil { + return false + } + if errors.Is(err, admin.ErrNoSuchUser) || errors.Is(err, admin.ErrNoSuchBucket) || + errors.Is(err, admin.ErrNoSuchKey) || errors.Is(err, admin.ErrNoSuchObject) { + return true + } + var apiErr smithy.APIError + if errors.As(err, &apiErr) && notFoundCodes[apiErr.ErrorCode()] { + return true + } + var respErr *awshttp.ResponseError + if errors.As(err, &respErr) && respErr.HTTPStatusCode() == http.StatusNotFound { + return true + } + return false } -// IsConflict reports whether err is a 409 from the dashboard. +// IsConflict reports whether err represents an already-exists conflict on either +// the admin or the S3 path. func IsConflict(err error) bool { - var a *APIError - return errors.As(err, &a) && a.Status == http.StatusConflict + if err == nil { + return false + } + if errors.Is(err, admin.ErrUserExists) || errors.Is(err, admin.ErrEmailExists) || + errors.Is(err, admin.ErrKeyExists) || errors.Is(err, admin.ErrBucketNotEmpty) { + return true + } + var apiErr smithy.APIError + if errors.As(err, &apiErr) && conflictCodes[apiErr.ErrorCode()] { + return true + } + var respErr *awshttp.ResponseError + if errors.As(err, &respErr) && respErr.HTTPStatusCode() == http.StatusConflict { + return true + } + return false } // NewClient validates cfg and builds a Client. func NewClient(cfg Config) (*Client, error) { - base := strings.TrimRight(cfg.BaseURL, "/") - if base == "" { - return nil, fmt.Errorf("ceph: dashboard base URL is required") + endpoint := strings.TrimRight(cfg.Endpoint, "/") + if endpoint == "" { + return nil, fmt.Errorf("ceph: radosgw endpoint is required") } - if cfg.Username == "" || cfg.Password == "" { - return nil, fmt.Errorf("ceph: dashboard username and password are required") + if cfg.AccessKey == "" || cfg.SecretKey == "" { + return nil, fmt.Errorf("ceph: radosgw admin access and secret key are required") } tlsCfg := &tls.Config{InsecureSkipVerify: cfg.Insecure} //nolint:gosec // opt-in via config @@ -99,124 +141,84 @@ func NewClient(cfg Config) (*Client, error) { if timeout == 0 { timeout = 30 * time.Second } + region := cfg.Region + if region == "" { + region = "default" + } + + httpClient := &http.Client{ + Timeout: timeout, + Transport: &http.Transport{TLSClientConfig: tlsCfg}, + } + + adminAPI, err := admin.New(endpoint, cfg.AccessKey, cfg.SecretKey, httpClient) + if err != nil { + return nil, fmt.Errorf("ceph: build admin client: %w", err) + } + + s3Client := s3.New(s3.Options{ + Region: region, + Credentials: credentials.NewStaticCredentialsProvider(cfg.AccessKey, cfg.SecretKey, ""), + HTTPClient: httpClient, + BaseEndpoint: aws.String(endpoint), + // radosgw serves buckets path-style, not virtual-host style. + UsePathStyle: true, + // radosgw (pre-Reef backports) rejects the SDK's default CRC32 / + // aws-chunked integrity protections; only send checksums when the API + // requires them. + RequestChecksumCalculation: aws.RequestChecksumCalculationWhenRequired, + ResponseChecksumValidation: aws.ResponseChecksumValidationWhenRequired, + }) return &Client{ - base: base, - user: cfg.Username, - pass: cfg.Password, - http: &http.Client{ - Timeout: timeout, - Transport: &http.Transport{TLSClientConfig: tlsCfg}, - }, + admin: adminAPI, + s3: s3Client, + region: region, + keyCache: map[string]aws.CredentialsProvider{}, }, nil } -// do performs an authenticated request, decoding a 2xx JSON body into out (when -// non-nil). On a 401 it drops the cached token, re-authenticates, and retries -// once. accept overrides the Accept header version when non-empty. -func (c *Client) do(ctx context.Context, method, path string, body, out any, accept string) error { - if accept == "" { - accept = defaultAccept - } - tok, err := c.ensureToken(ctx) - if err != nil { - return err - } - status, err := c.execute(ctx, method, path, body, accept, tok, out) - if status == http.StatusUnauthorized { - c.clearToken() - tok, err = c.ensureToken(ctx) +// asOwner returns a per-call S3 option that signs the request as the RGW user +// uid, looking up (and caching) the user's first key pair via the Admin Ops API. +// Signing S3 sub-resource operations as the bucket owner (rather than the admin +// user) makes the owner the bucket owner directly and keeps RGW's per-user S3 +// authorization intact. +func (c *Client) asOwner(ctx context.Context, uid string) (func(*s3.Options), error) { + c.mu.Lock() + provider, ok := c.keyCache[uid] + c.mu.Unlock() + if !ok { + user, err := c.GetUser(ctx, uid) if err != nil { - return err + return nil, err } - _, err = c.execute(ctx, method, path, body, accept, tok, out) + key, has := user.S3Key() + if !has { + return nil, fmt.Errorf("ceph: user %s has no S3 keys to sign bucket operations", uid) + } + provider = credentials.NewStaticCredentialsProvider(key.AccessKey, key.SecretKey, "") + c.mu.Lock() + c.keyCache[uid] = provider + c.mu.Unlock() } - return err + return func(o *s3.Options) { o.Credentials = provider }, nil } -func (c *Client) ensureToken(ctx context.Context) (string, error) { +// forgetIdentity drops any cached S3 credentials for uid, e.g. after its keys +// may have changed. +func (c *Client) forgetIdentity(uid string) { c.mu.Lock() - defer c.mu.Unlock() - if c.token != "" { - return c.token, nil - } - tok, err := c.login(ctx) - if err != nil { - return "", err - } - c.token = tok - return tok, nil -} - -func (c *Client) clearToken() { - c.mu.Lock() - c.token = "" + delete(c.keyCache, uid) c.mu.Unlock() } -func (c *Client) login(ctx context.Context) (string, error) { - var out struct { - Token string `json:"token"` - } - payload := map[string]string{"username": c.user, "password": c.pass} - if _, err := c.execute(ctx, http.MethodPost, "/api/auth", payload, defaultAccept, "", &out); err != nil { - return "", fmt.Errorf("dashboard login failed: %w", err) - } - if out.Token == "" { - return "", fmt.Errorf("dashboard login returned no token") - } - return out.Token, nil -} - -// execute runs a single request and returns the HTTP status. A non-2xx status -// yields an *APIError. token is sent as a bearer when non-empty. -func (c *Client) execute(ctx context.Context, method, path string, body any, accept, token string, out any) (int, error) { - var reader io.Reader - if body != nil { - b, err := json.Marshal(body) - if err != nil { - return 0, fmt.Errorf("marshal request body: %w", err) - } - reader = bytes.NewReader(b) - } - - req, err := http.NewRequestWithContext(ctx, method, c.base+path, reader) - if err != nil { - return 0, err - } - req.Header.Set("Accept", accept) - if body != nil { - req.Header.Set("Content-Type", "application/json") - } - if token != "" { - req.Header.Set("Authorization", "Bearer "+token) - } - - resp, err := c.http.Do(req) - if err != nil { - return 0, err - } - defer resp.Body.Close() - - data, _ := io.ReadAll(resp.Body) - if resp.StatusCode < 200 || resp.StatusCode >= 300 { - return resp.StatusCode, &APIError{ - Status: resp.StatusCode, - Method: method, - Path: path, - Body: strings.TrimSpace(string(data)), - } - } - if out != nil && len(data) > 0 { - if err := json.Unmarshal(data, out); err != nil { - return resp.StatusCode, fmt.Errorf("decode %s %s response: %w", method, path, err) - } - } - return resp.StatusCode, nil -} - -// Ping verifies connectivity and credentials by authenticating. +// Ping verifies connectivity and that the admin credentials sign correctly. It +// asks the Admin Ops API for a sentinel user: a NoSuchUser answer still proves +// the request authenticated, so only transport/auth errors fail the check. func (c *Client) Ping(ctx context.Context) error { - _, err := c.ensureToken(ctx) + _, err := c.GetUser(ctx, "cephrgw-operator-ping-nonexistent") + if err == nil || IsNotFound(err) { + return nil + } return err } diff --git a/internal/ceph/client_test.go b/internal/ceph/client_test.go new file mode 100644 index 0000000..8595f46 --- /dev/null +++ b/internal/ceph/client_test.go @@ -0,0 +1,81 @@ +package ceph + +import ( + "testing" + + s3types "github.com/aws/aws-sdk-go-v2/service/s3/types" + smithy "github.com/aws/smithy-go" + "github.com/ceph/go-ceph/rgw/admin" +) + +func TestNewClientValidation(t *testing.T) { + cases := []struct { + name string + cfg Config + wantErr bool + }{ + {"ok", Config{Endpoint: "https://rgw:443", AccessKey: "a", SecretKey: "s"}, false}, + {"no endpoint", Config{AccessKey: "a", SecretKey: "s"}, true}, + {"no access key", Config{Endpoint: "https://rgw:443", SecretKey: "s"}, true}, + {"no secret key", Config{Endpoint: "https://rgw:443", AccessKey: "a"}, true}, + {"bad ca", Config{Endpoint: "https://rgw:443", AccessKey: "a", SecretKey: "s", CACert: []byte("not pem")}, true}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + _, err := NewClient(tc.cfg) + if (err != nil) != tc.wantErr { + t.Fatalf("NewClient err=%v wantErr=%v", err, tc.wantErr) + } + }) + } +} + +func TestIsNotFound(t *testing.T) { + cases := []struct { + name string + err error + want bool + }{ + {"nil", nil, false}, + {"admin no such user", admin.ErrNoSuchUser, true}, + {"admin no such bucket", admin.ErrNoSuchBucket, true}, + {"admin no such key", admin.ErrNoSuchKey, true}, + {"s3 no such bucket", &s3types.NoSuchBucket{}, true}, + {"s3 no such key", &s3types.NoSuchKey{}, true}, + {"generic no such bucket policy", &smithy.GenericAPIError{Code: "NoSuchBucketPolicy"}, true}, + {"generic no such tag set", &smithy.GenericAPIError{Code: "NoSuchTagSet"}, true}, + {"admin user exists is not notfound", admin.ErrUserExists, false}, + {"unrelated", &smithy.GenericAPIError{Code: "AccessDenied"}, false}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + if got := IsNotFound(tc.err); got != tc.want { + t.Errorf("IsNotFound(%v)=%v want %v", tc.err, got, tc.want) + } + }) + } +} + +func TestIsConflict(t *testing.T) { + cases := []struct { + name string + err error + want bool + }{ + {"nil", nil, false}, + {"admin user exists", admin.ErrUserExists, true}, + {"admin bucket not empty", admin.ErrBucketNotEmpty, true}, + {"s3 bucket already owned by you", &s3types.BucketAlreadyOwnedByYou{}, true}, + {"s3 bucket already exists", &s3types.BucketAlreadyExists{}, true}, + {"generic bucket already exists", &smithy.GenericAPIError{Code: "BucketAlreadyExists"}, true}, + {"admin no such user is not conflict", admin.ErrNoSuchUser, false}, + {"unrelated", &smithy.GenericAPIError{Code: "AccessDenied"}, false}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + if got := IsConflict(tc.err); got != tc.want { + t.Errorf("IsConflict(%v)=%v want %v", tc.err, got, tc.want) + } + }) + } +} diff --git a/internal/ceph/policy.go b/internal/ceph/policy.go index 7c433f7..0e14a9a 100644 --- a/internal/ceph/policy.go +++ b/internal/ceph/policy.go @@ -139,7 +139,8 @@ func sid(prefix, uid string) string { return b.String() } -// BuildTagJSON renders bucket tags in the JSON form the dashboard expects. +// BuildTagJSON renders bucket tags as a {Key,Value} JSON list, the intermediate +// form SetBucketTags parses and re-encodes into the S3 Tagging XML document. func BuildTagJSON(tags map[string]string) (string, error) { if len(tags) == 0 { return "", nil diff --git a/internal/ceph/users.go b/internal/ceph/users.go index d8dc459..bf9e3bc 100644 --- a/internal/ceph/users.go +++ b/internal/ceph/users.go @@ -2,25 +2,25 @@ package ceph import ( "context" - "net/http" - "net/url" + + "github.com/ceph/go-ceph/rgw/admin" ) // UserKey is an S3 access/secret key pair belonging to an RGW user. type UserKey struct { - User string `json:"user"` - AccessKey string `json:"access_key"` - SecretKey string `json:"secret_key"` + User string + AccessKey string + SecretKey string } // User is the subset of an RGW user record the operator consumes. type User struct { - UID string `json:"user_id"` - DisplayName string `json:"display_name"` - Email string `json:"email"` - MaxBuckets int `json:"max_buckets"` - Suspended int `json:"suspended"` - Keys []UserKey `json:"keys"` + UID string + DisplayName string + Email string + MaxBuckets int + Suspended int + Keys []UserKey } // S3Key returns the first access/secret key pair, if any. @@ -40,92 +40,87 @@ type UserSpec struct { Suspended bool } -type createUserRequest struct { - UID string `json:"uid"` - DisplayName string `json:"display_name"` - Email string `json:"email,omitempty"` - MaxBuckets *int32 `json:"max_buckets,omitempty"` - Suspended bool `json:"suspended"` - GenerateKey bool `json:"generate_key"` +// fromAdminUser converts a go-ceph admin.User into the subset the operator uses. +func fromAdminUser(u admin.User) *User { + out := &User{ + UID: u.ID, + DisplayName: u.DisplayName, + Email: u.Email, + MaxBuckets: derefInt(u.MaxBuckets), + Suspended: derefInt(u.Suspended), + } + for _, k := range u.Keys { + out.Keys = append(out.Keys, UserKey{User: k.User, AccessKey: k.AccessKey, SecretKey: k.SecretKey}) + } + return out } -type updateUserRequest struct { - DisplayName string `json:"display_name"` - Email string `json:"email,omitempty"` - MaxBuckets *int32 `json:"max_buckets,omitempty"` - Suspended bool `json:"suspended"` -} - -// GetUser fetches an RGW user by uid, returning an *APIError with status 404 -// (see IsNotFound) when it does not exist. +// GetUser fetches an RGW user by uid, returning an error classified by +// IsNotFound (admin.ErrNoSuchUser) when it does not exist. func (c *Client) GetUser(ctx context.Context, uid string) (*User, error) { - var u User - if err := c.do(ctx, http.MethodGet, "/api/rgw/user/"+url.PathEscape(uid), nil, &u, ""); err != nil { + u, err := c.admin.GetUser(ctx, admin.User{ID: uid}) + if err != nil { return nil, err } - return &u, nil + return fromAdminUser(u), nil } -// CreateUser creates an RGW user, asking the dashboard to generate an S3 key -// pair. The returned User carries the generated keys. +// CreateUser creates an RGW user, asking radosgw to generate an S3 key pair. The +// returned User carries the generated keys. func (c *Client) CreateUser(ctx context.Context, spec UserSpec) (*User, error) { - req := createUserRequest{ - UID: spec.UID, + u, err := c.admin.CreateUser(ctx, admin.User{ + ID: spec.UID, DisplayName: firstNonEmpty(spec.DisplayName, spec.UID), Email: spec.Email, - MaxBuckets: spec.MaxBuckets, - Suspended: spec.Suspended, - GenerateKey: true, - } - var u User - if err := c.do(ctx, http.MethodPost, "/api/rgw/user", req, &u, ""); err != nil { + MaxBuckets: int32PtrToIntPtr(spec.MaxBuckets), + Suspended: boolToIntPtr(spec.Suspended), + GenerateKey: boolPtr(true), + }) + if err != nil { return nil, err } - return &u, nil + c.forgetIdentity(spec.UID) + return fromAdminUser(u), nil } // UpdateUser reconciles the mutable attributes of an existing RGW user. func (c *Client) UpdateUser(ctx context.Context, spec UserSpec) (*User, error) { - req := updateUserRequest{ + u, err := c.admin.ModifyUser(ctx, admin.User{ + ID: spec.UID, DisplayName: firstNonEmpty(spec.DisplayName, spec.UID), Email: spec.Email, - MaxBuckets: spec.MaxBuckets, - Suspended: spec.Suspended, - } - var u User - if err := c.do(ctx, http.MethodPut, "/api/rgw/user/"+url.PathEscape(spec.UID), req, &u, ""); err != nil { + MaxBuckets: int32PtrToIntPtr(spec.MaxBuckets), + Suspended: boolToIntPtr(spec.Suspended), + }) + if err != nil { return nil, err } - return &u, nil + return fromAdminUser(u), nil } -// DeleteUser removes an RGW user. A 404 is treated as success. +// DeleteUser removes an RGW user. A NoSuchUser response is treated as success. func (c *Client) DeleteUser(ctx context.Context, uid string) error { - err := c.do(ctx, http.MethodDelete, "/api/rgw/user/"+url.PathEscape(uid), nil, nil, "") + err := c.admin.RemoveUser(ctx, admin.User{ID: uid}) + c.forgetIdentity(uid) if IsNotFound(err) { return nil } return err } -type quotaRequest struct { - QuotaType string `json:"quota_type"` - Enabled bool `json:"enabled"` - MaxSizeKb int64 `json:"max_size_kb"` - MaxObjects int64 `json:"max_objects"` -} - // SetUserQuota applies a quota to a user. quotaType is "user" or "bucket" (the // latter sets the per-bucket default for buckets the user owns). A nil or // negative limit means unlimited for that dimension. func (c *Client) SetUserQuota(ctx context.Context, uid, quotaType string, enabled bool, maxSizeBytes, maxObjects *int64) error { - req := quotaRequest{ + maxSize := valueOr(maxSizeBytes, -1) + maxObj := valueOr(maxObjects, -1) + return c.admin.SetUserQuota(ctx, admin.QuotaSpec{ + UID: uid, QuotaType: quotaType, - Enabled: enabled, - MaxSizeKb: bytesToKb(maxSizeBytes), - MaxObjects: valueOr(maxObjects, -1), - } - return c.do(ctx, http.MethodPut, "/api/rgw/user/"+url.PathEscape(uid)+"/quota", req, nil, "") + Enabled: &enabled, + MaxSize: &maxSize, + MaxObjects: &maxObj, + }) } func firstNonEmpty(vals ...string) string { @@ -137,16 +132,34 @@ func firstNonEmpty(vals ...string) string { return "" } -func bytesToKb(b *int64) int64 { - if b == nil || *b < 0 { - return -1 - } - return *b / 1024 -} - func valueOr(v *int64, fallback int64) int64 { if v == nil || *v < 0 { return fallback } return *v } + +func derefInt(p *int) int { + if p == nil { + return 0 + } + return *p +} + +func boolPtr(b bool) *bool { return &b } + +func int32PtrToIntPtr(p *int32) *int { + if p == nil { + return nil + } + v := int(*p) + return &v +} + +func boolToIntPtr(b bool) *int { + v := 0 + if b { + v = 1 + } + return &v +}