mirror of
https://github.com/labring/sealos.git
synced 2026-08-28 17:22:42 +08:00
fix(objectstorage): tighten public read bucket policy (#7066)
* fix: tighten public read bucket policy * fix(objectstorage): remove list access from public readwrite * fix(objectstorage): preserve scaffold markers
This commit is contained in:
@@ -21,18 +21,15 @@ import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
stderrors "errors"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
"github.com/labring/sealos/controllers/pkg/utils/env"
|
||||
"github.com/minio/madmin-go/v3"
|
||||
"github.com/minio/minio-go/v7"
|
||||
|
||||
"github.com/labring/sealos/controllers/pkg/utils/env"
|
||||
|
||||
objectstoragev1 "github/labring/sealos/controllers/objectstorage/api/v1"
|
||||
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
"k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
@@ -72,12 +69,22 @@ const (
|
||||
//+kubebuilder:rbac:groups=objectstorage.sealos.io,resources=objectstoragebuckets/status,verbs=get;update;patch
|
||||
//+kubebuilder:rbac:groups=objectstorage.sealos.io,resources=objectstoragebuckets/finalizers,verbs=update
|
||||
|
||||
func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
|
||||
func (r *ObjectStorageBucketReconciler) Reconcile(
|
||||
ctx context.Context,
|
||||
req ctrl.Request,
|
||||
) (ctrl.Result, error) {
|
||||
// new OSClient
|
||||
if r.OSAdminClient == nil || r.OSClient == nil {
|
||||
secret := &corev1.Secret{}
|
||||
if err := r.Get(ctx, client.ObjectKey{Name: r.OSAdminSecret, Namespace: r.OSNamespace}, secret); err != nil {
|
||||
r.Logger.Error(err, "failed to get secret", "name", r.OSAdminSecret, "namespace", r.OSNamespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to get secret",
|
||||
"name",
|
||||
r.OSAdminSecret,
|
||||
"namespace",
|
||||
r.OSNamespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
@@ -109,7 +116,14 @@ func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.
|
||||
bucket := &objectstoragev1.ObjectStorageBucket{}
|
||||
if err := r.Get(ctx, client.ObjectKey{Name: req.Name, Namespace: namespace}, bucket); err != nil {
|
||||
if !errors.IsNotFound(err) {
|
||||
r.Logger.Error(err, "failed to get object storage bucket", "name", req.Name, "namespace", namespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to get object storage bucket",
|
||||
"name",
|
||||
req.Name,
|
||||
"namespace",
|
||||
namespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
@@ -123,7 +137,14 @@ func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.
|
||||
if serviceAccount.AccessKey == serviceAccountName {
|
||||
err := r.OSAdminClient.DeleteServiceAccount(ctx, serviceAccountName)
|
||||
if err != nil {
|
||||
r.Logger.Error(err, "failed to delete service account", "serviceAccountName", serviceAccountName, "bucketName", bucketName)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to delete service account",
|
||||
"serviceAccountName",
|
||||
serviceAccountName,
|
||||
"bucketName",
|
||||
bucketName,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
break
|
||||
@@ -146,7 +167,14 @@ func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.
|
||||
})
|
||||
for object := range objects {
|
||||
if err := r.OSClient.RemoveObject(ctx, bucketName, object.Key, minio.RemoveObjectOptions{}); err != nil {
|
||||
r.Logger.Error(err, "failed to remove object from bucket", "object", object.Key, "bucket", bucketName)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to remove object from bucket",
|
||||
"object",
|
||||
object.Key,
|
||||
"bucket",
|
||||
bucketName,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
}
|
||||
@@ -163,7 +191,14 @@ func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.
|
||||
user := &objectstoragev1.ObjectStorageUser{}
|
||||
if err := r.Get(ctx, client.ObjectKey{Name: username, Namespace: namespace}, user); err != nil {
|
||||
if !errors.IsNotFound(err) {
|
||||
r.Logger.Error(err, "failed to get object storage user", "name", username, "namespace", namespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to get object storage user",
|
||||
"name",
|
||||
username,
|
||||
"namespace",
|
||||
namespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
@@ -171,7 +206,14 @@ func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.
|
||||
user.Name = username
|
||||
user.Namespace = namespace
|
||||
if err := r.Create(ctx, user); err != nil {
|
||||
r.Logger.Error(err, "failed to create object storage user", "name", username, "namespace", namespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to create object storage user",
|
||||
"name",
|
||||
username,
|
||||
"namespace",
|
||||
namespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
}
|
||||
@@ -193,7 +235,14 @@ func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.
|
||||
|
||||
// set bucket policy
|
||||
if err := r.OSClient.SetBucketPolicy(ctx, bucketName, buildPolicy(bucket.Spec.Policy, bucketName)); err != nil {
|
||||
r.Logger.Error(err, "failed to set policy for bucket", "name", bucketName, "policy", bucket.Spec.Policy)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to set policy for bucket",
|
||||
"name",
|
||||
bucketName,
|
||||
"policy",
|
||||
bucket.Spec.Policy,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
@@ -206,7 +255,14 @@ func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.
|
||||
|
||||
if update {
|
||||
if err := r.Status().Update(ctx, bucket); err != nil {
|
||||
r.Logger.Error(err, "failed to update bucket status", "name", bucket.Name, "namespace", bucket.Namespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to update bucket status",
|
||||
"name",
|
||||
bucket.Name,
|
||||
"namespace",
|
||||
bucket.Namespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
}
|
||||
@@ -218,11 +274,19 @@ func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.
|
||||
userInfo, err := r.OSAdminClient.GetUserInfo(ctx, username)
|
||||
if err != nil {
|
||||
if err.Error() == userIsNotFound {
|
||||
r.Logger.V(1).Info("the minio user is being created", "user", username, "namespace", namespace)
|
||||
r.Logger.V(1).
|
||||
Info("the minio user is being created", "user", username, "namespace", namespace)
|
||||
return ctrl.Result{Requeue: true}, nil
|
||||
}
|
||||
|
||||
r.Logger.Error(err, "failed to get minio user info", "user", username, "namespace", namespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to get minio user info",
|
||||
"user",
|
||||
username,
|
||||
"namespace",
|
||||
namespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
@@ -253,7 +317,14 @@ func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.
|
||||
}
|
||||
sa, err = r.OSAdminClient.AddServiceAccount(ctx, saReq)
|
||||
if err != nil {
|
||||
r.Logger.Error(err, "failed to add service account", "serviceAccountName", serviceAccountName, "bucket", bucketName)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to add service account",
|
||||
"serviceAccountName",
|
||||
serviceAccountName,
|
||||
"bucket",
|
||||
bucketName,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
}
|
||||
@@ -267,12 +338,26 @@ func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.
|
||||
secret.Namespace = namespace
|
||||
if err := r.Get(ctx, client.ObjectKey{Name: secretName, Namespace: namespace}, secret); err != nil {
|
||||
if !errors.IsNotFound(err) {
|
||||
r.Logger.Error(err, "failed to get object storage key secret", "name", secretName, "namespace", namespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to get object storage key secret",
|
||||
"name",
|
||||
secretName,
|
||||
"namespace",
|
||||
namespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
if err := r.newObjectStorageKeySecret(ctx, secret, bucket, accessKey, secretKey); err != nil {
|
||||
r.Logger.Error(err, "failed to new object storage key secret", "name", secretName, "namespace", namespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to new object storage key secret",
|
||||
"name",
|
||||
secretName,
|
||||
"namespace",
|
||||
namespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
}
|
||||
@@ -281,11 +366,19 @@ func (r *ObjectStorageBucketReconciler) Reconcile(ctx context.Context, req ctrl.
|
||||
|
||||
if keySecretUpdated {
|
||||
if err := r.Update(ctx, secret); err != nil {
|
||||
r.Logger.Error(err, "failed to update object storage key secret", "name", secretName, "namespace", namespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to update object storage key secret",
|
||||
"name",
|
||||
secretName,
|
||||
"namespace",
|
||||
namespace,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
r.Logger.V(1).Info("[bucket] bucket info", "name", bucket.Status.Name, "size", bucket.Status.Size, "policy", bucket.Spec.Policy)
|
||||
r.Logger.V(1).
|
||||
Info("[bucket] bucket info", "name", bucket.Status.Name, "size", bucket.Status.Size, "policy", bucket.Spec.Policy)
|
||||
|
||||
return ctrl.Result{Requeue: true, RequeueAfter: r.OSBDetectionCycle}, nil
|
||||
}
|
||||
@@ -295,11 +388,9 @@ func buildPolicy(policy, bucketName string) string {
|
||||
case PrivateBucketPolicy:
|
||||
return `{"Version":"2012-10-17","Statement":[]}`
|
||||
case PublicReadBucketPolicy:
|
||||
return `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"AWS":["*"]},"Action":["s3:GetBucketLocation","s3:ListBucket"],"Resource":["arn:aws:s3:::` + bucketName + `"]},
|
||||
{"Effect":"Allow","Principal":{"AWS":["*"]},"Action":["s3:GetObject"],"Resource":["arn:aws:s3:::` + bucketName + `/*"]}]}`
|
||||
return `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"AWS":["*"]},"Action":["s3:GetObject"],"Resource":["arn:aws:s3:::` + bucketName + `/*"]}]}`
|
||||
case PublicReadwriteBucketPolicy:
|
||||
return `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"AWS":["*"]},"Action":["s3:ListBucketMultipartUploads","s3:GetBucketLocation","s3:ListBucket"],"Resource":["arn:aws:s3:::` + bucketName + `"]},
|
||||
{"Effect":"Allow","Principal":{"AWS":["*"]},"Action":["s3:PutObject","s3:AbortMultipartUpload","s3:DeleteObject","s3:GetObject","s3:ListMultipartUploadParts"],"Resource":["arn:aws:s3:::` + bucketName + `/*"]}]}`
|
||||
return `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"AWS":["*"]},"Action":["s3:GetObject","s3:PutObject","s3:DeleteObject","s3:AbortMultipartUpload","s3:ListMultipartUploadParts"],"Resource":["arn:aws:s3:::` + bucketName + `/*"]}]}`
|
||||
case BucketServiceAccountPolicy:
|
||||
return `{
|
||||
"Version": "2012-10-17",
|
||||
@@ -309,14 +400,22 @@ func buildPolicy(policy, bucketName string) string {
|
||||
"Action": [
|
||||
"s3:ListBucket",
|
||||
"s3:ListBucketMultipartUploads",
|
||||
"s3:ListMultipartUploadParts",
|
||||
"s3:GetBucketPolicy",
|
||||
"s3:GetBucketLocation",
|
||||
"s3:GetBucketTagging",
|
||||
"s3:PutBucketTagging",
|
||||
"s3:PutBucketTagging"
|
||||
],
|
||||
"Resource": [
|
||||
"arn:aws:s3:::` + bucketName + `"
|
||||
]
|
||||
},
|
||||
{
|
||||
"Effect": "Allow",
|
||||
"Action": [
|
||||
"s3:GetObject",
|
||||
"s3:PutObject",
|
||||
"s3:DeleteObject"
|
||||
"s3:DeleteObject",
|
||||
"s3:ListMultipartUploadParts"
|
||||
],
|
||||
"Resource": [
|
||||
"arn:aws:s3:::` + bucketName + `/*"
|
||||
@@ -333,7 +432,12 @@ func buildBucketName(name, namespace string) string {
|
||||
return strings.Split(namespace, "-")[1] + "-" + name
|
||||
}
|
||||
|
||||
func (r *ObjectStorageBucketReconciler) newObjectStorageKeySecret(ctx context.Context, secret *corev1.Secret, bucket *objectstoragev1.ObjectStorageBucket, accessKey, secretKey string) error {
|
||||
func (r *ObjectStorageBucketReconciler) newObjectStorageKeySecret(
|
||||
ctx context.Context,
|
||||
secret *corev1.Secret,
|
||||
bucket *objectstoragev1.ObjectStorageBucket,
|
||||
accessKey, secretKey string,
|
||||
) error {
|
||||
secret.Data = make(map[string][]byte)
|
||||
secret.Data[OSKeySecretAccessKey] = []byte(accessKey)
|
||||
secret.Data[OSKeySecretSecretKey] = []byte(secretKey)
|
||||
@@ -356,8 +460,11 @@ func (r *ObjectStorageBucketReconciler) newObjectStorageKeySecret(ctx context.Co
|
||||
return r.Create(ctx, secret)
|
||||
}
|
||||
|
||||
func (r *ObjectStorageBucketReconciler) initObjectStorageKeySecret(secret *corev1.Secret, accessKey, secretKey, bucketName string) bool {
|
||||
var updated = false
|
||||
func (r *ObjectStorageBucketReconciler) initObjectStorageKeySecret(
|
||||
secret *corev1.Secret,
|
||||
accessKey, secretKey, bucketName string,
|
||||
) bool {
|
||||
updated := false
|
||||
|
||||
if !bytes.Equal(secret.Data[OSKeySecretAccessKey], []byte(accessKey)) {
|
||||
secret.Data[OSKeySecretAccessKey] = []byte(accessKey)
|
||||
@@ -414,7 +521,9 @@ func (r *ObjectStorageBucketReconciler) SetupWithManager(mgr ctrl.Manager) error
|
||||
r.OSAdminSecret = oSAdminSecret
|
||||
|
||||
if internalEndpoint == "" || oSNamespace == "" || oSAdminSecret == "" {
|
||||
return fmt.Errorf("failed to get the endpoint or namespace or admin secret env of object storage")
|
||||
return stderrors.New(
|
||||
"failed to get the endpoint or namespace or admin secret env of object storage",
|
||||
)
|
||||
}
|
||||
|
||||
return ctrl.NewControllerManagedBy(mgr).
|
||||
|
||||
@@ -0,0 +1,134 @@
|
||||
/*
|
||||
Copyright 2023.
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package controllers
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"testing"
|
||||
)
|
||||
|
||||
type policyDocument struct {
|
||||
Statement []policyStatement `json:"Statement"`
|
||||
}
|
||||
|
||||
type policyStatement struct {
|
||||
Action []string `json:"Action"`
|
||||
Resource []string `json:"Resource"`
|
||||
}
|
||||
|
||||
func TestBuildPolicyPublicReadDoesNotAllowAnonymousList(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
policy := mustParsePolicy(t, buildPolicy(PublicReadBucketPolicy, "test-bucket"))
|
||||
actions := collectActions(policy)
|
||||
|
||||
assertHasAction(t, actions, "s3:GetObject")
|
||||
assertNoAction(t, actions, "s3:ListBucket")
|
||||
assertNoAction(t, actions, "s3:DeleteObject")
|
||||
}
|
||||
|
||||
func TestBuildPolicyPublicReadwriteDoesNotAllowAnonymousList(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
policy := mustParsePolicy(t, buildPolicy(PublicReadwriteBucketPolicy, "test-bucket"))
|
||||
actions := collectActions(policy)
|
||||
|
||||
assertHasAction(t, actions, "s3:GetObject")
|
||||
assertHasAction(t, actions, "s3:PutObject")
|
||||
assertHasAction(t, actions, "s3:DeleteObject")
|
||||
assertHasAction(t, actions, "s3:AbortMultipartUpload")
|
||||
assertHasAction(t, actions, "s3:ListMultipartUploadParts")
|
||||
assertNoAction(t, actions, "s3:GetBucketLocation")
|
||||
assertNoAction(t, actions, "s3:ListBucket")
|
||||
assertNoAction(t, actions, "s3:ListBucketMultipartUploads")
|
||||
}
|
||||
|
||||
func TestBuildPolicyBucketServiceAccountSeparatesBucketAndObjectResources(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const bucketName = "test-bucket"
|
||||
|
||||
policy := mustParsePolicy(t, buildPolicy(BucketServiceAccountPolicy, bucketName))
|
||||
if len(policy.Statement) != 2 {
|
||||
t.Fatalf("expected 2 statements, got %d", len(policy.Statement))
|
||||
}
|
||||
|
||||
bucketStatement := policy.Statement[0]
|
||||
assertResources(t, bucketStatement.Resource, "arn:aws:s3:::"+bucketName)
|
||||
assertHasAction(t, actionSet(bucketStatement.Action), "s3:ListBucket")
|
||||
assertHasAction(t, actionSet(bucketStatement.Action), "s3:GetBucketLocation")
|
||||
assertNoAction(t, actionSet(bucketStatement.Action), "s3:GetObject")
|
||||
|
||||
objectStatement := policy.Statement[1]
|
||||
assertResources(t, objectStatement.Resource, "arn:aws:s3:::"+bucketName+"/*")
|
||||
assertHasAction(t, actionSet(objectStatement.Action), "s3:GetObject")
|
||||
assertHasAction(t, actionSet(objectStatement.Action), "s3:PutObject")
|
||||
assertHasAction(t, actionSet(objectStatement.Action), "s3:DeleteObject")
|
||||
assertNoAction(t, actionSet(objectStatement.Action), "s3:ListBucket")
|
||||
}
|
||||
|
||||
func mustParsePolicy(t *testing.T, raw string) policyDocument {
|
||||
t.Helper()
|
||||
|
||||
var policy policyDocument
|
||||
if err := json.Unmarshal([]byte(raw), &policy); err != nil {
|
||||
t.Fatalf("policy is invalid JSON: %v", err)
|
||||
}
|
||||
return policy
|
||||
}
|
||||
|
||||
func collectActions(policy policyDocument) map[string]struct{} {
|
||||
actions := map[string]struct{}{}
|
||||
for _, statement := range policy.Statement {
|
||||
for _, action := range statement.Action {
|
||||
actions[action] = struct{}{}
|
||||
}
|
||||
}
|
||||
return actions
|
||||
}
|
||||
|
||||
func actionSet(actions []string) map[string]struct{} {
|
||||
actionMap := map[string]struct{}{}
|
||||
for _, action := range actions {
|
||||
actionMap[action] = struct{}{}
|
||||
}
|
||||
return actionMap
|
||||
}
|
||||
|
||||
func assertHasAction(t *testing.T, actions map[string]struct{}, action string) {
|
||||
t.Helper()
|
||||
|
||||
if _, ok := actions[action]; !ok {
|
||||
t.Fatalf("expected action %q", action)
|
||||
}
|
||||
}
|
||||
|
||||
func assertNoAction(t *testing.T, actions map[string]struct{}, action string) {
|
||||
t.Helper()
|
||||
|
||||
if _, ok := actions[action]; ok {
|
||||
t.Fatalf("did not expect action %q", action)
|
||||
}
|
||||
}
|
||||
|
||||
func assertResources(t *testing.T, resources []string, want string) {
|
||||
t.Helper()
|
||||
|
||||
if len(resources) != 1 || resources[0] != want {
|
||||
t.Fatalf("expected resources [%q], got %v", want, resources)
|
||||
}
|
||||
}
|
||||
@@ -19,19 +19,18 @@ package controllers
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
stderrors "errors"
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
"github.com/minio/madmin-go/v3"
|
||||
"github.com/minio/minio-go/v7"
|
||||
|
||||
myObjectStorage "github.com/labring/sealos/controllers/pkg/objectstorage"
|
||||
"github.com/labring/sealos/controllers/pkg/utils/env"
|
||||
|
||||
"github.com/minio/madmin-go/v3"
|
||||
"github.com/minio/minio-go/v7"
|
||||
objectstoragev1 "github/labring/sealos/controllers/objectstorage/api/v1"
|
||||
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
"k8s.io/apimachinery/pkg/api/errors"
|
||||
"k8s.io/apimachinery/pkg/api/resource"
|
||||
@@ -89,13 +88,17 @@ const (
|
||||
//+kubebuilder:rbac:groups=core,resources=resourcequotas,verbs=get;list;watch;create;update;patch;delete
|
||||
//+kubebuilder:rbac:groups=core,resources=resourcequotas/status,verbs=get;list;watch;create;update;patch;delete
|
||||
|
||||
func (r *ObjectStorageUserReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
|
||||
func (r *ObjectStorageUserReconciler) Reconcile(
|
||||
ctx context.Context,
|
||||
req ctrl.Request,
|
||||
) (ctrl.Result, error) {
|
||||
username := req.Name
|
||||
userNamespace := req.Namespace
|
||||
|
||||
// check object storage user name if correct or not
|
||||
if username != strings.Split(userNamespace, "-")[1] {
|
||||
r.Logger.V(1).Info("object storage user name is not correspond to the namespace", "name", username, "namespace", userNamespace)
|
||||
r.Logger.V(1).
|
||||
Info("object storage user name is not correspond to the namespace", "name", username, "namespace", userNamespace)
|
||||
return ctrl.Result{}, nil
|
||||
}
|
||||
|
||||
@@ -103,7 +106,14 @@ func (r *ObjectStorageUserReconciler) Reconcile(ctx context.Context, req ctrl.Re
|
||||
if r.OSAdminClient == nil || r.OSClient == nil {
|
||||
secret := &corev1.Secret{}
|
||||
if err := r.Get(ctx, client.ObjectKey{Name: r.OSAdminSecret, Namespace: r.OSNamespace}, secret); err != nil {
|
||||
r.Logger.Error(err, "failed to get secret", "name", r.OSAdminSecret, "namespace", r.OSNamespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to get secret",
|
||||
"name",
|
||||
r.OSAdminSecret,
|
||||
"namespace",
|
||||
r.OSNamespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
@@ -130,12 +140,26 @@ func (r *ObjectStorageUserReconciler) Reconcile(ctx context.Context, req ctrl.Re
|
||||
user := &objectstoragev1.ObjectStorageUser{}
|
||||
if err := r.Get(ctx, client.ObjectKey{Name: username, Namespace: userNamespace}, user); err != nil {
|
||||
if !errors.IsNotFound(err) {
|
||||
r.Logger.Error(err, "failed to get object storage user", "name", username, "namespace", userNamespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to get object storage user",
|
||||
"name",
|
||||
username,
|
||||
"namespace",
|
||||
userNamespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
if err := r.deleteObjectStorageUser(ctx, username, userNamespace); err != nil {
|
||||
r.Logger.Error(err, "failed to delete object storage user", "name", username, "namespace", userNamespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to delete object storage user",
|
||||
"name",
|
||||
username,
|
||||
"namespace",
|
||||
userNamespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
@@ -144,7 +168,14 @@ func (r *ObjectStorageUserReconciler) Reconcile(ctx context.Context, req ctrl.Re
|
||||
|
||||
resourceQuota := &corev1.ResourceQuota{}
|
||||
if err := r.Get(ctx, client.ObjectKey{Name: ResourceQuotaPrefix + userNamespace, Namespace: userNamespace}, resourceQuota); err != nil {
|
||||
r.Logger.Error(err, "failed to get resource quota", "name", ResourceQuotaPrefix+userNamespace, "namespace", userNamespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to get resource quota",
|
||||
"name",
|
||||
ResourceQuotaPrefix+userNamespace,
|
||||
"namespace",
|
||||
userNamespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
@@ -183,18 +214,33 @@ func (r *ObjectStorageUserReconciler) Reconcile(ctx context.Context, req ctrl.Re
|
||||
if err := r.OSAdminClient.SetUser(ctx, accessKey, secretKey, madmin.AccountEnabled); err != nil {
|
||||
r.Logger.Error(err, "failed to set user secret key", "name", accessKey)
|
||||
}
|
||||
r.Logger.V(1).Info("[user] password change info", "name", user.Name, "spec secret key version", user.Spec.SecretKeyVersion)
|
||||
r.Logger.V(1).
|
||||
Info("[user] password change info", "name", user.Name, "spec secret key version", user.Spec.SecretKeyVersion)
|
||||
}
|
||||
|
||||
secret := &corev1.Secret{}
|
||||
if err := r.Get(ctx, client.ObjectKey{Name: OSKeySecret, Namespace: userNamespace}, secret); err != nil {
|
||||
if !errors.IsNotFound(err) {
|
||||
r.Logger.Error(err, "failed to get object storage key secret", "name", OSKeySecret, "namespace", userNamespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to get object storage key secret",
|
||||
"name",
|
||||
OSKeySecret,
|
||||
"namespace",
|
||||
userNamespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
if err := r.newObjectStorageKeySecret(ctx, secret, user, accessKey, secretKey); err != nil {
|
||||
r.Logger.Error(err, "failed to new object storage key secret", "name", OSKeySecret, "namespace", userNamespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to new object storage key secret",
|
||||
"name",
|
||||
OSKeySecret,
|
||||
"namespace",
|
||||
userNamespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
}
|
||||
@@ -203,14 +249,28 @@ func (r *ObjectStorageUserReconciler) Reconcile(ctx context.Context, req ctrl.Re
|
||||
|
||||
if keySecretUpdated {
|
||||
if err := r.Update(ctx, secret); err != nil {
|
||||
r.Logger.Error(err, "failed to update object storage key secret", "name", OSKeySecret, "namespace", userNamespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to update object storage key secret",
|
||||
"name",
|
||||
OSKeySecret,
|
||||
"namespace",
|
||||
userNamespace,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
// check whether the space used exceeds the quota
|
||||
size, objectsCount, err := myObjectStorage.GetUserObjectStorageSize(r.OSClient, user.Name)
|
||||
if err != nil {
|
||||
r.Logger.Error(err, "failed to get user space used", "name", username, "namespace", userNamespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to get user space used",
|
||||
"name",
|
||||
username,
|
||||
"namespace",
|
||||
userNamespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
|
||||
@@ -224,7 +284,14 @@ func (r *ObjectStorageUserReconciler) Reconcile(ctx context.Context, req ctrl.Re
|
||||
if used.String() != stringSize {
|
||||
resourceQuota.Status.Used[ResourceObjectStorageSize] = resource.MustParse(stringSize)
|
||||
if err := r.Status().Update(ctx, resourceQuota); err != nil {
|
||||
r.Logger.Error(err, "failed to update status", "name", resourceQuota.Name, "namespace", userNamespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to update status",
|
||||
"name",
|
||||
resourceQuota.Name,
|
||||
"namespace",
|
||||
userNamespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
}
|
||||
@@ -236,12 +303,20 @@ func (r *ObjectStorageUserReconciler) Reconcile(ctx context.Context, req ctrl.Re
|
||||
|
||||
if updated {
|
||||
if err := r.Status().Update(ctx, user); err != nil {
|
||||
r.Logger.Error(err, "failed to update status", "name", username, "namespace", userNamespace)
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to update status",
|
||||
"name",
|
||||
username,
|
||||
"namespace",
|
||||
userNamespace,
|
||||
)
|
||||
return ctrl.Result{}, err
|
||||
}
|
||||
}
|
||||
|
||||
r.Logger.V(1).Info("[user] user info", "name", user.Name, "quota", user.Status.Quota, "size", size, "objectsCount", user.Status.ObjectsCount)
|
||||
r.Logger.V(1).
|
||||
Info("[user] user info", "name", user.Name, "quota", user.Status.Quota, "size", size, "objectsCount", user.Status.ObjectsCount)
|
||||
|
||||
if r.QuotaEnabled && size > user.Status.Quota {
|
||||
if err := r.addUserToGroup(ctx, accessKey, UserDenyWriteGroup); err != nil {
|
||||
@@ -258,7 +333,10 @@ func (r *ObjectStorageUserReconciler) Reconcile(ctx context.Context, req ctrl.Re
|
||||
return ctrl.Result{Requeue: true, RequeueAfter: r.OSUDetectionCycle}, nil
|
||||
}
|
||||
|
||||
func (r *ObjectStorageUserReconciler) NewObjectStorageUser(ctx context.Context, accessKey, secretKey string) error {
|
||||
func (r *ObjectStorageUserReconciler) NewObjectStorageUser(
|
||||
ctx context.Context,
|
||||
accessKey, secretKey string,
|
||||
) error {
|
||||
if err := r.OSAdminClient.AddUser(ctx, accessKey, secretKey); err != nil {
|
||||
r.Logger.Error(err, "failed to create object storage user")
|
||||
return err
|
||||
@@ -272,7 +350,12 @@ func (r *ObjectStorageUserReconciler) NewObjectStorageUser(ctx context.Context,
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *ObjectStorageUserReconciler) newObjectStorageKeySecret(ctx context.Context, secret *corev1.Secret, user *objectstoragev1.ObjectStorageUser, accessKey, secretKey string) error {
|
||||
func (r *ObjectStorageUserReconciler) newObjectStorageKeySecret(
|
||||
ctx context.Context,
|
||||
secret *corev1.Secret,
|
||||
user *objectstoragev1.ObjectStorageUser,
|
||||
accessKey, secretKey string,
|
||||
) error {
|
||||
secret.SetName(OSKeySecret)
|
||||
secret.SetNamespace(user.Namespace)
|
||||
|
||||
@@ -297,7 +380,10 @@ func (r *ObjectStorageUserReconciler) newObjectStorageKeySecret(ctx context.Cont
|
||||
return r.Create(ctx, secret)
|
||||
}
|
||||
|
||||
func (r *ObjectStorageUserReconciler) addUserToGroup(ctx context.Context, user string, group string) error {
|
||||
func (r *ObjectStorageUserReconciler) addUserToGroup(
|
||||
ctx context.Context,
|
||||
user, group string,
|
||||
) error {
|
||||
newGroupDesc := madmin.GroupAddRemove{}
|
||||
|
||||
groupDesc, err := r.OSAdminClient.GetGroupDescription(ctx, group)
|
||||
@@ -310,7 +396,7 @@ func (r *ObjectStorageUserReconciler) addUserToGroup(ctx context.Context, user s
|
||||
newGroupDesc.Members = member
|
||||
newGroupDesc.Status = madmin.GroupStatus(groupDesc.Status)
|
||||
|
||||
var isExist = false
|
||||
isExist := false
|
||||
for _, member := range groupDesc.Members {
|
||||
if user == member {
|
||||
isExist = true
|
||||
@@ -328,7 +414,10 @@ func (r *ObjectStorageUserReconciler) addUserToGroup(ctx context.Context, user s
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *ObjectStorageUserReconciler) removeUserFromGroup(ctx context.Context, user string, group string) error {
|
||||
func (r *ObjectStorageUserReconciler) removeUserFromGroup(
|
||||
ctx context.Context,
|
||||
user, group string,
|
||||
) error {
|
||||
newGroupDesc := madmin.GroupAddRemove{}
|
||||
|
||||
groupDesc, err := r.OSAdminClient.GetGroupDescription(ctx, group)
|
||||
@@ -351,10 +440,22 @@ func (r *ObjectStorageUserReconciler) removeUserFromGroup(ctx context.Context, u
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *ObjectStorageUserReconciler) deleteObjectStorageUser(ctx context.Context, username, userNamespace string) error {
|
||||
func (r *ObjectStorageUserReconciler) deleteObjectStorageUser(
|
||||
ctx context.Context,
|
||||
username, userNamespace string,
|
||||
) error {
|
||||
// delete all bucket cr of user
|
||||
if err := r.Client.DeleteAllOf(ctx, &objectstoragev1.ObjectStorageBucket{}, client.InNamespace(userNamespace)); client.IgnoreNotFound(err) != nil {
|
||||
r.Logger.Error(err, "failed to delete all bucket of user", "name", username, "namespace", userNamespace)
|
||||
if err := r.DeleteAllOf(ctx, &objectstoragev1.ObjectStorageBucket{}, client.InNamespace(userNamespace)); client.IgnoreNotFound(
|
||||
err,
|
||||
) != nil {
|
||||
r.Logger.Error(
|
||||
err,
|
||||
"failed to delete all bucket of user",
|
||||
"name",
|
||||
username,
|
||||
"namespace",
|
||||
userNamespace,
|
||||
)
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -378,8 +479,12 @@ func (r *ObjectStorageUserReconciler) deleteObjectStorageUser(ctx context.Contex
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *ObjectStorageUserReconciler) initObjectStorageUser(user *objectstoragev1.ObjectStorageUser, username string, quota int64) bool {
|
||||
var updated = false
|
||||
func (r *ObjectStorageUserReconciler) initObjectStorageUser(
|
||||
user *objectstoragev1.ObjectStorageUser,
|
||||
username string,
|
||||
quota int64,
|
||||
) bool {
|
||||
updated := false
|
||||
|
||||
if user.Status.Quota != quota {
|
||||
user.Status.Quota = quota
|
||||
@@ -409,8 +514,11 @@ func (r *ObjectStorageUserReconciler) initObjectStorageUser(user *objectstoragev
|
||||
return updated
|
||||
}
|
||||
|
||||
func (r *ObjectStorageUserReconciler) initObjectStorageKeySecret(secret *corev1.Secret, accessKey, secretKey string) bool {
|
||||
var updated = false
|
||||
func (r *ObjectStorageUserReconciler) initObjectStorageKeySecret(
|
||||
secret *corev1.Secret,
|
||||
accessKey, secretKey string,
|
||||
) bool {
|
||||
updated := false
|
||||
|
||||
if !bytes.Equal(secret.Data[OSKeySecretAccessKey], []byte(accessKey)) {
|
||||
secret.Data[OSKeySecretAccessKey] = []byte(accessKey)
|
||||
@@ -453,7 +561,7 @@ func ConvertBytesToString(bytes int64) string {
|
||||
value = float64(bytes) / (1 << 10)
|
||||
unit = "Ki"
|
||||
default:
|
||||
return fmt.Sprintf("%d", bytes)
|
||||
return strconv.FormatInt(bytes, 10)
|
||||
}
|
||||
|
||||
return fmt.Sprintf("%.0f%s", value, unit)
|
||||
@@ -479,8 +587,11 @@ func (r *ObjectStorageUserReconciler) SetupWithManager(mgr ctrl.Manager) error {
|
||||
oSAdminSecret := env.GetEnvWithDefault(OSAdminSecret, "")
|
||||
r.OSAdminSecret = oSAdminSecret
|
||||
|
||||
if internalEndpoint == "" || externalEndpoint == "" || oSNamespace == "" || oSAdminSecret == "" {
|
||||
return fmt.Errorf("failed to get the endpoint or namespace or admin secret env of object storage")
|
||||
if internalEndpoint == "" || externalEndpoint == "" || oSNamespace == "" ||
|
||||
oSAdminSecret == "" {
|
||||
return stderrors.New(
|
||||
"failed to get the endpoint or namespace or admin secret env of object storage",
|
||||
)
|
||||
}
|
||||
|
||||
quotaEnabled := env.GetBoolWithDefault(QuotaEnabled, true)
|
||||
|
||||
@@ -22,24 +22,24 @@ import (
|
||||
|
||||
. "github.com/onsi/ginkgo/v2"
|
||||
. "github.com/onsi/gomega"
|
||||
|
||||
objectstoragev1 "github/labring/sealos/controllers/objectstorage/api/v1"
|
||||
"k8s.io/client-go/kubernetes/scheme"
|
||||
"k8s.io/client-go/rest"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
"sigs.k8s.io/controller-runtime/pkg/envtest"
|
||||
logf "sigs.k8s.io/controller-runtime/pkg/log"
|
||||
"sigs.k8s.io/controller-runtime/pkg/log/zap"
|
||||
|
||||
objectstoragev1 "github/labring/sealos/controllers/objectstorage/api/v1"
|
||||
//+kubebuilder:scaffold:imports
|
||||
"sigs.k8s.io/controller-runtime/pkg/log/zap"
|
||||
)
|
||||
|
||||
// These tests use Ginkgo (BDD-style Go testing framework). Refer to
|
||||
// http://onsi.github.io/ginkgo/ to learn more about Ginkgo.
|
||||
|
||||
var cfg *rest.Config
|
||||
var k8sClient client.Client
|
||||
var testEnv *envtest.Environment
|
||||
var (
|
||||
cfg *rest.Config
|
||||
k8sClient client.Client
|
||||
testEnv *envtest.Environment
|
||||
)
|
||||
|
||||
func TestAPIs(t *testing.T) {
|
||||
RegisterFailHandler(Fail)
|
||||
@@ -70,7 +70,6 @@ var _ = BeforeSuite(func() {
|
||||
k8sClient, err = client.New(cfg, client.Options{Scheme: scheme.Scheme})
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
Expect(k8sClient).NotTo(BeNil())
|
||||
|
||||
})
|
||||
|
||||
var _ = AfterSuite(func() {
|
||||
|
||||
@@ -20,19 +20,17 @@ import (
|
||||
"flag"
|
||||
"os"
|
||||
|
||||
_ "k8s.io/client-go/plugin/pkg/client/auth"
|
||||
|
||||
objectstoragev1 "github/labring/sealos/controllers/objectstorage/api/v1"
|
||||
objectstoragecontrollers "github/labring/sealos/controllers/objectstorage/controllers"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
|
||||
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
|
||||
_ "k8s.io/client-go/plugin/pkg/client/auth"
|
||||
ctrl "sigs.k8s.io/controller-runtime"
|
||||
"sigs.k8s.io/controller-runtime/pkg/healthz"
|
||||
"sigs.k8s.io/controller-runtime/pkg/log/zap"
|
||||
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
|
||||
|
||||
objectstoragev1 "github/labring/sealos/controllers/objectstorage/api/v1"
|
||||
"github/labring/sealos/controllers/objectstorage/controllers"
|
||||
//+kubebuilder:scaffold:imports
|
||||
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
|
||||
)
|
||||
|
||||
var (
|
||||
@@ -51,8 +49,20 @@ func main() {
|
||||
var metricsAddr string
|
||||
var enableLeaderElection bool
|
||||
var probeAddr string
|
||||
flag.StringVar(&metricsAddr, "metrics-bind-address", ":8080", "The address the metric endpoint binds to.")
|
||||
flag.StringVar(&probeAddr, "health-probe-bind-address", ":8081", "The address the probe endpoint binds to.")
|
||||
flag.StringVar(
|
||||
&metricsAddr,
|
||||
"metrics-bind-address",
|
||||
":8080",
|
||||
"The address the metric endpoint binds to.",
|
||||
)
|
||||
// Accept the legacy monitor flag passed by existing deployments.
|
||||
flag.String("monitor-bind-address", ":9090", "The address the monitor endpoint binds to.")
|
||||
flag.StringVar(
|
||||
&probeAddr,
|
||||
"health-probe-bind-address",
|
||||
":8081",
|
||||
"The address the probe endpoint binds to.",
|
||||
)
|
||||
flag.BoolVar(&enableLeaderElection, "leader-elect", false,
|
||||
"Enable leader election for controller manager. "+
|
||||
"Enabling this will ensure there is only one active controller manager.")
|
||||
@@ -89,14 +99,14 @@ func main() {
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
if err = (&controllers.ObjectStorageUserReconciler{
|
||||
if err = (&objectstoragecontrollers.ObjectStorageUserReconciler{
|
||||
Client: mgr.GetClient(),
|
||||
Scheme: mgr.GetScheme(),
|
||||
}).SetupWithManager(mgr); err != nil {
|
||||
setupLog.Error(err, "unable to create controller", "controller", "ObjectStorageUser")
|
||||
os.Exit(1)
|
||||
}
|
||||
if err = (&controllers.ObjectStorageBucketReconciler{
|
||||
if err = (&objectstoragecontrollers.ObjectStorageBucketReconciler{
|
||||
Client: mgr.GetClient(),
|
||||
Scheme: mgr.GetScheme(),
|
||||
}).SetupWithManager(mgr); err != nil {
|
||||
|
||||
Reference in New Issue
Block a user