s3api: resolve volume data encryption in CopyObject SSE flows (#11646) (#11683)

* s3api: resolve volume data encryption in CopyObject SSE flows (#11646)

* s3api: honor bucket-default KMS key on copy and fix transformed-upload metadata

- Synthesize the destination bucket's default encryption as request
  headers before any SSE evaluation, so the configured KMS key ID and
  bucket-key setting reach the copy paths instead of only a boolean.
- uploadTransformedChunkData returns the upload result so callers record
  the uploader's cipher key AND compression decision; a wrongly cleared
  IsCompressed made transformed copies unreadable.
- decompressChunkVolumeCipher fails loudly when a compressed chunk does
  not decompress, instead of uploading still-compressed bytes marked
  uncompressed.
- copyMultipartSSECChunk now strips the volume cipher and re-encrypts on
  upload like the other transform paths.
- The ciphered inner upload now uses the caller's private BytesBuffer.

* s3api: extract copy bucket-default header synthesis for coverage

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* s3api: test bucket-default encryption header synthesis on copy

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* s3api: validate copy encryption headers before bucket defaults and resolve empty KMS key

* s3api: name the AWS-managed SSE-KMS default key once

---------

Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
authored and GitHub committed 2026-10-10 10:59:36 +08:00
1 parent ea65746947
commit d3dc03c85a
8 files changed
+364 -6

No files matched your search

+4
View File
@@ -10,6 +10,10 @@ const (
SSEAlgorithmAES256 = "AES256"
SSEAlgorithmKMS = "aws:kms"
// SSEKMSDefaultKeyID is the AWS-managed key SSE-KMS requests resolve to
// when no key ID is given, matching AWS's aws/s3 managed key alias.
SSEKMSDefaultKeyID = "alias/aws/s3"
// SSE type identifiers for response headers and internal processing
SSETypeC = "SSE-C"
SSETypeKMS = "SSE-KMS"
+4
View File
@@ -865,6 +865,10 @@ func DetermineUnifiedCopyStrategy(state *EncryptionState, srcMetadata map[string
return CopyStrategyReencrypt, nil
}
if state.SrcSSES3 && state.DstSSES3 {
return CopyStrategyDirect, nil
}
// Encrypt: plain → encrypted
if !state.IsSourceEncrypted() && state.IsTargetEncrypted() {
return CopyStrategyEncrypt, nil
@@ -268,3 +268,109 @@ func TestGetEncryptedStreamFromVolumesRangesInlineContent(t *testing.T) {
})
}
}
func TestDetermineUnifiedCopyStrategySSES3(t *testing.T) {
testCases := []struct {
name string
state *EncryptionState
expected UnifiedCopyStrategy
}{
{
name: "SSE-S3 to SSE-S3 direct copy",
state: &EncryptionState{
SrcSSES3: true,
DstSSES3: true,
},
expected: CopyStrategyDirect,
},
{
name: "SSE-S3 to plain decrypt copy",
state: &EncryptionState{
SrcSSES3: true,
DstSSES3: false,
},
expected: CopyStrategyDecrypt,
},
{
name: "Plain to SSE-S3 encrypt copy",
state: &EncryptionState{
SrcSSES3: false,
DstSSES3: true,
},
expected: CopyStrategyEncrypt,
},
}
for _, tc := range testCases {
t.Run(tc.name, func(t *testing.T) {
strategy, err := DetermineUnifiedCopyStrategy(tc.state, nil, nil)
if err != nil {
t.Fatalf("DetermineUnifiedCopyStrategy failed: %v", err)
}
if strategy != tc.expected {
t.Errorf("expected strategy %v, got %v", tc.expected, strategy)
}
})
}
}
func TestDecryptChunkVolumeCipher(t *testing.T) {
s3a := &S3ApiServer{}
plainData := []byte("hello-seaweedfs-encrypted-volume-data-verification-test")
t.Run("unencrypted chunk passes through", func(t *testing.T) {
chunk := &filer_pb.FileChunk{
CipherKey: nil,
}
got, err := s3a.decryptChunkVolumeCipher(plainData, chunk)
if err != nil {
t.Fatalf("decryptChunkVolumeCipher failed: %v", err)
}
if !bytes.Equal(got, plainData) {
t.Fatalf("expected %q, got %q", plainData, got)
}
})
t.Run("volume cipher encrypted chunk decrypts correctly", func(t *testing.T) {
cipherKey := util.GenCipherKey()
encrypted, err := util.Encrypt(plainData, cipherKey)
if err != nil {
t.Fatalf("Encrypt failed: %v", err)
}
chunk := &filer_pb.FileChunk{
CipherKey: cipherKey,
}
got, err := s3a.decryptChunkVolumeCipher(encrypted, chunk)
if err != nil {
t.Fatalf("decryptChunkVolumeCipher failed: %v", err)
}
if !bytes.Equal(got, plainData) {
t.Fatalf("expected %q, got %q", plainData, got)
}
})
t.Run("volume cipher encrypted and compressed chunk decrypts correctly", func(t *testing.T) {
cipherKey := util.GenCipherKey()
compressed, err := util.GzipData(plainData)
if err != nil {
t.Fatalf("GzipData failed: %v", err)
}
encrypted, err := util.Encrypt(compressed, cipherKey)
if err != nil {
t.Fatalf("Encrypt failed: %v", err)
}
chunk := &filer_pb.FileChunk{
CipherKey: cipherKey,
IsCompressed: true,
}
got, err := s3a.decryptChunkVolumeCipher(encrypted, chunk)
if err != nil {
t.Fatalf("decryptChunkVolumeCipher failed: %v", err)
}
if !bytes.Equal(got, plainData) {
t.Fatalf("expected %q, got %q", plainData, got)
}
})
}
+115 -5
View File
@@ -159,6 +159,30 @@ func (s3a *S3ApiServer) CopyObjectHandler(w http.ResponseWriter, r *http.Request
return
}
// Validate the client's encryption headers as sent: synthesizing bucket
// defaults first would launder a malformed request (a KMS option without
// aws:kms, or a repeated header) into an accepted one.
if errCode := ValidateRequestEncryption(r.Header); errCode != s3err.ErrNone {
s3err.WriteErrorResponse(w, r, errCode)
return
}
// A copy request without SSE headers takes the destination bucket's
// default encryption, including its configured KMS key and bucket-key
// setting. Surface them as headers so every downstream parse sees them.
if !IsSSECRequest(r) && r.Header.Get(s3_constants.AmzServerSideEncryption) == "" {
encryptionConfig, encErr := s3a.GetBucketEncryptionConfig(dstBucket)
switch {
case encErr == nil:
applyCopyBucketDefaultEncryptionHeaders(r, encryptionConfig)
case errors.Is(encErr, ErrNoEncryptionConfig):
default:
glog.Errorf("CopyObjectHandler: read encryption config for bucket %s: %v", dstBucket, encErr)
s3err.WriteErrorResponse(w, r, s3err.ErrInternalError)
return
}
}
if !isValidDirective(r.Header.Get(s3_constants.AmzUserMetaDirective)) {
s3err.WriteErrorResponse(w, r, s3err.ErrInvalidMetadataDirective)
return
@@ -1783,6 +1807,47 @@ func (s3a *S3ApiServer) uploadChunkData(chunkData []byte, assignResult *filer_pb
return nil
}
func (s3a *S3ApiServer) uploadTransformedChunkData(chunkData []byte, assignResult *filer_pb.AssignVolumeResponse) (*operation.UploadResult, error) {
dstUrl := fmt.Sprintf("http://%s/%s", assignResult.Location.Url, assignResult.FileId)
if assignResult.Fsync {
dstUrl += "?fsync=true"
}
uploadOption := &operation.UploadOption{
UploadUrl: dstUrl,
Cipher: s3a.cipher,
IsInputCompressed: false,
MimeType: "",
PairMap: nil,
Jwt: security.EncodedJwt(assignResult.Auth),
BytesBuffer: bytes.NewBuffer(make([]byte, 0, len(chunkData))),
}
uploader, err := operation.NewUploader()
if err != nil {
return nil, fmt.Errorf("create uploader: %w", err)
}
uploadResult, err := uploader.UploadData(context.Background(), chunkData, uploadOption)
if err != nil {
return nil, fmt.Errorf("upload transformed chunk: %w", err)
}
return uploadResult, nil
}
func (s3a *S3ApiServer) decryptChunkVolumeCipher(data []byte, chunk *filer_pb.FileChunk) ([]byte, error) {
if len(chunk.CipherKey) == 0 {
return data, nil
}
decrypted, err := util.Decrypt(data, util.CipherKey(chunk.CipherKey))
if err != nil {
return nil, fmt.Errorf("decrypt volume cipher: %w", err)
}
if chunk.IsCompressed {
if decrypted, err = util.DecompressData(decrypted); err != nil {
return nil, fmt.Errorf("decompress chunk: %w", err)
}
}
return decrypted, nil
}
// multipartFramingOverhead reserves space for the multipart wrapper
// upload_content writes around chunkData (boundary + Content-Disposition +
// optional Content-Type/Content-Encoding/Content-MD5 headers + trailing
@@ -2012,6 +2077,10 @@ func (s3a *S3ApiServer) copyMultipartSSEKMSChunk(chunk *filer_pb.FileChunk, sour
if err != nil {
return nil, fmt.Errorf("download encrypted chunk data: %w", err)
}
encryptedData, err = s3a.decryptChunkVolumeCipher(encryptedData, chunk)
if err != nil {
return nil, err
}
var finalData []byte
@@ -2074,9 +2143,14 @@ func (s3a *S3ApiServer) copyMultipartSSEKMSChunk(chunk *filer_pb.FileChunk, sour
}
// Upload the final data
if err := s3a.uploadChunkData(finalData, assignResult, false); err != nil {
uploadResult, err := s3a.uploadTransformedChunkData(finalData, assignResult)
if err != nil {
return nil, fmt.Errorf("upload chunk data: %w", err)
}
if uploadResult != nil {
dstChunk.CipherKey = uploadResult.CipherKey
dstChunk.IsCompressed = uploadResult.Gzip > 0
}
// Update chunk size
dstChunk.Size = uint64(len(finalData))
@@ -2109,6 +2183,10 @@ func (s3a *S3ApiServer) copyMultipartSSECChunk(chunk *filer_pb.FileChunk, copySo
if err != nil {
return nil, nil, fmt.Errorf("download encrypted chunk data: %w", err)
}
encryptedData, err = s3a.decryptChunkVolumeCipher(encryptedData, chunk)
if err != nil {
return nil, nil, err
}
var finalData []byte
var destIV []byte
@@ -2204,9 +2282,14 @@ func (s3a *S3ApiServer) copyMultipartSSECChunk(chunk *filer_pb.FileChunk, copySo
}
// Upload the final data
if err := s3a.uploadChunkData(finalData, assignResult, false); err != nil {
uploadResult, err := s3a.uploadTransformedChunkData(finalData, assignResult)
if err != nil {
return nil, nil, fmt.Errorf("upload chunk data: %w", err)
}
if uploadResult != nil {
dstChunk.CipherKey = uploadResult.CipherKey
dstChunk.IsCompressed = uploadResult.Gzip > 0
}
// Update chunk size
dstChunk.Size = uint64(len(finalData))
@@ -2403,6 +2486,10 @@ func (s3a *S3ApiServer) copyCrossEncryptionChunk(chunk *filer_pb.FileChunk, sour
if err != nil {
return nil, fmt.Errorf("download encrypted chunk data: %w", err)
}
encryptedData, err = s3a.decryptChunkVolumeCipher(encryptedData, chunk)
if err != nil {
return nil, err
}
var finalData []byte
@@ -2593,9 +2680,14 @@ func (s3a *S3ApiServer) copyCrossEncryptionChunk(chunk *filer_pb.FileChunk, sour
// For unencrypted destination, finalData remains as decrypted plaintext
// Upload the final data
if err := s3a.uploadChunkData(finalData, assignResult, false); err != nil {
uploadResult, err := s3a.uploadTransformedChunkData(finalData, assignResult)
if err != nil {
return nil, fmt.Errorf("upload chunk data: %w", err)
}
if uploadResult != nil {
dstChunk.CipherKey = uploadResult.CipherKey
dstChunk.IsCompressed = uploadResult.Gzip > 0
}
// Update chunk size
dstChunk.Size = uint64(len(finalData))
@@ -2748,6 +2840,10 @@ func (s3a *S3ApiServer) copyChunkWithReencryption(chunk *filer_pb.FileChunk, cop
if err != nil {
return nil, fmt.Errorf("download encrypted chunk data: %w", err)
}
encryptedData, err = s3a.decryptChunkVolumeCipher(encryptedData, chunk)
if err != nil {
return nil, err
}
var finalData []byte
@@ -2806,9 +2902,14 @@ func (s3a *S3ApiServer) copyChunkWithReencryption(chunk *filer_pb.FileChunk, cop
}
// Upload the processed data
if err := s3a.uploadChunkData(finalData, assignResult, false); err != nil {
uploadResult, err := s3a.uploadTransformedChunkData(finalData, assignResult)
if err != nil {
return nil, fmt.Errorf("upload processed chunk data: %w", err)
}
if uploadResult != nil {
dstChunk.CipherKey = uploadResult.CipherKey
dstChunk.IsCompressed = uploadResult.Gzip > 0
}
return dstChunk, nil
}
@@ -3079,6 +3180,10 @@ func (s3a *S3ApiServer) copyChunkWithSSEKMSReencryption(chunk *filer_pb.FileChun
if err != nil {
return nil, fmt.Errorf("download chunk data: %w", err)
}
chunkData, err = s3a.decryptChunkVolumeCipher(chunkData, chunk)
if err != nil {
return nil, err
}
var finalData []byte
@@ -3162,9 +3267,14 @@ func (s3a *S3ApiServer) copyChunkWithSSEKMSReencryption(chunk *filer_pb.FileChun
}
// Upload the processed data
if err := s3a.uploadChunkData(finalData, assignResult, false); err != nil {
uploadResult, err := s3a.uploadTransformedChunkData(finalData, assignResult)
if err != nil {
return nil, fmt.Errorf("upload processed chunk data: %w", err)
}
if uploadResult != nil {
dstChunk.CipherKey = uploadResult.CipherKey
dstChunk.IsCompressed = uploadResult.Gzip > 0
}
glog.V(3).Infof("Successfully processed SSE-KMS chunk re-encryption: src_key=%s, dst_key=%s, size=%d→%d",
getKeyIDString(sourceSSEKey), destKeyID, len(chunkData), len(finalData))
@@ -0,0 +1,105 @@
package s3api
import (
"net/http"
"net/http/httptest"
"testing"
"github.com/gorilla/mux"
"github.com/seaweedfs/seaweedfs/weed/pb/s3_pb"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
)
// Multi-chunk copy paths read SSE settings from request headers only, so the
// destination bucket default must be synthesized into headers to reach them.
func TestApplyCopyBucketDefaultEncryptionHeaders(t *testing.T) {
tests := []struct {
name string
cfg *s3_pb.EncryptionConfiguration
wantSse string
wantKmsKey string
wantBktKey string
}{
{name: "nil config", cfg: nil},
{
name: "kms with key and bucket key",
cfg: &s3_pb.EncryptionConfiguration{
SseAlgorithm: "aws:kms",
KmsKeyId: "arn:aws:kms:us-east-1:123:key/abc",
BucketKeyEnabled: true,
},
wantSse: "aws:kms",
wantKmsKey: "arn:aws:kms:us-east-1:123:key/abc",
wantBktKey: "true",
},
{
name: "kms without key resolves the aws/s3 default",
cfg: &s3_pb.EncryptionConfiguration{
SseAlgorithm: "aws:kms",
},
wantSse: "aws:kms",
wantKmsKey: "alias/aws/s3",
},
{
name: "aes256",
cfg: &s3_pb.EncryptionConfiguration{SseAlgorithm: "AES256"},
wantSse: "AES256",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
r, _ := http.NewRequest("PUT", "/dst", nil)
applyCopyBucketDefaultEncryptionHeaders(r, tt.cfg)
if got := r.Header.Get(s3_constants.AmzServerSideEncryption); got != tt.wantSse {
t.Fatalf("SSE header = %q, want %q", got, tt.wantSse)
}
if got := r.Header.Get(s3_constants.AmzServerSideEncryptionAwsKmsKeyId); got != tt.wantKmsKey {
t.Fatalf("KMS key header = %q, want %q", got, tt.wantKmsKey)
}
if got := r.Header.Get(s3_constants.AmzServerSideEncryptionBucketKeyEnabled); got != tt.wantBktKey {
t.Fatalf("bucket-key header = %q, want %q", got, tt.wantBktKey)
}
})
}
}
// The client's encryption headers must validate before bucket defaults are
// synthesized: a KMS option without aws:kms, or a repeated header, would
// otherwise be laundered into an accepted request by the default headers.
func TestCopyObjectHandlerRejectsMalformedEncryptionBeforeDefaults(t *testing.T) {
for _, tt := range []struct {
name string
headers map[string][]string
}{
{
name: "kms key id without algorithm",
headers: map[string][]string{
s3_constants.AmzServerSideEncryptionAwsKmsKeyId: {"arn:aws:kms:us-east-1:123:key/abc"},
},
},
{
name: "repeated algorithm header",
headers: map[string][]string{
s3_constants.AmzServerSideEncryption: {"AES256", "aws:kms"},
},
},
} {
t.Run(tt.name, func(t *testing.T) {
s3a := newHeadBucketTestServer(t, &fakeLookupFiler{})
r, _ := http.NewRequest(http.MethodPut, "/dst-b/o", nil)
r = mux.SetURLVars(r, map[string]string{"bucket": "dst-b", "object": "o"})
r.Header.Set("X-Amz-Copy-Source", "/src-b/k")
for name, values := range tt.headers {
for _, v := range values {
r.Header.Add(name, v)
}
}
w := httptest.NewRecorder()
s3a.CopyObjectHandler(w, r)
if w.Code != http.StatusBadRequest {
t.Fatalf("status = %d, want 400; body: %s", w.Code, w.Body.String())
}
})
}
}
@@ -7,6 +7,8 @@ import (
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/s3_pb"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
weed_server "github.com/seaweedfs/seaweedfs/weed/server"
)
@@ -185,3 +187,29 @@ func (s3a *S3ApiServer) applyCopyBucketDefaultEncryption(state *EncryptionState,
}
}
}
// applyCopyBucketDefaultEncryptionHeaders surfaces a bucket's default
// encryption as request headers, so multi-chunk copy paths that only read
// request headers honor it too.
func applyCopyBucketDefaultEncryptionHeaders(r *http.Request, cfg *s3_pb.EncryptionConfiguration) {
if cfg == nil {
return
}
switch cfg.SseAlgorithm {
case EncryptionTypeKMS:
r.Header.Set(s3_constants.AmzServerSideEncryption, "aws:kms")
// aws:kms with no configured key resolves the AWS default, the same
// fallback applySSEKMSDefaultEncryption uses for uploads; an empty
// key ID here would let the copy paths take their plaintext branch.
keyID := cfg.KmsKeyId
if keyID == "" {
keyID = s3_constants.SSEKMSDefaultKeyID
}
r.Header.Set(s3_constants.AmzServerSideEncryptionAwsKmsKeyId, keyID)
if cfg.BucketKeyEnabled {
r.Header.Set(s3_constants.AmzServerSideEncryptionBucketKeyEnabled, "true")
}
case EncryptionTypeAES256:
r.Header.Set(s3_constants.AmzServerSideEncryption, "AES256")
}
}
+1 -1
View File
@@ -2246,7 +2246,7 @@ func (s3a *S3ApiServer) applySSEKMSDefaultEncryption(bucket string, r *http.Requ
// Use the KMS key ID from bucket configuration, or default if not specified
keyID := encryptionConfig.KmsKeyId
if keyID == "" {
keyID = "alias/aws/s3" // AWS default KMS key for S3
keyID = s3_constants.SSEKMSDefaultKeyID
}
// Check if bucket key is enabled in configuration