feat: enhance bucket credential retrieval to support read/write operations and improve caching logic (#46)

Signed-off-by: Noste <83548733+Noooste@users.noreply.github.com>
This commit is contained in:
Noste
2026-05-15 15:37:07 +02:00
committed by GitHub
parent c8337de3a8
commit c8cb3c4923
2 changed files with 258 additions and 53 deletions
+66 -45
View File
@@ -52,56 +52,77 @@ func NewS3Service(cfg *config.GarageConfig, adminService AdminService) *S3Servic
}
}
func (s *S3Service) getBucketCredentials(ctx context.Context, bucketName string) (*credentials.Credentials, error) {
cacheKey := fmt.Sprintf("key:%s", bucketName)
cacheData := utils.GlobalCache.Get(cacheKey)
// Operation is a bitmask of S3 permissions a call needs. Combine with bitwise
// OR (e.g. OpRead | OpWrite) when more than one is required.
type Operation byte
if cacheData != nil {
return cacheData.(*credentials.Credentials), nil
const (
OpRead Operation = 0x1
OpWrite Operation = 0x2
)
// satisfies reports whether perms grants every bit set in op.
func (op Operation) satisfies(perms models.BucketKeyPermission) bool {
if op&OpRead != 0 && !perms.Read {
return false
}
if op&OpWrite != 0 && !perms.Write {
return false
}
return true
}
func setKeyInCache(bucketName string, permissions models.BucketKeyPermission, creds *credentials.Credentials) {
canWrite := permissions.Write
canRead := permissions.Read
if canWrite {
key := fmt.Sprintf("key:%s:%d", bucketName, OpWrite)
utils.GlobalCache.Set(key, creds, time.Hour)
}
if canRead {
key := fmt.Sprintf("key:%s:%d", bucketName, OpRead)
utils.GlobalCache.Set(key, creds, time.Hour)
}
if canRead && canWrite {
key := fmt.Sprintf("key:%s:%d", bucketName, OpRead|OpWrite)
utils.GlobalCache.Set(key, creds, time.Hour)
}
}
func (s *S3Service) getBucketCredentials(ctx context.Context, bucketName string, op Operation) (*credentials.Credentials, error) {
cacheKey := fmt.Sprintf("key:%s:%d", bucketName, op)
if cached := utils.GlobalCache.Get(cacheKey); cached != nil {
return cached.(*credentials.Credentials), nil
}
// Get bucket info from Garage Admin API
bucketInfo, err := s.adminService.GetBucketInfoByAlias(ctx, bucketName)
if err != nil {
return nil, fmt.Errorf("failed to get bucket info: %w", err)
}
// Find a key with read and write permissions
var accessKeyID, secretAccessKey string
for _, keyInfo := range bucketInfo.Keys {
if !keyInfo.Permissions.Read || !keyInfo.Permissions.Write {
if !op.satisfies(keyInfo.Permissions) {
continue
}
// Get key details with secret
keyDetails, err := s.adminService.GetKeyInfo(ctx, keyInfo.AccessKeyID, true)
if err != nil {
return nil, fmt.Errorf("failed to get key info: %w", err)
}
if keyDetails.SecretAccessKey != nil {
accessKeyID = keyDetails.AccessKeyID
secretAccessKey = *keyDetails.SecretAccessKey
break
if err != nil || keyDetails.SecretAccessKey == nil {
continue
}
creds := credentials.NewStaticV4(keyDetails.AccessKeyID, *keyDetails.SecretAccessKey, "")
setKeyInCache(bucketName, keyInfo.Permissions, creds)
return creds, nil
}
if accessKeyID == "" || secretAccessKey == "" {
return nil, fmt.Errorf("no valid credentials found for bucket %s", bucketName)
}
// Create credentials
creds := credentials.NewStaticV4(accessKeyID, secretAccessKey, "")
// Cache credentials for 1 hour
utils.GlobalCache.Set(cacheKey, creds, time.Hour)
return creds, nil
return nil, fmt.Errorf("no valid credentials found for bucket %s", bucketName)
}
// getMinioClient creates a MinIO client for a specific bucket with dynamic credentials
func (s *S3Service) getMinioClient(ctx context.Context, bucketName string) (*minio.Client, error) {
creds, err := s.getBucketCredentials(ctx, bucketName)
// getMinioClient creates a MinIO client for a specific bucket with credentials
// that satisfy op.
func (s *S3Service) getMinioClient(ctx context.Context, bucketName string, op Operation) (*minio.Client, error) {
creds, err := s.getBucketCredentials(ctx, bucketName, op)
if err != nil {
return nil, fmt.Errorf("cannot get credentials for bucket %s: %w", bucketName, err)
}
@@ -151,7 +172,7 @@ func (s *S3Service) ListBuckets(ctx context.Context) (*models.BucketListResponse
// CreateBucket creates a new bucket in Garage
func (s *S3Service) CreateBucket(ctx context.Context, bucketName string) error {
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpRead|OpWrite)
if err != nil {
return fmt.Errorf("failed to get MinIO client for bucket %s: %w", bucketName, err)
}
@@ -172,7 +193,7 @@ func (s *S3Service) CreateBucket(ctx context.Context, bucketName string) error {
// DeleteBucket deletes a bucket from Garage
func (s *S3Service) DeleteBucket(ctx context.Context, bucketName string) error {
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpRead|OpWrite)
if err != nil {
return fmt.Errorf("failed to get MinIO client for bucket %s: %w", bucketName, err)
}
@@ -192,7 +213,7 @@ func (s *S3Service) DeleteBucket(ctx context.Context, bucketName string) error {
// ListObjects lists objects in a bucket with optional prefix filter and pagination
func (s *S3Service) ListObjects(ctx context.Context, bucketName, prefix string, maxKeys int, continuationToken string) (*models.ObjectListResponse, error) {
// Get bucket-specific MinIO client
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpRead)
if err != nil {
return nil, fmt.Errorf("failed to get MinIO client for bucket %s: %w", bucketName, err)
}
@@ -313,7 +334,7 @@ func (s *S3Service) ListObjects(ctx context.Context, bucketName, prefix string,
// UploadObject uploads an object to a bucket
func (s *S3Service) UploadObject(ctx context.Context, bucketName, key string, body io.Reader, contentType string) (*models.ObjectUploadResponse, error) {
// Get bucket-specific MinIO client
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpWrite)
if err != nil {
return nil, fmt.Errorf("failed to get MinIO client for bucket %s: %w", bucketName, err)
}
@@ -351,7 +372,7 @@ func (s *S3Service) UploadObject(ctx context.Context, bucketName, key string, bo
// size=0 forces a single PutObject request with Content-Length: 0, which
// Garage accepts as a directory marker.
func (s *S3Service) CreateDirectoryMarker(ctx context.Context, bucketName, key string) (*models.ObjectUploadResponse, error) {
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpWrite)
if err != nil {
return nil, fmt.Errorf("failed to get MinIO client for bucket %s: %w", bucketName, err)
}
@@ -381,7 +402,7 @@ func (s *S3Service) CreateDirectoryMarker(ctx context.Context, bucketName, key s
// GetObject retrieves an object from a bucket
func (s *S3Service) GetObject(ctx context.Context, bucketName, key string) (io.ReadCloser, *models.ObjectInfo, error) {
// Get bucket-specific MinIO client
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpRead)
if err != nil {
return nil, nil, fmt.Errorf("failed to get MinIO client for bucket %s: %w", bucketName, err)
}
@@ -421,7 +442,7 @@ func (s *S3Service) GetObject(ctx context.Context, bucketName, key string) (io.R
// DeleteObject deletes an object from a bucket
func (s *S3Service) DeleteObject(ctx context.Context, bucketName, key string) error {
// Get bucket-specific MinIO client
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpWrite)
if err != nil {
return fmt.Errorf("failed to get MinIO client for bucket %s: %w", bucketName, err)
}
@@ -441,7 +462,7 @@ func (s *S3Service) DeleteObject(ctx context.Context, bucketName, key string) er
// ObjectExists checks if an object exists in a bucket
func (s *S3Service) ObjectExists(ctx context.Context, bucketName, key string) (bool, error) {
// Get bucket-specific MinIO client
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpRead)
if err != nil {
return false, fmt.Errorf("failed to get MinIO client for bucket %s: %w", bucketName, err)
}
@@ -469,7 +490,7 @@ func (s *S3Service) ObjectExists(ctx context.Context, bucketName, key string) (b
// GetObjectMetadata retrieves metadata for an object without downloading it
func (s *S3Service) GetObjectMetadata(ctx context.Context, bucketName, key string) (*models.ObjectInfo, error) {
// Get bucket-specific MinIO client
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpRead)
if err != nil {
return nil, fmt.Errorf("failed to get MinIO client for bucket %s: %w", bucketName, err)
}
@@ -505,7 +526,7 @@ func (s *S3Service) DeleteMultipleObjects(ctx context.Context, bucketName string
}
// Get bucket-specific MinIO client
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpWrite)
if err != nil {
return fmt.Errorf("failed to get MinIO client for bucket %s: %w", bucketName, err)
}
@@ -540,7 +561,7 @@ func (s *S3Service) DeleteMultipleObjects(ctx context.Context, bucketName string
// This is useful for sharing files without exposing credentials
func (s *S3Service) GetPresignedURL(ctx context.Context, bucketName, key string, expiresIn time.Duration) (string, error) {
// Get bucket-specific MinIO client
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpRead)
if err != nil {
return "", fmt.Errorf("failed to get MinIO client for bucket %s: %w", bucketName, err)
}
@@ -579,7 +600,7 @@ func (s *S3Service) UploadMultipleObjects(ctx context.Context, bucketName string
results := make([]UploadResult, len(files))
// Get bucket-specific MinIO client once for all uploads
client, err := s.getMinioClient(ctx, bucketName)
client, err := s.getMinioClient(ctx, bucketName, OpWrite)
if err != nil {
// If we can't get the client, all uploads fail
for i := range files {
+192 -8
View File
@@ -3,6 +3,7 @@ package services
import (
"context"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"strings"
@@ -80,7 +81,9 @@ func uniqueBucket(t *testing.T) string {
t.Helper()
name := "test-bucket-" + t.Name()
t.Cleanup(func() {
utils.GlobalCache.Delete("key:" + name)
for _, op := range []Operation{OpRead, OpWrite, OpRead | OpWrite} {
utils.GlobalCache.Delete(fmt.Sprintf("key:%s:%d", name, op))
}
})
return name
}
@@ -111,7 +114,7 @@ func TestGetBucketCredentials_HappyPath(t *testing.T) {
})
s3, _ := adminBackedS3(t, mux)
creds, err := s3.getBucketCredentials(context.Background(), bucket)
creds, err := s3.getBucketCredentials(context.Background(), bucket, OpRead|OpWrite)
if err != nil {
t.Fatalf("getBucketCredentials: %v", err)
}
@@ -152,7 +155,7 @@ func TestGetBucketCredentials_CachesAcrossCalls(t *testing.T) {
s3, _ := adminBackedS3(t, mux)
for i := range 3 {
if _, err := s3.getBucketCredentials(context.Background(), bucket); err != nil {
if _, err := s3.getBucketCredentials(context.Background(), bucket, OpRead|OpWrite); err != nil {
t.Fatalf("call %d: %v", i, err)
}
}
@@ -164,7 +167,82 @@ func TestGetBucketCredentials_CachesAcrossCalls(t *testing.T) {
}
}
func TestGetBucketCredentials_SkipsKeysWithoutReadOrWrite(t *testing.T) {
func TestGetBucketCredentials_RWKeyWarmsAllTiers(t *testing.T) {
bucket := uniqueBucket(t)
secret := "rw-secret"
var bucketCalls, keyCalls int
mux := http.NewServeMux()
mux.HandleFunc("/v2/GetBucketInfo", func(w http.ResponseWriter, r *http.Request) {
bucketCalls++
_ = json.NewEncoder(w).Encode(&models.GarageBucketInfo{
ID: "bid",
Keys: []models.BucketKeyInfo{
{AccessKeyID: "RW", Permissions: models.BucketKeyPermission{Read: true, Write: true}},
},
})
})
mux.HandleFunc("/v2/GetKeyInfo", func(w http.ResponseWriter, r *http.Request) {
keyCalls++
_ = json.NewEncoder(w).Encode(&models.GarageKeyInfo{
AccessKeyID: "RW",
SecretAccessKey: &secret,
})
})
s3, _ := adminBackedS3(t, mux)
// Prime via OpRead — should populate OpRead, OpWrite, and OpRead|OpWrite.
if _, err := s3.getBucketCredentials(context.Background(), bucket, OpRead); err != nil {
t.Fatalf("prime OpRead: %v", err)
}
for _, op := range []Operation{OpWrite, OpRead | OpWrite, OpRead} {
if _, err := s3.getBucketCredentials(context.Background(), bucket, op); err != nil {
t.Fatalf("op %d: %v", op, err)
}
}
if bucketCalls != 1 {
t.Errorf("GetBucketInfo called %d times, want 1 (RW key should warm every tier)", bucketCalls)
}
if keyCalls != 1 {
t.Errorf("GetKeyInfo called %d times, want 1", keyCalls)
}
}
// A read-only key must NOT populate the write or RW cache slots, otherwise an
// OpWrite call would receive credentials the cluster will reject.
func TestGetBucketCredentials_ReadOnlyKeyDoesNotPoisonWriteCache(t *testing.T) {
bucket := uniqueBucket(t)
secret := "ro-secret"
mux := http.NewServeMux()
mux.HandleFunc("/v2/GetBucketInfo", func(w http.ResponseWriter, r *http.Request) {
_ = json.NewEncoder(w).Encode(&models.GarageBucketInfo{
ID: "bid",
Keys: []models.BucketKeyInfo{
{AccessKeyID: "READ-ONLY", Permissions: models.BucketKeyPermission{Read: true, Write: false}},
},
})
})
mux.HandleFunc("/v2/GetKeyInfo", func(w http.ResponseWriter, r *http.Request) {
_ = json.NewEncoder(w).Encode(&models.GarageKeyInfo{
AccessKeyID: "READ-ONLY",
SecretAccessKey: &secret,
})
})
s3, _ := adminBackedS3(t, mux)
// Warm OpRead cache with the read-only key.
if _, err := s3.getBucketCredentials(context.Background(), bucket, OpRead); err != nil {
t.Fatalf("prime OpRead: %v", err)
}
// OpWrite must still fail — the read-only key must not have leaked into
// the write cache slot.
if _, err := s3.getBucketCredentials(context.Background(), bucket, OpWrite); err == nil {
t.Fatal("OpWrite served credentials from a read-only key; cache was poisoned")
}
}
func TestGetBucketCredentials_OpReadWriteSkipsKeysMissingAnyBit(t *testing.T) {
bucket := uniqueBucket(t)
secret := "good-secret"
@@ -191,7 +269,7 @@ func TestGetBucketCredentials_SkipsKeysWithoutReadOrWrite(t *testing.T) {
})
s3, _ := adminBackedS3(t, mux)
creds, err := s3.getBucketCredentials(context.Background(), bucket)
creds, err := s3.getBucketCredentials(context.Background(), bucket, OpRead|OpWrite)
if err != nil {
t.Fatalf("getBucketCredentials: %v", err)
}
@@ -204,6 +282,112 @@ func TestGetBucketCredentials_SkipsKeysWithoutReadOrWrite(t *testing.T) {
}
}
// Regression test for issue #44: read-only buckets must remain browsable when
// no read+write key is assigned. Before the fix, this case returned
// "no valid credentials found for bucket music" and the UI broke entirely.
func TestGetBucketCredentials_ReadOnlyFallsBackToReadKey(t *testing.T) {
bucket := uniqueBucket(t)
secret := "ro-secret"
mux := http.NewServeMux()
mux.HandleFunc("/v2/GetBucketInfo", func(w http.ResponseWriter, r *http.Request) {
_ = json.NewEncoder(w).Encode(&models.GarageBucketInfo{
ID: "bid",
Keys: []models.BucketKeyInfo{
{AccessKeyID: "READ-ONLY", Permissions: models.BucketKeyPermission{Read: true, Write: false}},
},
})
})
mux.HandleFunc("/v2/GetKeyInfo", func(w http.ResponseWriter, r *http.Request) {
_ = json.NewEncoder(w).Encode(&models.GarageKeyInfo{
AccessKeyID: "READ-ONLY",
SecretAccessKey: &secret,
})
})
s3, _ := adminBackedS3(t, mux)
creds, err := s3.getBucketCredentials(context.Background(), bucket, OpRead)
if err != nil {
t.Fatalf("getBucketCredentials: %v", err)
}
v, err := creds.GetWithContext(nil)
if err != nil {
t.Fatalf("creds.GetWithContext: %v", err)
}
if v.AccessKeyID != "READ-ONLY" {
t.Errorf("AccessKeyID = %q, want READ-ONLY", v.AccessKeyID)
}
}
// Even with only a read-only key available, asking for write credentials must
// still fail loudly so uploads/deletes return a meaningful error instead of
// silently using a key that the cluster will reject.
func TestGetBucketCredentials_ReadOnlyBucketRejectsWriteRequest(t *testing.T) {
bucket := uniqueBucket(t)
secret := "ro-secret"
mux := http.NewServeMux()
mux.HandleFunc("/v2/GetBucketInfo", func(w http.ResponseWriter, r *http.Request) {
_ = json.NewEncoder(w).Encode(&models.GarageBucketInfo{
ID: "bid",
Keys: []models.BucketKeyInfo{
{AccessKeyID: "READ-ONLY", Permissions: models.BucketKeyPermission{Read: true, Write: false}},
},
})
})
mux.HandleFunc("/v2/GetKeyInfo", func(w http.ResponseWriter, r *http.Request) {
_ = json.NewEncoder(w).Encode(&models.GarageKeyInfo{
AccessKeyID: "READ-ONLY",
SecretAccessKey: &secret,
})
})
s3, _ := adminBackedS3(t, mux)
_, err := s3.getBucketCredentials(context.Background(), bucket, OpRead|OpWrite)
if err == nil {
t.Fatal("expected error when only a read-only key exists, got nil")
}
if !strings.Contains(err.Error(), "no valid credentials") {
t.Errorf("expected 'no valid credentials' in error, got %v", err)
}
}
// Mirror of issue #44 for write-only buckets: uploads must still succeed with a
// write-only key, even though no key grants read access.
func TestGetBucketCredentials_WriteOnlyFallsBackToWriteKey(t *testing.T) {
bucket := uniqueBucket(t)
secret := "wo-secret"
mux := http.NewServeMux()
mux.HandleFunc("/v2/GetBucketInfo", func(w http.ResponseWriter, r *http.Request) {
_ = json.NewEncoder(w).Encode(&models.GarageBucketInfo{
ID: "bid",
Keys: []models.BucketKeyInfo{
{AccessKeyID: "WRITE-ONLY", Permissions: models.BucketKeyPermission{Read: false, Write: true}},
},
})
})
mux.HandleFunc("/v2/GetKeyInfo", func(w http.ResponseWriter, r *http.Request) {
_ = json.NewEncoder(w).Encode(&models.GarageKeyInfo{
AccessKeyID: "WRITE-ONLY",
SecretAccessKey: &secret,
})
})
s3, _ := adminBackedS3(t, mux)
creds, err := s3.getBucketCredentials(context.Background(), bucket, OpWrite)
if err != nil {
t.Fatalf("getBucketCredentials: %v", err)
}
v, err := creds.GetWithContext(nil)
if err != nil {
t.Fatalf("creds.GetWithContext: %v", err)
}
if v.AccessKeyID != "WRITE-ONLY" {
t.Errorf("AccessKeyID = %q, want WRITE-ONLY", v.AccessKeyID)
}
}
func TestGetBucketCredentials_NoEligibleKeyReturnsError(t *testing.T) {
bucket := uniqueBucket(t)
@@ -219,7 +403,7 @@ func TestGetBucketCredentials_NoEligibleKeyReturnsError(t *testing.T) {
})
s3, _ := adminBackedS3(t, mux)
_, err := s3.getBucketCredentials(context.Background(), bucket)
_, err := s3.getBucketCredentials(context.Background(), bucket, OpRead)
if err == nil {
t.Fatal("expected error when bucket has no keys, got nil")
}
@@ -254,7 +438,7 @@ func TestGetBucketCredentials_KeyWithoutSecretIsSkipped(t *testing.T) {
})
s3, _ := adminBackedS3(t, mux)
creds, err := s3.getBucketCredentials(context.Background(), bucket)
creds, err := s3.getBucketCredentials(context.Background(), bucket, OpRead|OpWrite)
if err != nil {
t.Fatalf("getBucketCredentials: %v", err)
}
@@ -277,7 +461,7 @@ func TestGetBucketCredentials_AdminErrorPropagates(t *testing.T) {
})
s3, _ := adminBackedS3(t, mux)
_, err := s3.getBucketCredentials(context.Background(), bucket)
_, err := s3.getBucketCredentials(context.Background(), bucket, OpRead|OpWrite)
if err == nil {
t.Fatal("expected error when admin call fails, got nil")
}