Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
100 changes: 42 additions & 58 deletions internal/controller/bucket/acl.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,66 +5,54 @@ import (

"github.com/aws/aws-sdk-go-v2/aws"
s3types "github.com/aws/aws-sdk-go-v2/service/s3/types"
"github.com/crossplane/crossplane-runtime/v2/pkg/errors"
"github.com/go-logr/logr"

"github.com/linode/provider-ceph/apis/provider-ceph/v1alpha1"
apisv1alpha1 "github.com/linode/provider-ceph/apis/v1alpha1"
"github.com/linode/provider-ceph/internal/backendstore"
"github.com/linode/provider-ceph/internal/consts"
"github.com/linode/provider-ceph/internal/controller/s3clienthandler"
"github.com/linode/provider-ceph/internal/otel/traces"
"github.com/linode/provider-ceph/internal/rgw"

"go.opentelemetry.io/otel"
)

// ACLClient is the client for API methods and reconciling the ACL
type ACLClient struct {
backendStore *backendstore.BackendStore
s3ClientHandler *s3clienthandler.Handler
log logr.Logger
BaseSubresourceClient
}

// NewACLClient creates the client for ACL
func NewACLClient(b *backendstore.BackendStore, h *s3clienthandler.Handler, l logr.Logger) *ACLClient {
return &ACLClient{backendStore: b, s3ClientHandler: h, log: l}
return &ACLClient{BaseSubresourceClient: NewBaseSubresourceClient(b, h, l)}
}

