Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 619aa6751d | |||
| 4b0430f0df | |||
| 253105f914 | |||
| 466514063a | |||
| ab21378f8f |
@@ -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
|
||||
|
||||
@@ -31,16 +40,36 @@ user for that grant and writes its keys into `spec.secretName` (default
|
||||
`<name>-rgw`). If `userRef` names an existing `ObjectStoreUser`, that user's own
|
||||
credential Secret is reused and only the policy is extended.
|
||||
|
||||
#### Fine-grained grants
|
||||
|
||||
The level is the ergonomic default; four optional fields on `BucketAccess`
|
||||
refine it (see `config/samples/04-access-fine-grained.yaml`):
|
||||
|
||||
- `spec.paths` — scope object access to key prefixes; each becomes the resource
|
||||
`<bucket>/<prefix>*`. The bucket-level `ListBucket` still spans the whole
|
||||
bucket.
|
||||
- `spec.actions` — grant exactly these S3 actions instead of the level's set (on
|
||||
the bucket and its, optionally prefixed, objects).
|
||||
- `spec.conditions` — `sourceIPs` (an `aws:SourceIp` CIDR allowlist) and
|
||||
`secureTransportOnly` (require TLS).
|
||||
- `spec.rawStatements` — an escape hatch of raw S3 policy statements
|
||||
(`effect`/`actions`/`resources`/`conditions`) merged for this grant's
|
||||
principal. When set, `level`, `actions`, `paths` and `conditions` are ignored;
|
||||
resources without an `arn:` prefix are treated as bucket-relative key prefixes.
|
||||
|
||||
RGW honours S3 bucket policy on **Reef 18.2+ / Squid**; condition-key support is
|
||||
a subset of AWS, so validate exotic conditions against your cluster.
|
||||
|
||||
The `Bucket` controller renders the policy as the **union of every ready
|
||||
`BucketAccess`** that targets it, so the result is convergent regardless of the
|
||||
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 +79,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 +87,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 +134,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.
|
||||
|
||||
@@ -45,6 +45,72 @@ type BucketAccessSpec struct {
|
||||
// dedicated user it creates (UserRef empty). Defaults to "<name>-rgw".
|
||||
// +optional
|
||||
SecretName string `json:"secretName,omitempty"`
|
||||
|
||||
// Paths optionally scopes object-level access to these key prefixes within
|
||||
// the bucket; each becomes the resource "<bucket>/<prefix>*". Empty grants
|
||||
// the whole bucket. The bucket-level ListBucket action always applies to the
|
||||
// whole bucket. Ignored when RawStatements is set.
|
||||
// +optional
|
||||
Paths []string `json:"paths,omitempty"`
|
||||
|
||||
// Actions optionally overrides the S3 actions granted by Level. When set,
|
||||
// exactly these actions are granted, on the bucket and its (optionally
|
||||
// prefixed) objects. Ignored when RawStatements is set.
|
||||
// +optional
|
||||
Actions []string `json:"actions,omitempty"`
|
||||
|
||||
// Conditions optionally restricts when the grant applies (e.g. source IPs,
|
||||
// TLS required). Ignored when RawStatements is set.
|
||||
// +optional
|
||||
Conditions *AccessConditions `json:"conditions,omitempty"`
|
||||
|
||||
// RawStatements is an escape hatch for arbitrary S3 policy statements, merged
|
||||
// into the bucket policy for this grant's principal. When set, Level,
|
||||
// Actions, Paths and Conditions on this object are ignored; the operator only
|
||||
// fills in the Principal (this grant's user) when a statement omits one.
|
||||
// +optional
|
||||
RawStatements []PolicyStatement `json:"rawStatements,omitempty"`
|
||||
}
|
||||
|
||||
// AccessConditions restricts when a grant applies. Each field maps to an S3
|
||||
// policy condition and, when several are set, all must hold (they are AND'd).
|
||||
type AccessConditions struct {
|
||||
// SourceIPs restricts the grant to requests from these CIDRs (or single
|
||||
// addresses), via the S3 aws:SourceIp condition.
|
||||
// +optional
|
||||
SourceIPs []string `json:"sourceIPs,omitempty"`
|
||||
|
||||
// SecureTransportOnly requires the request to use TLS, via the S3
|
||||
// aws:SecureTransport condition.
|
||||
// +optional
|
||||
SecureTransportOnly bool `json:"secureTransportOnly,omitempty"`
|
||||
}
|
||||
|
||||
// PolicyStatement is a raw S3 bucket-policy statement, exposed for grants that
|
||||
// need control beyond Level/Actions/Paths/Conditions.
|
||||
type PolicyStatement struct {
|
||||
// Sid is an optional statement id. The operator derives one when empty.
|
||||
// +optional
|
||||
Sid string `json:"sid,omitempty"`
|
||||
|
||||
// Effect is Allow or Deny. Defaults to Allow.
|
||||
// +kubebuilder:validation:Enum=Allow;Deny
|
||||
// +kubebuilder:default=Allow
|
||||
// +optional
|
||||
Effect string `json:"effect,omitempty"`
|
||||
|
||||
// Actions are the S3 actions the statement covers (e.g. s3:GetObject).
|
||||
Actions []string `json:"actions"`
|
||||
|
||||
// Resources are S3 resource ARNs, or bucket-relative key prefixes when they
|
||||
// do not start with "arn:". Empty means the whole bucket and its objects.
|
||||
// +optional
|
||||
Resources []string `json:"resources,omitempty"`
|
||||
|
||||
// Conditions is the raw S3 condition block: operator -> condition key ->
|
||||
// values, e.g. {"IpAddress": {"aws:SourceIp": ["10.0.0.0/8"]}}.
|
||||
// +optional
|
||||
Conditions map[string]map[string][]string `json:"conditions,omitempty"`
|
||||
}
|
||||
|
||||
// BucketAccessStatus reports observed grant state.
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -9,6 +9,26 @@ import (
|
||||
runtime "k8s.io/apimachinery/pkg/runtime"
|
||||
)
|
||||
|
||||
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
|
||||
func (in *AccessConditions) DeepCopyInto(out *AccessConditions) {
|
||||
*out = *in
|
||||
if in.SourceIPs != nil {
|
||||
in, out := &in.SourceIPs, &out.SourceIPs
|
||||
*out = make([]string, len(*in))
|
||||
copy(*out, *in)
|
||||
}
|
||||
}
|
||||
|
||||
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new AccessConditions.
|
||||
func (in *AccessConditions) DeepCopy() *AccessConditions {
|
||||
if in == nil {
|
||||
return nil
|
||||
}
|
||||
out := new(AccessConditions)
|
||||
in.DeepCopyInto(out)
|
||||
return out
|
||||
}
|
||||
|
||||
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
|
||||
func (in *Bucket) DeepCopyInto(out *Bucket) {
|
||||
*out = *in
|
||||
@@ -41,7 +61,7 @@ func (in *BucketAccess) DeepCopyInto(out *BucketAccess) {
|
||||
*out = *in
|
||||
out.TypeMeta = in.TypeMeta
|
||||
in.ObjectMeta.DeepCopyInto(&out.ObjectMeta)
|
||||
out.Spec = in.Spec
|
||||
in.Spec.DeepCopyInto(&out.Spec)
|
||||
in.Status.DeepCopyInto(&out.Status)
|
||||
}
|
||||
|
||||
@@ -98,6 +118,28 @@ func (in *BucketAccessList) DeepCopyObject() runtime.Object {
|
||||
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
|
||||
func (in *BucketAccessSpec) DeepCopyInto(out *BucketAccessSpec) {
|
||||
*out = *in
|
||||
if in.Paths != nil {
|
||||
in, out := &in.Paths, &out.Paths
|
||||
*out = make([]string, len(*in))
|
||||
copy(*out, *in)
|
||||
}
|
||||
if in.Actions != nil {
|
||||
in, out := &in.Actions, &out.Actions
|
||||
*out = make([]string, len(*in))
|
||||
copy(*out, *in)
|
||||
}
|
||||
if in.Conditions != nil {
|
||||
in, out := &in.Conditions, &out.Conditions
|
||||
*out = new(AccessConditions)
|
||||
(*in).DeepCopyInto(*out)
|
||||
}
|
||||
if in.RawStatements != nil {
|
||||
in, out := &in.RawStatements, &out.RawStatements
|
||||
*out = make([]PolicyStatement, len(*in))
|
||||
for i := range *in {
|
||||
(*in)[i].DeepCopyInto(&(*out)[i])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new BucketAccessSpec.
|
||||
@@ -349,6 +391,58 @@ func (in *ObjectStoreUserStatus) DeepCopy() *ObjectStoreUserStatus {
|
||||
return out
|
||||
}
|
||||
|
||||
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
|
||||
func (in *PolicyStatement) DeepCopyInto(out *PolicyStatement) {
|
||||
*out = *in
|
||||
if in.Actions != nil {
|
||||
in, out := &in.Actions, &out.Actions
|
||||
*out = make([]string, len(*in))
|
||||
copy(*out, *in)
|
||||
}
|
||||
if in.Resources != nil {
|
||||
in, out := &in.Resources, &out.Resources
|
||||
*out = make([]string, len(*in))
|
||||
copy(*out, *in)
|
||||
}
|
||||
if in.Conditions != nil {
|
||||
in, out := &in.Conditions, &out.Conditions
|
||||
*out = make(map[string]map[string][]string, len(*in))
|
||||
for key, val := range *in {
|
||||
var outVal map[string][]string
|
||||
if val == nil {
|
||||
(*out)[key] = nil
|
||||
} else {
|
||||
inVal := (*in)[key]
|
||||
in, out := &inVal, &outVal
|
||||
*out = make(map[string][]string, len(*in))
|
||||
for key, val := range *in {
|
||||
var outVal []string
|
||||
if val == nil {
|
||||
(*out)[key] = nil
|
||||
} else {
|
||||
inVal := (*in)[key]
|
||||
in, out := &inVal, &outVal
|
||||
*out = make([]string, len(*in))
|
||||
copy(*out, *in)
|
||||
}
|
||||
(*out)[key] = outVal
|
||||
}
|
||||
}
|
||||
(*out)[key] = outVal
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new PolicyStatement.
|
||||
func (in *PolicyStatement) DeepCopy() *PolicyStatement {
|
||||
if in == nil {
|
||||
return nil
|
||||
}
|
||||
out := new(PolicyStatement)
|
||||
in.DeepCopyInto(out)
|
||||
return out
|
||||
}
|
||||
|
||||
// DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil.
|
||||
func (in *Quota) DeepCopyInto(out *Quota) {
|
||||
*out = *in
|
||||
|
||||
+24
-14
@@ -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
|
||||
}
|
||||
|
||||
@@ -60,10 +60,36 @@ spec:
|
||||
operator provisions a dedicated user for this grant and writes its keys into
|
||||
a Secret; otherwise it grants an existing ObjectStoreUser.
|
||||
properties:
|
||||
actions:
|
||||
description: |-
|
||||
Actions optionally overrides the S3 actions granted by Level. When set,
|
||||
exactly these actions are granted, on the bucket and its (optionally
|
||||
prefixed) objects. Ignored when RawStatements is set.
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
bucketRef:
|
||||
description: BucketRef names the Bucket (in this namespace) to grant
|
||||
access to.
|
||||
type: string
|
||||
conditions:
|
||||
description: |-
|
||||
Conditions optionally restricts when the grant applies (e.g. source IPs,
|
||||
TLS required). Ignored when RawStatements is set.
|
||||
properties:
|
||||
secureTransportOnly:
|
||||
description: |-
|
||||
SecureTransportOnly requires the request to use TLS, via the S3
|
||||
aws:SecureTransport condition.
|
||||
type: boolean
|
||||
sourceIPs:
|
||||
description: |-
|
||||
SourceIPs restricts the grant to requests from these CIDRs (or single
|
||||
addresses), via the S3 aws:SourceIp condition.
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
type: object
|
||||
level:
|
||||
description: Level is the access level to grant.
|
||||
enum:
|
||||
@@ -71,6 +97,65 @@ spec:
|
||||
- read-write
|
||||
- full
|
||||
type: string
|
||||
paths:
|
||||
description: |-
|
||||
Paths optionally scopes object-level access to these key prefixes within
|
||||
the bucket; each becomes the resource "<bucket>/<prefix>*". Empty grants
|
||||
the whole bucket. The bucket-level ListBucket action always applies to the
|
||||
whole bucket. Ignored when RawStatements is set.
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
rawStatements:
|
||||
description: |-
|
||||
RawStatements is an escape hatch for arbitrary S3 policy statements, merged
|
||||
into the bucket policy for this grant's principal. When set, Level,
|
||||
Actions, Paths and Conditions on this object are ignored; the operator only
|
||||
fills in the Principal (this grant's user) when a statement omits one.
|
||||
items:
|
||||
description: |-
|
||||
PolicyStatement is a raw S3 bucket-policy statement, exposed for grants that
|
||||
need control beyond Level/Actions/Paths/Conditions.
|
||||
properties:
|
||||
actions:
|
||||
description: Actions are the S3 actions the statement covers
|
||||
(e.g. s3:GetObject).
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
conditions:
|
||||
additionalProperties:
|
||||
additionalProperties:
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
type: object
|
||||
description: |-
|
||||
Conditions is the raw S3 condition block: operator -> condition key ->
|
||||
values, e.g. {"IpAddress": {"aws:SourceIp": ["10.0.0.0/8"]}}.
|
||||
type: object
|
||||
effect:
|
||||
default: Allow
|
||||
description: Effect is Allow or Deny. Defaults to Allow.
|
||||
enum:
|
||||
- Allow
|
||||
- Deny
|
||||
type: string
|
||||
resources:
|
||||
description: |-
|
||||
Resources are S3 resource ARNs, or bucket-relative key prefixes when they
|
||||
do not start with "arn:". Empty means the whole bucket and its objects.
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
sid:
|
||||
description: Sid is an optional statement id. The operator derives
|
||||
one when empty.
|
||||
type: string
|
||||
required:
|
||||
- actions
|
||||
type: object
|
||||
type: array
|
||||
secretName:
|
||||
description: |-
|
||||
SecretName is the Secret the operator writes credentials into for the
|
||||
|
||||
@@ -52,7 +52,7 @@ spec:
|
||||
spec:
|
||||
description: |-
|
||||
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.
|
||||
properties:
|
||||
|
||||
+86
-1
@@ -61,10 +61,36 @@ spec:
|
||||
operator provisions a dedicated user for this grant and writes its keys into
|
||||
a Secret; otherwise it grants an existing ObjectStoreUser.
|
||||
properties:
|
||||
actions:
|
||||
description: |-
|
||||
Actions optionally overrides the S3 actions granted by Level. When set,
|
||||
exactly these actions are granted, on the bucket and its (optionally
|
||||
prefixed) objects. Ignored when RawStatements is set.
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
bucketRef:
|
||||
description: BucketRef names the Bucket (in this namespace) to grant
|
||||
access to.
|
||||
type: string
|
||||
conditions:
|
||||
description: |-
|
||||
Conditions optionally restricts when the grant applies (e.g. source IPs,
|
||||
TLS required). Ignored when RawStatements is set.
|
||||
properties:
|
||||
secureTransportOnly:
|
||||
description: |-
|
||||
SecureTransportOnly requires the request to use TLS, via the S3
|
||||
aws:SecureTransport condition.
|
||||
type: boolean
|
||||
sourceIPs:
|
||||
description: |-
|
||||
SourceIPs restricts the grant to requests from these CIDRs (or single
|
||||
addresses), via the S3 aws:SourceIp condition.
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
type: object
|
||||
level:
|
||||
description: Level is the access level to grant.
|
||||
enum:
|
||||
@@ -72,6 +98,65 @@ spec:
|
||||
- read-write
|
||||
- full
|
||||
type: string
|
||||
paths:
|
||||
description: |-
|
||||
Paths optionally scopes object-level access to these key prefixes within
|
||||
the bucket; each becomes the resource "<bucket>/<prefix>*". Empty grants
|
||||
the whole bucket. The bucket-level ListBucket action always applies to the
|
||||
whole bucket. Ignored when RawStatements is set.
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
rawStatements:
|
||||
description: |-
|
||||
RawStatements is an escape hatch for arbitrary S3 policy statements, merged
|
||||
into the bucket policy for this grant's principal. When set, Level,
|
||||
Actions, Paths and Conditions on this object are ignored; the operator only
|
||||
fills in the Principal (this grant's user) when a statement omits one.
|
||||
items:
|
||||
description: |-
|
||||
PolicyStatement is a raw S3 bucket-policy statement, exposed for grants that
|
||||
need control beyond Level/Actions/Paths/Conditions.
|
||||
properties:
|
||||
actions:
|
||||
description: Actions are the S3 actions the statement covers
|
||||
(e.g. s3:GetObject).
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
conditions:
|
||||
additionalProperties:
|
||||
additionalProperties:
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
type: object
|
||||
description: |-
|
||||
Conditions is the raw S3 condition block: operator -> condition key ->
|
||||
values, e.g. {"IpAddress": {"aws:SourceIp": ["10.0.0.0/8"]}}.
|
||||
type: object
|
||||
effect:
|
||||
default: Allow
|
||||
description: Effect is Allow or Deny. Defaults to Allow.
|
||||
enum:
|
||||
- Allow
|
||||
- Deny
|
||||
type: string
|
||||
resources:
|
||||
description: |-
|
||||
Resources are S3 resource ARNs, or bucket-relative key prefixes when they
|
||||
do not start with "arn:". Empty means the whole bucket and its objects.
|
||||
items:
|
||||
type: string
|
||||
type: array
|
||||
sid:
|
||||
description: Sid is an optional statement id. The operator derives
|
||||
one when empty.
|
||||
type: string
|
||||
required:
|
||||
- actions
|
||||
type: object
|
||||
type: array
|
||||
secretName:
|
||||
description: |-
|
||||
SecretName is the Secret the operator writes credentials into for the
|
||||
@@ -463,7 +548,7 @@ spec:
|
||||
spec:
|
||||
description: |-
|
||||
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.
|
||||
properties:
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
# Fine-grained grants. Each of these refines the coarse read-only/read-write/full
|
||||
# levels with prefix scoping, action overrides, conditions, or raw statements.
|
||||
|
||||
# 1. Prefix-scoped read-write: this workload may read/write objects only under
|
||||
# the "uploads/" and "tmp/" key prefixes (bucket-level ListBucket still spans
|
||||
# the whole bucket).
|
||||
apiVersion: ceph.unkin.net/v1alpha1
|
||||
kind: BucketAccess
|
||||
metadata:
|
||||
name: app-data-uploader
|
||||
namespace: default
|
||||
spec:
|
||||
bucketRef: app-data
|
||||
level: read-write
|
||||
secretName: app-data-uploader-rgw
|
||||
paths:
|
||||
- uploads/
|
||||
- tmp/
|
||||
---
|
||||
# 2. Read-only from inside the cluster only: restrict the grant to a source CIDR
|
||||
# and require TLS.
|
||||
apiVersion: ceph.unkin.net/v1alpha1
|
||||
kind: BucketAccess
|
||||
metadata:
|
||||
name: app-data-internal-ro
|
||||
namespace: default
|
||||
spec:
|
||||
bucketRef: app-data
|
||||
level: read-only
|
||||
secretName: app-data-internal-ro-rgw
|
||||
conditions:
|
||||
sourceIPs:
|
||||
- 10.0.0.0/8
|
||||
secureTransportOnly: true
|
||||
---
|
||||
# 3. Explicit action set: grant exactly these actions instead of a level's
|
||||
# canned set (level is still required but its actions are ignored).
|
||||
apiVersion: ceph.unkin.net/v1alpha1
|
||||
kind: BucketAccess
|
||||
metadata:
|
||||
name: app-data-getput
|
||||
namespace: default
|
||||
spec:
|
||||
bucketRef: app-data
|
||||
level: read-only
|
||||
secretName: app-data-getput-rgw
|
||||
actions:
|
||||
- s3:GetObject
|
||||
- s3:PutObject
|
||||
---
|
||||
# 4. Raw statements escape hatch: full control over the policy statement. Level,
|
||||
# actions, paths and conditions are ignored; the operator only injects the
|
||||
# Principal (this grant's user). Resources without an "arn:" prefix are
|
||||
# treated as bucket-relative key prefixes.
|
||||
apiVersion: ceph.unkin.net/v1alpha1
|
||||
kind: BucketAccess
|
||||
metadata:
|
||||
name: app-data-raw
|
||||
namespace: default
|
||||
spec:
|
||||
bucketRef: app-data
|
||||
level: read-only
|
||||
secretName: app-data-raw-rgw
|
||||
rawStatements:
|
||||
- effect: Allow
|
||||
actions:
|
||||
- s3:GetObject
|
||||
resources:
|
||||
- public/
|
||||
- effect: Deny
|
||||
actions:
|
||||
- s3:DeleteObject
|
||||
resources:
|
||||
- locked/
|
||||
conditions:
|
||||
Bool:
|
||||
aws:SecureTransport:
|
||||
- "false"
|
||||
+82
-110
@@ -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).
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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=
|
||||
|
||||
@@ -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
|
||||
|
||||
+147
-74
@@ -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 ("<zonegroup>:<placement>"), or "" for default placement.
|
||||
func locationConstraint(zonegroup, placement string) string {
|
||||
loc := zonegroup
|
||||
if placement != "" {
|
||||
loc = zonegroup + ":" + placement
|
||||
}
|
||||
return loc
|
||||
}
|
||||
|
||||
+152
-150
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
+169
-37
@@ -3,6 +3,7 @@ package ceph
|
||||
import (
|
||||
"encoding/json"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
)
|
||||
|
||||
@@ -14,10 +15,40 @@ const (
|
||||
LevelFull = "full"
|
||||
)
|
||||
|
||||
// Grant couples an RGW user id with the access level to grant it on a bucket.
|
||||
// GrantConditions restricts when a grant's statements apply. The zero value adds
|
||||
// no conditions.
|
||||
type GrantConditions struct {
|
||||
// SourceIPs restricts the grant to these CIDRs (S3 aws:SourceIp).
|
||||
SourceIPs []string
|
||||
// SecureTransportOnly requires TLS (S3 aws:SecureTransport).
|
||||
SecureTransportOnly bool
|
||||
}
|
||||
|
||||
// RawStatement is a caller-supplied S3 policy statement for a grant.
|
||||
type RawStatement struct {
|
||||
Sid string
|
||||
Effect string
|
||||
Actions []string
|
||||
Resources []string
|
||||
Condition map[string]map[string][]string
|
||||
}
|
||||
|
||||
// Grant couples an RGW user id with the access it should have on a bucket. The
|
||||
// simple form is a Level; Paths, Actions and Conditions refine it, and Raw
|
||||
// replaces it entirely with caller-supplied statements.
|
||||
type Grant struct {
|
||||
UID string
|
||||
Level string
|
||||
// Paths scopes object-level access to these key prefixes; empty = whole
|
||||
// bucket.
|
||||
Paths []string
|
||||
// Actions overrides the level's action set; empty = derive from Level.
|
||||
Actions []string
|
||||
// Conditions optionally restricts when the grant applies.
|
||||
Conditions *GrantConditions
|
||||
// Raw, when non-empty, replaces Level/Actions/Paths/Conditions with these
|
||||
// statements (the operator still fills in a Principal when one is omitted).
|
||||
Raw []RawStatement
|
||||
}
|
||||
|
||||
type policyDocument struct {
|
||||
@@ -26,11 +57,12 @@ type policyDocument struct {
|
||||
}
|
||||
|
||||
type policyStatement struct {
|
||||
Sid string `json:"Sid"`
|
||||
Effect string `json:"Effect"`
|
||||
Principal map[string][]string `json:"Principal"`
|
||||
Action []string `json:"Action"`
|
||||
Resource []string `json:"Resource"`
|
||||
Sid string `json:"Sid,omitempty"`
|
||||
Effect string `json:"Effect"`
|
||||
Principal map[string][]string `json:"Principal,omitempty"`
|
||||
Action []string `json:"Action"`
|
||||
Resource []string `json:"Resource"`
|
||||
Condition map[string]map[string][]string `json:"Condition,omitempty"`
|
||||
}
|
||||
|
||||
// bucket-level and object-level S3 actions per access level.
|
||||
@@ -68,7 +100,7 @@ var objectActions = map[string][]string{
|
||||
}
|
||||
|
||||
// BuildBucketPolicy renders a deterministic S3 bucket policy granting each
|
||||
// principal its requested level. It returns "" when there are no grants so the
|
||||
// principal its requested access. It returns "" when there are no grants so the
|
||||
// caller can clear the policy.
|
||||
func BuildBucketPolicy(bucket string, grants []Grant) (string, error) {
|
||||
if len(grants) == 0 {
|
||||
@@ -85,38 +117,10 @@ func BuildBucketPolicy(bucket string, grants []Grant) (string, error) {
|
||||
})
|
||||
|
||||
bucketARN := "arn:aws:s3:::" + bucket
|
||||
objectARN := bucketARN + "/*"
|
||||
|
||||
doc := policyDocument{Version: "2012-10-17"}
|
||||
for _, g := range sorted {
|
||||
principal := map[string][]string{"AWS": {"arn:aws:iam:::user/" + g.UID}}
|
||||
switch g.Level {
|
||||
case LevelFull:
|
||||
doc.Statement = append(doc.Statement, policyStatement{
|
||||
Sid: sid("full", g.UID),
|
||||
Effect: "Allow",
|
||||
Principal: principal,
|
||||
Action: []string{"s3:*"},
|
||||
Resource: []string{bucketARN, objectARN},
|
||||
})
|
||||
default:
|
||||
doc.Statement = append(doc.Statement,
|
||||
policyStatement{
|
||||
Sid: sid(g.Level+"-bkt", g.UID),
|
||||
Effect: "Allow",
|
||||
Principal: principal,
|
||||
Action: bucketActions[g.Level],
|
||||
Resource: []string{bucketARN},
|
||||
},
|
||||
policyStatement{
|
||||
Sid: sid(g.Level+"-obj", g.UID),
|
||||
Effect: "Allow",
|
||||
Principal: principal,
|
||||
Action: objectActions[g.Level],
|
||||
Resource: []string{objectARN},
|
||||
},
|
||||
)
|
||||
}
|
||||
doc.Statement = append(doc.Statement, statementsForGrant(bucketARN, g)...)
|
||||
}
|
||||
|
||||
b, err := json.Marshal(doc)
|
||||
@@ -126,6 +130,133 @@ func BuildBucketPolicy(bucket string, grants []Grant) (string, error) {
|
||||
return string(b), nil
|
||||
}
|
||||
|
||||
// statementsForGrant renders the policy statements for a single grant.
|
||||
func statementsForGrant(bucketARN string, g Grant) []policyStatement {
|
||||
principal := map[string][]string{"AWS": {"arn:aws:iam:::user/" + g.UID}}
|
||||
|
||||
if len(g.Raw) > 0 {
|
||||
out := make([]policyStatement, 0, len(g.Raw))
|
||||
for i, rs := range g.Raw {
|
||||
st := policyStatement{
|
||||
Sid: firstNonEmpty(rs.Sid, sid("raw", g.UID)+strconv.Itoa(i)),
|
||||
Effect: firstNonEmpty(rs.Effect, "Allow"),
|
||||
Principal: principal,
|
||||
Action: rs.Actions,
|
||||
Resource: resolveResources(bucketARN, rs.Resources),
|
||||
Condition: rs.Condition,
|
||||
}
|
||||
out = append(out, st)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
cond := buildCondition(g.Conditions)
|
||||
objectARNs := objectResources(bucketARN, g.Paths)
|
||||
|
||||
if len(g.Actions) > 0 {
|
||||
return []policyStatement{{
|
||||
Sid: sid("custom", g.UID),
|
||||
Effect: "Allow",
|
||||
Principal: principal,
|
||||
Action: g.Actions,
|
||||
Resource: append([]string{bucketARN}, objectARNs...),
|
||||
Condition: cond,
|
||||
}}
|
||||
}
|
||||
|
||||
if g.Level == LevelFull {
|
||||
return []policyStatement{{
|
||||
Sid: sid("full", g.UID),
|
||||
Effect: "Allow",
|
||||
Principal: principal,
|
||||
Action: []string{"s3:*"},
|
||||
Resource: append([]string{bucketARN}, objectARNs...),
|
||||
Condition: cond,
|
||||
}}
|
||||
}
|
||||
|
||||
return []policyStatement{
|
||||
{
|
||||
Sid: sid(g.Level+"-bkt", g.UID),
|
||||
Effect: "Allow",
|
||||
Principal: principal,
|
||||
Action: bucketActions[g.Level],
|
||||
Resource: []string{bucketARN},
|
||||
Condition: cond,
|
||||
},
|
||||
{
|
||||
Sid: sid(g.Level+"-obj", g.UID),
|
||||
Effect: "Allow",
|
||||
Principal: principal,
|
||||
Action: objectActions[g.Level],
|
||||
Resource: objectARNs,
|
||||
Condition: cond,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// objectResources renders the object-level resource ARNs for a grant: the whole
|
||||
// bucket ("<bucket>/*") when no paths are given, or one "<bucket>/<prefix>*" per
|
||||
// prefix (deduplicated and sorted for determinism).
|
||||
func objectResources(bucketARN string, paths []string) []string {
|
||||
if len(paths) == 0 {
|
||||
return []string{bucketARN + "/*"}
|
||||
}
|
||||
seen := map[string]struct{}{}
|
||||
out := make([]string, 0, len(paths))
|
||||
for _, p := range paths {
|
||||
p = strings.TrimPrefix(p, "/")
|
||||
arn := bucketARN + "/" + p + "*"
|
||||
if _, dup := seen[arn]; dup {
|
||||
continue
|
||||
}
|
||||
seen[arn] = struct{}{}
|
||||
out = append(out, arn)
|
||||
}
|
||||
sort.Strings(out)
|
||||
return out
|
||||
}
|
||||
|
||||
// resolveResources renders raw-statement resources: entries that already look
|
||||
// like ARNs pass through verbatim; bucket-relative prefixes become
|
||||
// "<bucket>/<prefix>*". An empty list defaults to the whole bucket and objects.
|
||||
func resolveResources(bucketARN string, resources []string) []string {
|
||||
if len(resources) == 0 {
|
||||
return []string{bucketARN, bucketARN + "/*"}
|
||||
}
|
||||
out := make([]string, 0, len(resources))
|
||||
for _, r := range resources {
|
||||
switch {
|
||||
case strings.HasPrefix(r, "arn:"):
|
||||
out = append(out, r)
|
||||
case r == "" || r == "/":
|
||||
out = append(out, bucketARN+"/*")
|
||||
default:
|
||||
out = append(out, bucketARN+"/"+strings.TrimPrefix(r, "/")+"*")
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// buildCondition renders the S3 condition block for a grant, or nil when there
|
||||
// is nothing to add.
|
||||
func buildCondition(c *GrantConditions) map[string]map[string][]string {
|
||||
if c == nil {
|
||||
return nil
|
||||
}
|
||||
cond := map[string]map[string][]string{}
|
||||
if len(c.SourceIPs) > 0 {
|
||||
cond["IpAddress"] = map[string][]string{"aws:SourceIp": c.SourceIPs}
|
||||
}
|
||||
if c.SecureTransportOnly {
|
||||
cond["Bool"] = map[string][]string{"aws:SecureTransport": {"true"}}
|
||||
}
|
||||
if len(cond) == 0 {
|
||||
return nil
|
||||
}
|
||||
return cond
|
||||
}
|
||||
|
||||
// sid builds a policy statement id that only contains characters S3 accepts.
|
||||
func sid(prefix, uid string) string {
|
||||
var b strings.Builder
|
||||
@@ -139,7 +270,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
|
||||
|
||||
@@ -88,3 +88,155 @@ func TestBuildBucketPolicyStructure(t *testing.T) {
|
||||
t.Fatal("reader principal ARN missing")
|
||||
}
|
||||
}
|
||||
|
||||
// parsedPolicy is a fuller parse of a rendered policy for the fine-grained tests.
|
||||
type parsedPolicy struct {
|
||||
Statement []struct {
|
||||
Sid string `json:"Sid"`
|
||||
Effect string `json:"Effect"`
|
||||
Principal map[string][]string `json:"Principal"`
|
||||
Action []string `json:"Action"`
|
||||
Resource []string `json:"Resource"`
|
||||
Condition map[string]map[string][]string `json:"Condition"`
|
||||
} `json:"Statement"`
|
||||
}
|
||||
|
||||
func parsePolicy(t *testing.T, raw string) parsedPolicy {
|
||||
t.Helper()
|
||||
var doc parsedPolicy
|
||||
if err := json.Unmarshal([]byte(raw), &doc); err != nil {
|
||||
t.Fatalf("policy is not valid JSON: %v\n%s", err, raw)
|
||||
}
|
||||
return doc
|
||||
}
|
||||
|
||||
func TestBuildBucketPolicyPaths(t *testing.T) {
|
||||
raw, err := BuildBucketPolicy("data", []Grant{
|
||||
{UID: "reader", Level: LevelReadOnly, Paths: []string{"team-a/", "/shared/inbox/"}},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
doc := parsePolicy(t, raw)
|
||||
|
||||
var objResources []string
|
||||
for _, s := range doc.Statement {
|
||||
for _, a := range s.Action {
|
||||
if a == "s3:GetObject" {
|
||||
objResources = s.Resource
|
||||
}
|
||||
}
|
||||
}
|
||||
want := map[string]bool{
|
||||
"arn:aws:s3:::data/shared/inbox/*": false,
|
||||
"arn:aws:s3:::data/team-a/*": false,
|
||||
}
|
||||
if len(objResources) != len(want) {
|
||||
t.Fatalf("expected %d object resources, got %v", len(want), objResources)
|
||||
}
|
||||
for _, r := range objResources {
|
||||
if _, ok := want[r]; !ok {
|
||||
t.Fatalf("unexpected object resource %q (leading slash not trimmed?)", r)
|
||||
}
|
||||
want[r] = true
|
||||
}
|
||||
for r, seen := range want {
|
||||
if !seen {
|
||||
t.Fatalf("missing object resource %q", r)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildBucketPolicyActionsOverride(t *testing.T) {
|
||||
raw, err := BuildBucketPolicy("data", []Grant{
|
||||
{UID: "svc", Level: LevelReadOnly, Actions: []string{"s3:GetObject", "s3:PutObject"}},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
doc := parsePolicy(t, raw)
|
||||
if len(doc.Statement) != 1 {
|
||||
t.Fatalf("expected a single custom-action statement, got %d", len(doc.Statement))
|
||||
}
|
||||
s := doc.Statement[0]
|
||||
if len(s.Action) != 2 || s.Action[0] != "s3:GetObject" || s.Action[1] != "s3:PutObject" {
|
||||
t.Fatalf("actions not taken verbatim: %v", s.Action)
|
||||
}
|
||||
// Custom-action statement lists both the bucket and object resources.
|
||||
if len(s.Resource) != 2 || s.Resource[0] != "arn:aws:s3:::data" || s.Resource[1] != "arn:aws:s3:::data/*" {
|
||||
t.Fatalf("unexpected resources: %v", s.Resource)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildBucketPolicyConditions(t *testing.T) {
|
||||
raw, err := BuildBucketPolicy("data", []Grant{
|
||||
{UID: "reader", Level: LevelReadOnly, Conditions: &GrantConditions{
|
||||
SourceIPs: []string{"10.0.0.0/8"},
|
||||
SecureTransportOnly: true,
|
||||
}},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
doc := parsePolicy(t, raw)
|
||||
for _, s := range doc.Statement {
|
||||
if s.Condition == nil {
|
||||
t.Fatalf("statement %q missing condition block", s.Sid)
|
||||
}
|
||||
if ip := s.Condition["IpAddress"]["aws:SourceIp"]; len(ip) != 1 || ip[0] != "10.0.0.0/8" {
|
||||
t.Fatalf("unexpected SourceIp condition: %v", s.Condition["IpAddress"])
|
||||
}
|
||||
if tls := s.Condition["Bool"]["aws:SecureTransport"]; len(tls) != 1 || tls[0] != "true" {
|
||||
t.Fatalf("unexpected SecureTransport condition: %v", s.Condition["Bool"])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildBucketPolicyRawStatements(t *testing.T) {
|
||||
raw, err := BuildBucketPolicy("data", []Grant{
|
||||
{UID: "svc", Level: LevelReadOnly, Raw: []RawStatement{
|
||||
{
|
||||
Effect: "Deny",
|
||||
Actions: []string{"s3:DeleteObject"},
|
||||
Resources: []string{"locked/", "arn:aws:s3:::other/*"},
|
||||
},
|
||||
}},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
doc := parsePolicy(t, raw)
|
||||
if len(doc.Statement) != 1 {
|
||||
t.Fatalf("expected 1 raw statement, got %d", len(doc.Statement))
|
||||
}
|
||||
s := doc.Statement[0]
|
||||
if s.Effect != "Deny" {
|
||||
t.Fatalf("raw effect not honoured: %q", s.Effect)
|
||||
}
|
||||
// The operator fills in the principal; bucket-relative prefixes are expanded
|
||||
// while explicit ARNs pass through.
|
||||
if len(s.Principal["AWS"]) != 1 || !strings.HasSuffix(s.Principal["AWS"][0], "user/svc") {
|
||||
t.Fatalf("raw statement principal not injected: %v", s.Principal)
|
||||
}
|
||||
if len(s.Resource) != 2 || s.Resource[0] != "arn:aws:s3:::data/locked/*" || s.Resource[1] != "arn:aws:s3:::other/*" {
|
||||
t.Fatalf("unexpected raw resources: %v", s.Resource)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildBucketPolicyFineGrainedDeterministic(t *testing.T) {
|
||||
grants := []Grant{
|
||||
{UID: "reader", Level: LevelReadOnly, Paths: []string{"a/", "b/"}},
|
||||
{UID: "svc", Level: LevelReadWrite, Conditions: &GrantConditions{SecureTransportOnly: true}},
|
||||
}
|
||||
a, err := BuildBucketPolicy("data", grants)
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
b, err := BuildBucketPolicy("data", []Grant{grants[1], grants[0]})
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
if a != b {
|
||||
t.Fatalf("fine-grained policy is order-dependent:\n a=%s\n b=%s", a, b)
|
||||
}
|
||||
}
|
||||
|
||||
+83
-70
@@ -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
|
||||
}
|
||||
|
||||
@@ -2,6 +2,7 @@ package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
|
||||
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
@@ -176,17 +177,52 @@ func (r *BucketReconciler) collectGrants(ctx context.Context, namespace, bucketR
|
||||
if ba.Status.UID == "" {
|
||||
continue
|
||||
}
|
||||
key := ba.Status.UID + "|" + string(ba.Spec.Level)
|
||||
g := grantFromAccess(ba.Status.UID, ba)
|
||||
key := grantKey(g)
|
||||
if _, dup := seen[key]; dup {
|
||||
continue
|
||||
}
|
||||
seen[key] = struct{}{}
|
||||
principals[ba.Status.UID] = struct{}{}
|
||||
grants = append(grants, ceph.Grant{UID: ba.Status.UID, Level: string(ba.Spec.Level)})
|
||||
grants = append(grants, g)
|
||||
}
|
||||
return grants, len(principals), nil
|
||||
}
|
||||
|
||||
// grantFromAccess translates a BucketAccess spec into the ceph grant model,
|
||||
// carrying the fine-grained scoping (paths, actions, conditions, raw statements).
|
||||
func grantFromAccess(uid string, ba *v1alpha1.BucketAccess) ceph.Grant {
|
||||
g := ceph.Grant{
|
||||
UID: uid,
|
||||
Level: string(ba.Spec.Level),
|
||||
Paths: ba.Spec.Paths,
|
||||
Actions: ba.Spec.Actions,
|
||||
}
|
||||
if c := ba.Spec.Conditions; c != nil {
|
||||
g.Conditions = &ceph.GrantConditions{
|
||||
SourceIPs: c.SourceIPs,
|
||||
SecureTransportOnly: c.SecureTransportOnly,
|
||||
}
|
||||
}
|
||||
for _, s := range ba.Spec.RawStatements {
|
||||
g.Raw = append(g.Raw, ceph.RawStatement{
|
||||
Sid: s.Sid,
|
||||
Effect: s.Effect,
|
||||
Actions: s.Actions,
|
||||
Resources: s.Resources,
|
||||
Condition: s.Conditions,
|
||||
})
|
||||
}
|
||||
return g
|
||||
}
|
||||
|
||||
// grantKey is a stable fingerprint of a grant used to collapse duplicate
|
||||
// BucketAccess objects that would render identical policy statements.
|
||||
func grantKey(g ceph.Grant) string {
|
||||
b, _ := json.Marshal(g)
|
||||
return string(b)
|
||||
}
|
||||
|
||||
func (r *BucketReconciler) pending(ctx context.Context, b *v1alpha1.Bucket, reason, msg string) (ctrl.Result, error) {
|
||||
b.Status.Phase = "Pending"
|
||||
b.Status.ObservedGeneration = b.Generation
|
||||
|
||||
Reference in New Issue
Block a user