s3api: apply bucket default encryption when volume data encryption is enabled (#11681)

* s3api: apply bucket default encryption when volume data encryption is enabled

When -s3.encryptVolumeData (s3a.cipher) is enabled, putToFiler skipped
checking and applying bucket default encryption due to a '!s3a.cipher'
guard. Volume-level data encryption and object-level Server-Side
Encryption (SSE-S3 / SSE-KMS) operate at different layers, and explicit
SSE headers already work alongside volume encryption.

Remove the '!s3a.cipher' guard so PutObject without explicit SSE headers
inherits bucket default encryption regardless of volume data encryption.
Add regression test TestPutObjectAppliesBucketDefaultEncryptionWithVolumeCipher.

Signed-off-by: Tyagiquamar <mohdquamartyagi@gmail.com>

* s3api: decrypt the volume cipher on direct SSE chunk reads

fetchFullChunk, fetchChunkViewData, and createEncryptedChunkReader fetched
raw bytes over HTTP, so volume-encrypted chunks reached SSE-S3/KMS/C
decryptors still ciphered. Route them through fetchChunkData: encrypted or
compressed chunks go through RetriedFetchChunkData (cipher-aware, slices
plaintext space for views); plain chunks keep the streaming range read.

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

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

* s3api: cover cipher-aware chunk reads with a fake volume server

The fake volume now serves stored GETs with Range support, and
TestFetchChunkDataDecryptsVolumeCipher verifies full-chunk and view reads
return plaintext for ciphered chunks while plain chunks still slice via HTTP
ranges.

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

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

* s3api: guard fake volume server stored map with its mutex

The HTTP handler goroutine read v.stored while test goroutines wrote
it, a data race go test -race can flag. Lock v.mu around the map read
and the test writes.

---------

Signed-off-by: Tyagiquamar <mohdquamartyagi@gmail.com>
Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
Co-authored-by: Chris Lu <chris.lu@gmail.com>
This commit is contained in:
authored and GitHub committed 2026-10-10 23:41:17 +08:00
1 parent c96eb7b8ee
commit d050464316
4 files changed
+243 -84

No files matched your search

@@ -0,0 +1,116 @@
package s3api
import (
"bytes"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"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"
)
// Issue 11647: Bucket default encryption (SSE-S3) must be applied when
// -s3.encryptVolumeData (s3a.cipher) is enabled, matching explicit SSE headers.
func TestPutObjectAppliesBucketDefaultEncryptionWithVolumeCipher(t *testing.T) {
// Configure test key manager with a super key for SSE-S3 encryption.
km := GetSSES3KeyManager()
oldSuperKey := km.superKey
km.superKey = make([]byte, 32)
for i := range km.superKey {
km.superKey[i] = byte(i + 1)
}
t.Cleanup(func() {
km.superKey = oldSuperKey
})
testCases := []struct {
name string
enableCipher bool
}{
{
name: "volume encryption enabled (s3a.cipher=true)",
enableCipher: true,
},
{
name: "volume encryption disabled (s3a.cipher=false)",
enableCipher: false,
},
}
for _, tc := range testCases {
t.Run(tc.name, func(t *testing.T) {
volume := startFakeVolumeServer(t)
filerImpl := &ambiguousPutFiler{
volume: volume,
entries: map[string]*filer_pb.Entry{},
apply: true,
}
s3a := newPutTestServer(t, startFakeFiler(t, filerImpl))
s3a.cipher = tc.enableCipher
// Configure bucket "b" with default AES256 (SSE-S3) encryption.
s3a.bucketConfigCache = NewBucketConfigCache(time.Minute)
s3a.bucketConfigCache.Set("b", &BucketConfig{
Name: "b",
Encryption: &s3_pb.EncryptionConfiguration{
SseAlgorithm: "AES256",
},
})
// PUT request without explicit SSE headers.
r := httptest.NewRequest(http.MethodPut, "/b/plain.txt", nil)
filePath := "/buckets/b/plain.txt"
etag, code, sseMeta := s3a.putToFiler(r, filePath, strings.NewReader("hello seaweedfs"), "b", "plain.txt", 1, 0, nil, false, "")
if code != s3err.ErrNone {
t.Fatalf("putToFiler returned error code %v, want %v", code, s3err.ErrNone)
}
if etag == "" {
t.Fatal("expected non-empty etag")
}
// Verify returned SSE response metadata.
if sseMeta.SSEType != s3_constants.SSETypeS3 {
t.Fatalf("expected SSE response metadata type %s, got %s", s3_constants.SSETypeS3, sseMeta.SSEType)
}
// Verify entry saved on the filer.
entry, ok := filerImpl.entries[filePath]
if !ok {
t.Fatalf("entry not found on filer at %s", filePath)
}
sseHeaderVal, hasSSE := entry.Extended[s3_constants.AmzServerSideEncryption]
if !hasSSE || !bytes.Equal(sseHeaderVal, []byte("AES256")) {
t.Fatalf("expected entry to have %s=AES256, got hasSSE=%v val=%s", s3_constants.AmzServerSideEncryption, hasSSE, string(sseHeaderVal))
}
if len(entry.Extended[s3_constants.SeaweedFSSSES3Key]) == 0 {
t.Fatal("expected entry to have stored SSE-S3 key metadata")
}
if len(entry.Chunks) == 0 {
t.Fatal("expected entry to have at least one chunk")
}
for i, chunk := range entry.Chunks {
if chunk.SseType != filer_pb.SSEType_SSE_S3 {
t.Errorf("chunk %d: expected SseType SSE_S3, got %v", i, chunk.SseType)
}
if len(chunk.SseMetadata) == 0 {
t.Errorf("chunk %d: expected non-empty SseMetadata", i)
}
if tc.enableCipher && len(chunk.CipherKey) == 0 {
t.Errorf("chunk %d: expected volume CipherKey when s3a.cipher is true", i)
}
if !tc.enableCipher && len(chunk.CipherKey) != 0 {
t.Errorf("chunk %d: expected empty volume CipherKey when s3a.cipher is false", i)
}
}
})
}
}
+40 -82
View File
@@ -1944,7 +1944,7 @@ func (s3a *S3ApiServer) decryptSSECChunkView(ctx context.Context, fileChunk *fil
// Fetch FULL encrypted chunk
// Note: Fetching full chunk is necessary for proper CTR decryption stream
fullChunkReader, err := s3a.fetchFullChunk(ctx, chunkView.FileId)
fullChunkReader, err := s3a.fetchFullChunk(ctx, chunkView)
if err != nil {
return nil, fmt.Errorf("failed to fetch full chunk: %w", err)
}
@@ -2013,7 +2013,7 @@ func (s3a *S3ApiServer) decryptSSEKMSChunkView(ctx context.Context, fileChunk *f
}
// Fetch FULL encrypted chunk
fullChunkReader, err := s3a.fetchFullChunk(ctx, chunkView.FileId)
fullChunkReader, err := s3a.fetchFullChunk(ctx, chunkView)
if err != nil {
return nil, fmt.Errorf("failed to fetch full chunk: %w", err)
}
@@ -2078,7 +2078,7 @@ func (s3a *S3ApiServer) decryptSSES3ChunkView(ctx context.Context, fileChunk *fi
}
// Fetch FULL encrypted chunk (necessary for proper CTR decryption stream)
fullChunkReader, err := s3a.fetchFullChunk(ctx, chunkView.FileId)
fullChunkReader, err := s3a.fetchFullChunk(ctx, chunkView)
if err != nil {
return nil, fmt.Errorf("failed to fetch full chunk: %w", err)
}
@@ -2124,7 +2124,7 @@ func (s3a *S3ApiServer) decryptSSES3ChunkView(ctx context.Context, fileChunk *fi
}
// Fetch FULL encrypted chunk
fullChunkReader, err := s3a.fetchFullChunk(ctx, chunkView.FileId)
fullChunkReader, err := s3a.fetchFullChunk(ctx, chunkView)
if err != nil {
return nil, fmt.Errorf("failed to fetch full chunk: %w", err)
}
@@ -2160,8 +2160,37 @@ func (s3a *S3ApiServer) decryptSSES3ChunkView(ctx context.Context, fileChunk *fi
return &rc{Reader: limitedReader, Closer: fullChunkReader}, nil
}
// fetchFullChunk fetches the complete encrypted chunk from volume server
func (s3a *S3ApiServer) fetchFullChunk(ctx context.Context, fileId string) (io.ReadCloser, error) {
// fetchFullChunk fetches the complete chunk from the volume server
func (s3a *S3ApiServer) fetchFullChunk(ctx context.Context, chunkView *filer.ChunkView) (io.ReadCloser, error) {
return s3a.fetchChunkData(ctx, chunkView.FileId, chunkView.CipherKey, chunkView.IsGzipped, 0, int64(chunkView.ChunkSize), true)
}
// fetchChunkViewData fetches data for a chunk view (with range)
func (s3a *S3ApiServer) fetchChunkViewData(ctx context.Context, chunkView *filer.ChunkView) (io.ReadCloser, error) {
return s3a.fetchChunkData(ctx, chunkView.FileId, chunkView.CipherKey, chunkView.IsGzipped, chunkView.OffsetInChunk, int64(chunkView.ViewSize), chunkView.IsFullChunk())
}
// fetchChunkData reads [offset, offset+size) of a chunk in plaintext space:
// the volume cipher is decrypted and compression undone by
// RetriedFetchChunkData. Unencrypted chunks take the streaming range read.
func (s3a *S3ApiServer) fetchChunkData(ctx context.Context, fileId string, cipherKey []byte, isCompressed bool, offset int64, size int64, isFullChunk bool) (io.ReadCloser, error) {
if len(cipherKey) > 0 || isCompressed {
lookupFileIdFn := s3a.createLookupFileIdFunction()
urlStrings, err := lookupFileIdFn(ctx, fileId)
if err != nil || len(urlStrings) == 0 {
return nil, fmt.Errorf("failed to lookup chunk %s: %w", fileId, err)
}
buffer := make([]byte, size)
n, err := util_http.RetriedFetchChunkData(ctx, buffer, urlStrings, cipherKey, isCompressed, isFullChunk, offset, fileId, nil)
if err != nil {
return nil, fmt.Errorf("failed to fetch chunk %s: %w", fileId, err)
}
return io.NopCloser(bytes.NewReader(buffer[:n])), nil
}
return s3a.readChunkRange(ctx, fileId, offset, size, isFullChunk)
}
func (s3a *S3ApiServer) readChunkRange(ctx context.Context, fileId string, offset int64, size int64, isFullChunk bool) (io.ReadCloser, error) {
// Lookup the volume server URLs for this chunk
lookupFileIdFn := s3a.createLookupFileIdFunction()
urlStrings, err := lookupFileIdFn(ctx, fileId)
@@ -2169,63 +2198,21 @@ func (s3a *S3ApiServer) fetchFullChunk(ctx context.Context, fileId string) (io.R
return nil, fmt.Errorf("failed to lookup chunk %s: %w", fileId, err)
}
// Use the first URL
// Use the first URL (already contains complete URL with fileId)
chunkUrl := urlStrings[0]
// Generate JWT for volume server authentication (uses config loaded once at startup)
jwt := filer.JwtForVolumeServer(fileId)
// Create request WITHOUT Range header to get full chunk
req, err := http.NewRequestWithContext(ctx, "GET", chunkUrl, nil)
if err != nil {
return nil, fmt.Errorf("failed to create request: %w", err)
}
// Set JWT for authentication
if jwt != "" {
req.Header.Set("Authorization", security.BearerPrefix+jwt)
}
// Use shared HTTP client
resp, err := volumeServerHTTPClient.Do(req)
if err != nil {
return nil, fmt.Errorf("failed to fetch chunk: %w", err)
}
if resp.StatusCode != http.StatusOK {
resp.Body.Close()
return nil, fmt.Errorf("unexpected status code %d for chunk %s", resp.StatusCode, fileId)
}
return resp.Body, nil
}
// fetchChunkViewData fetches encrypted data for a chunk view (with range)
func (s3a *S3ApiServer) fetchChunkViewData(ctx context.Context, chunkView *filer.ChunkView) (io.ReadCloser, error) {
// Lookup the volume server URLs for this chunk
lookupFileIdFn := s3a.createLookupFileIdFunction()
urlStrings, err := lookupFileIdFn(ctx, chunkView.FileId)
if err != nil || len(urlStrings) == 0 {
return nil, fmt.Errorf("failed to lookup chunk %s: %w", chunkView.FileId, err)
}
// Use the first URL (already contains complete URL with fileId)
chunkUrl := urlStrings[0]
// Generate JWT for volume server authentication (uses config loaded once at startup)
jwt := filer.JwtForVolumeServer(chunkView.FileId)
// Create request with Range header for the chunk view
// chunkUrl already contains the complete URL including fileId
req, err := http.NewRequestWithContext(ctx, "GET", chunkUrl, nil)
if err != nil {
return nil, fmt.Errorf("failed to create request: %w", err)
}
// Set Range header to fetch only the needed portion of the chunk
if !chunkView.IsFullChunk() {
rangeEnd := chunkView.OffsetInChunk + int64(chunkView.ViewSize) - 1
req.Header.Set("Range", fmt.Sprintf("bytes=%d-%d", chunkView.OffsetInChunk, rangeEnd))
if !isFullChunk {
req.Header.Set("Range", fmt.Sprintf("bytes=%d-%d", offset, offset+size-1))
}
// Set JWT for authentication
@@ -2241,7 +2228,7 @@ func (s3a *S3ApiServer) fetchChunkViewData(ctx context.Context, chunkView *filer
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusPartialContent {
resp.Body.Close()
return nil, fmt.Errorf("unexpected status code %d for chunk %s", resp.StatusCode, chunkView.FileId)
return nil, fmt.Errorf("unexpected status code %d for chunk %s", resp.StatusCode, fileId)
}
return resp.Body, nil
@@ -3262,36 +3249,7 @@ func (l *lazyMultipartChunkReader) Close() error {
// createEncryptedChunkReader creates a reader for a single encrypted chunk
// Context propagation ensures cancellation if the S3 client disconnects
func (s3a *S3ApiServer) createEncryptedChunkReader(ctx context.Context, chunk *filer_pb.FileChunk) (io.ReadCloser, error) {
// Get chunk URL
srcUrl, err := s3a.lookupVolumeUrl(chunk.GetFileIdString())
if err != nil {
return nil, fmt.Errorf("lookup volume URL for chunk %s: %v", chunk.GetFileIdString(), err)
}
// Create HTTP request with context for cancellation propagation
req, err := http.NewRequestWithContext(ctx, "GET", srcUrl, nil)
if err != nil {
return nil, fmt.Errorf("create HTTP request for chunk: %v", err)
}
// Attach volume server JWT for authentication (uses config loaded once at startup)
jwt := filer.JwtForVolumeServer(chunk.GetFileIdString())
if jwt != "" {
req.Header.Set("Authorization", security.BearerPrefix+jwt)
}
// Use shared HTTP client with connection pooling
resp, err := volumeServerHTTPClient.Do(req)
if err != nil {
return nil, fmt.Errorf("execute HTTP request for chunk: %v", err)
}
if resp.StatusCode != http.StatusOK {
resp.Body.Close()
return nil, fmt.Errorf("HTTP request for chunk failed: %d", resp.StatusCode)
}
return resp.Body, nil
return s3a.fetchChunkData(ctx, chunk.GetFileIdString(), chunk.CipherKey, chunk.IsCompressed, 0, int64(chunk.Size), true)
}
// MultipartSSEReader wraps multiple readers and ensures all underlying readers are properly closed
+1 -1
View File
@@ -564,7 +564,7 @@ func (s3a *S3ApiServer) putToFiler(r *http.Request, filePath string, dataReader
// Apply bucket default encryption if no explicit encryption was provided
// This implements AWS S3 behavior where bucket default encryption automatically applies
if !hasExplicitEncryption(customerKey, sseKMSKey, sseS3Key) && !s3a.cipher {
if !hasExplicitEncryption(customerKey, sseKMSKey, sseS3Key) {
glog.V(4).Infof("putToFiler: no explicit encryption detected, checking for bucket default encryption")
// Apply bucket default encryption and get the result
@@ -1,6 +1,7 @@
package s3api
import (
"bytes"
"context"
"fmt"
"io"
@@ -17,10 +18,12 @@ import (
"google.golang.org/grpc/status"
"google.golang.org/protobuf/proto"
"github.com/seaweedfs/seaweedfs/weed/filer"
"github.com/seaweedfs/seaweedfs/weed/pb"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
"github.com/seaweedfs/seaweedfs/weed/util"
"github.com/seaweedfs/seaweedfs/weed/wdclient"
)
@@ -34,6 +37,7 @@ type fakeVolumeServer struct {
mu sync.Mutex
deletedFids []string
stored map[string][]byte
}
func (f *fakeVolumeServer) BatchDelete(_ context.Context, req *volume_server_pb.BatchDeleteRequest) (*volume_server_pb.BatchDeleteResponse, error) {
@@ -55,8 +59,30 @@ func (f *fakeVolumeServer) deleted() []string {
func startFakeVolumeServer(t *testing.T) *fakeVolumeServer {
t.Helper()
v := &fakeVolumeServer{}
v := &fakeVolumeServer{stored: map[string][]byte{}}
upload := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
fid := strings.TrimPrefix(r.URL.Path, "/")
if r.Method == http.MethodGet {
v.mu.Lock()
data, ok := v.stored[fid]
v.mu.Unlock()
if !ok {
w.WriteHeader(http.StatusNotFound)
return
}
if rg := r.Header.Get("Range"); rg != "" {
var start, end int
fmt.Sscanf(rg, "bytes=%d-%d", &start, &end)
if end >= len(data) {
end = len(data) - 1
}
w.WriteHeader(http.StatusPartialContent)
w.Write(data[start : end+1])
return
}
w.Write(data)
return
}
io.Copy(io.Discard, r.Body)
w.Header().Set("Content-MD5", r.Header.Get("Content-MD5"))
w.WriteHeader(http.StatusCreated)
@@ -322,3 +348,62 @@ func TestPutToFilerUnverifiableCreateKeepsChunks(t *testing.T) {
t.Fatalf("chunks were deleted while the create outcome was unverifiable: %v", deleted)
}
}
// A volume-encrypted chunk must be decrypted before SSE decryption sees it:
// fetchChunkData feeds ciphered chunks through the cipher-aware read path and
// slices plaintext space, while plain chunks keep streaming range reads.
func TestFetchChunkDataDecryptsVolumeCipher(t *testing.T) {
volume := startFakeVolumeServer(t)
filerImpl := &ambiguousPutFiler{volume: volume, entries: map[string]*filer_pb.Entry{}}
s3a := newPutTestServer(t, startFakeFiler(t, filerImpl))
plaintext := []byte("0123456789abcdefghijklmnopqrstuvwxyz")
cipherKey := util.GenCipherKey()
ciphertext, err := util.Encrypt(plaintext, cipherKey)
if err != nil {
t.Fatal(err)
}
fid := "3,01637037d6"
volume.mu.Lock()
volume.stored[fid] = ciphertext
volume.mu.Unlock()
full, err := s3a.fetchFullChunk(context.Background(), &filer.ChunkView{
FileId: fid, ChunkSize: uint64(len(plaintext)), CipherKey: cipherKey,
})
if err != nil {
t.Fatal(err)
}
got, _ := io.ReadAll(full)
full.Close()
if !bytes.Equal(got, plaintext) {
t.Fatalf("full chunk read = %q, want %q", got, plaintext)
}
view, err := s3a.fetchChunkViewData(context.Background(), &filer.ChunkView{
FileId: fid, OffsetInChunk: 5, ViewSize: 4, ChunkSize: uint64(len(plaintext)), CipherKey: cipherKey,
})
if err != nil {
t.Fatal(err)
}
got, _ = io.ReadAll(view)
view.Close()
if !bytes.Equal(got, plaintext[5:9]) {
t.Fatalf("ranged ciphered read = %q, want %q", got, plaintext[5:9])
}
volume.mu.Lock()
volume.stored[fid] = plaintext
volume.mu.Unlock()
plain, err := s3a.fetchChunkViewData(context.Background(), &filer.ChunkView{
FileId: fid, OffsetInChunk: 5, ViewSize: 4, ChunkSize: uint64(len(plaintext)),
})
if err != nil {
t.Fatal(err)
}
got, _ = io.ReadAll(plain)
plain.Close()
if !bytes.Equal(got, plaintext[5:9]) {
t.Fatalf("ranged plain read = %q, want %q", got, plaintext[5:9])
}
}