diff --git a/weed/s3api/s3api_bucket_default_encryption_test.go b/weed/s3api/s3api_bucket_default_encryption_test.go new file mode 100644 index 000000000..ece46b4a2 --- /dev/null +++ b/weed/s3api/s3api_bucket_default_encryption_test.go @@ -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) + } + } + }) + } +} diff --git a/weed/s3api/s3api_object_handlers.go b/weed/s3api/s3api_object_handlers.go index 0065eb863..c3da49513 100644 --- a/weed/s3api/s3api_object_handlers.go +++ b/weed/s3api/s3api_object_handlers.go @@ -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 diff --git a/weed/s3api/s3api_object_handlers_put.go b/weed/s3api/s3api_object_handlers_put.go index 6eaa355c4..7783029f5 100644 --- a/weed/s3api/s3api_object_handlers_put.go +++ b/weed/s3api/s3api_object_handlers_put.go @@ -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 diff --git a/weed/s3api/s3api_object_handlers_put_ambiguous_test.go b/weed/s3api/s3api_object_handlers_put_ambiguous_test.go index f6ee1e30b..07b131f44 100644 --- a/weed/s3api/s3api_object_handlers_put_ambiguous_test.go +++ b/weed/s3api/s3api_object_handlers_put_ambiguous_test.go @@ -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]) + } +}