func (l *ACLClient) Observe(ctx context.Context, bucket *v1alpha1.Bucket, backendNames []string) (ResourceStatus, error) {
_, span := otel.Tracer("").Start(ctx, "bucket.ACLClient.Observe")
defer span.End()

observationChan := make(chan ResourceStatus)

for _, backendName := range backendNames {
beName := backendName
go func() {
if l.backendStore.GetBackendHealthStatus(backendName) == apisv1alpha1.HealthStatusUnhealthy {
// If a backend is marked as unhealthy, we can ignore it for now by returning NoAction.
// The backend may be down for some time and we do not want to block Create/Update/Delete
// calls on other backends. By returning NoAction here, we would never pass the Observe
// phase until the backend becomes Healthy or Disabled.
observationChan <- NoAction

return
}
observationChan <- l.observeBackend(ctx, bucket, beName)
}()
}
func (a *ACLClient) Observe(ctx context.Context, bucket *v1alpha1.Bucket, backendNames []string) (ResourceStatus, error) {
return a.BaseSubresourceClient.Observe(ctx, bucket, backendNames, a)
}

for i := 0; i < len(backendNames); i++ {
observation := <-observationChan
if observation == NeedsUpdate || observation == NeedsDeletion {
return observation, nil
}
}
func (a *ACLClient) Handle(ctx context.Context, b *v1alpha1.Bucket, backendName string, bb *bucketBackends) error {
return a.BaseSubresourceClient.Handle(ctx, b, backendName, bb, a)
}

// Implement Subresource interface

func (a *ACLClient) GetLogger() logr.Logger {
return a.log
}

return Updated, nil
func (a *ACLClient) GetBackendStore() *backendstore.BackendStore {
return a.backendStore
}

func (l *ACLClient) observeBackend(ctx context.Context, bucket *v1alpha1.Bucket, backendName string) ResourceStatus {
_, log := traces.InjectTraceAndLogger(ctx, l.log)
func (a *ACLClient) GetS3ClientHandler() *s3clienthandler.Handler {
return a.s3ClientHandler
}

func (a *ACLClient) GetObserveErrorMsg() string {
return errObserveAcl
}

func (a *ACLClient) ObserveBackend(ctx context.Context, bucket *v1alpha1.Bucket, backendName string) (ResourceStatus, error) {
_, log := traces.InjectTraceAndLogger(ctx, a.log)

log.V(1).Info("Observing subresource acl on backend", consts.KeyBucketName, bucket.Name, consts.KeyBackendName, backendName)

Expand All @@ -73,7 +61,7 @@ func (l *ACLClient) observeBackend(ctx context.Context, bucket *v1alpha1.Bucket,
if s3types.ObjectOwnership(aws.ToString(bucket.Spec.ForProvider.ObjectOwnership)) == s3types.ObjectOwnershipBucketOwnerEnforced {
log.V(1).Info("Access control limits are disabled - no action required", consts.KeyBucketName, bucket.Name, consts.KeyBackendName, backendName)

return Updated
return Updated, nil
}

if bucket.Spec.ForProvider.ACL == nil &&
Expand All @@ -85,42 +73,38 @@ func (l *ACLClient) observeBackend(ctx context.Context, bucket *v1alpha1.Bucket,
bucket.Spec.ForProvider.GrantReadACP == nil {
log.V(1).Info("No acl or access control policy or grants requested - no action required", consts.KeyBucketName, bucket.Name, consts.KeyBackendName, backendName)

return Updated
return Updated, nil
}

return NeedsUpdate
return NeedsUpdate, nil
}

func (l *ACLClient) Handle(ctx context.Context, b *v1alpha1.Bucket, backendName string, bb *bucketBackends) error {
ctx, span := otel.Tracer("").Start(ctx, "bucket.ACLClient.Handle")
defer span.End()
// Implement Subresource interface

if l.backendStore.GetBackendHealthStatus(backendName) == apisv1alpha1.HealthStatusUnhealthy {
traces.SetAndRecordError(span, errUnhealthyBackend)
func (a *ACLClient) GetHandleErrorMsg() string {
return errHandleAcl
}

return errUnhealthyBackend
}
func (a *ACLClient) GetSubresourceName() string {
return "ACLClient"
}

switch l.observeBackend(ctx, b, backendName) {
func (a *ACLClient) HandleObservation(ctx context.Context, observation ResourceStatus, bucket *v1alpha1.Bucket, backendName string, bb *bucketBackends) error {
switch observation {
case NoAction, Updated:
return nil
case NeedsUpdate, NeedsDeletion:
Comment on lines +92 to 96
if err := l.createOrUpdate(ctx, b, backendName); err != nil {
err = errors.Wrap(err, errHandleAcl)
traces.SetAndRecordError(span, err)

return err
}
return a.createOrUpdate(ctx, bucket, backendName)
}

return nil
}

func (l *ACLClient) createOrUpdate(ctx context.Context, b *v1alpha1.Bucket, backendName string) error {
ctx, log := traces.InjectTraceAndLogger(ctx, l.log)
func (a *ACLClient) createOrUpdate(ctx context.Context, b *v1alpha1.Bucket, backendName string) error {
ctx, log := traces.InjectTraceAndLogger(ctx, a.log)

log.Info("Updating acl", consts.KeyBucketName, b.Name, consts.KeyBackendName, backendName)
s3Client, err := l.s3ClientHandler.GetS3Client(ctx, b, backendName)
s3Client, err := a.s3ClientHandler.GetS3Client(ctx, b, backendName)
if err != nil {
return err
}
Expand Down
4 changes: 3 additions & 1 deletion internal/controller/bucket/acl_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ import (
s3types "github.com/aws/aws-sdk-go-v2/service/s3/types"
"github.com/go-logr/logr"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"

"github.com/linode/provider-ceph/apis/provider-ceph/v1alpha1"
apisv1alpha1 "github.com/linode/provider-ceph/apis/v1alpha1"
Expand Down Expand Up @@ -240,7 +241,8 @@ func TestACLObserveBackend(t *testing.T) {
s3clienthandler.WithBackendStore(tc.fields.backendStore)),
logr.Discard())

got := c.observeBackend(context.Background(), tc.args.bucket, tc.args.backendName)
got, err := c.ObserveBackend(context.Background(), tc.args.bucket, tc.args.backendName)
require.NoError(t, err, "unexpected error")
assert.Equal(t, tc.want.status, got, "unexpected status")
})
}
Expand Down
167 changes: 167 additions & 0 deletions internal/controller/bucket/base_subresource.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,167 @@
package bucket

import (
"context"

"github.com/crossplane/crossplane-runtime/v2/pkg/errors"
"github.com/go-logr/logr"

"github.com/linode/provider-ceph/apis/provider-ceph/v1alpha1"
apisv1alpha1 "github.com/linode/provider-ceph/apis/v1alpha1"
"github.com/linode/provider-ceph/internal/backendstore"
"github.com/linode/provider-ceph/internal/consts"
"github.com/linode/provider-ceph/internal/controller/s3clienthandler"
"github.com/linode/provider-ceph/internal/otel/traces"

"go.opentelemetry.io/otel"
)

// Subresource provides the plugin interface for observing and handling backend state.
type Subresource interface {
// ObserveBackend observes the resource on a specific backend.
// Should return NoAction/Updated/NeedsUpdate/NeedsDeletion and any error.
ObserveBackend(ctx context.Context, bucket *v1alpha1.Bucket, backendName string) (ResourceStatus, error)

// GetObserveErrorMsg returns the error message for observe failures.
GetObserveErrorMsg() string

// SkipObservation returns true if observation should be skipped for this subresource.
// By default (false), all subresources are observed. Override to implement conditional observation.
SkipObservation(bucket *v1alpha1.Bucket) bool

// HandleObservation handles the resource on a specific backend based on the observation result.
HandleObservation(ctx context.Context, observation ResourceStatus, bucket *v1alpha1.Bucket, backendName string, bb *bucketBackends) error

// GetHandleErrorMsg returns the error message for handle failures.
GetHandleErrorMsg() string

// GetLogger returns the logger for this subresource.
GetLogger() logr.Logger

// GetBackendStore returns the backend store.
GetBackendStore() *backendstore.BackendStore

// GetS3ClientHandler returns the s3 client handler.
GetS3ClientHandler() *s3clienthandler.Handler

// GetSubresourceName returns the name of the subresource for tracing.
GetSubresourceName() string
}

// BaseSubresourceClient provides common logic for all subresource clients.
type BaseSubresourceClient struct {
backendStore *backendstore.BackendStore
s3ClientHandler *s3clienthandler.Handler
log logr.Logger
}

// NewBaseSubresourceClient creates a new base subresource client.
func NewBaseSubresourceClient(b *backendstore.BackendStore, h *s3clienthandler.Handler, l logr.Logger) BaseSubresourceClient {
return BaseSubresourceClient{
backendStore: b,
s3ClientHandler: h,
log: l,
}
}

// SkipObservation provides a default implementation that does not skip observations.
// Subresources can override this to implement conditional observation logic.
func (b *BaseSubresourceClient) SkipObservation(bucket *v1alpha1.Bucket) bool {
return false // By default, observe all subresources
}

// Observe implements the common observe pattern for all subresources.
// It handles concurrent observation of all backends and consolidates results.
func (b *BaseSubresourceClient) Observe(ctx context.Context, bucket *v1alpha1.Bucket, backendNames []string, subresource Subresource) (ResourceStatus, error) {
ctx, span := otel.Tracer("").Start(ctx, "bucket."+subresource.GetSubresourceName()+".Observe")
defer span.End()
ctx, log := traces.InjectTraceAndLogger(ctx, subresource.GetLogger())

// Check if this subresource should be skipped
if subresource.SkipObservation(bucket) {
log.V(1).Info(subresource.GetSubresourceName() + " observation skipped")

return Updated, nil
}

observationChan := make(chan ResourceStatus, len(backendNames))
errChan := make(chan error, len(backendNames))

for _, backendName := range backendNames {
beName := backendName
go func() {
if subresource.GetBackendStore().GetBackendHealthStatus(beName) == apisv1alpha1.HealthStatusUnhealthy {
// If a backend is marked as unhealthy, we can ignore it for now by returning NoAction.
Comment on lines +90 to +94
// The backend may be down for some time and we do not want to block Create/Update/Delete
// calls on other backends. By returning NoAction here, we would never pass the Observe
// phase until the backend becomes Healthy or Disabled.
observationChan <- NoAction

return
}

observation, err := subresource.ObserveBackend(ctx, bucket, beName)
if err != nil {
errChan <- err

return
}
observationChan <- observation
}()
}

for i := 0; i < len(backendNames); i++ {
select {
case <-ctx.Done():
log.Info("Context timeout during bucket "+subresource.GetSubresourceName()+" observation", consts.KeyBucketName, bucket.Name)
err := errors.Wrap(ctx.Err(), subresource.GetObserveErrorMsg())
traces.SetAndRecordError(span, err)

return NeedsUpdate, err
case observation := <-observationChan:
if observation == NeedsUpdate || observation == NeedsDeletion {
return observation, nil
}
case err := <-errChan:
err = errors.Wrap(err, subresource.GetObserveErrorMsg())
traces.SetAndRecordError(span, err)

return NeedsUpdate, err
}
}

return Updated, nil
}

// Handle implements the common handle pattern for all subresources.
// It performs health checks and delegates to the handler for observation-specific logic.
func (b *BaseSubresourceClient) Handle(ctx context.Context, bucket *v1alpha1.Bucket, backendName string, bb *bucketBackends, subresource Subresource) error {
ctx, span := otel.Tracer("").Start(ctx, "bucket."+subresource.GetSubresourceName()+".Handle")
defer span.End()

// Check if this subresource should be skipped
if subresource.SkipObservation(bucket) {
return nil
}

if subresource.GetBackendStore().GetBackendHealthStatus(backendName) == apisv1alpha1.HealthStatusUnhealthy {
traces.SetAndRecordError(span, errUnhealthyBackend)

return errUnhealthyBackend
}

observation, err := subresource.ObserveBackend(ctx, bucket, backendName)
if err != nil {
err = errors.Wrap(err, subresource.GetHandleErrorMsg())
traces.SetAndRecordError(span, err)

return err
}

err = subresource.HandleObservation(ctx, observation, bucket, backendName, bb)
if err != nil {
traces.SetAndRecordError(span, err)
}

return err
}
Loading
Loading