diff --git a/.github/workflows/s3tests.yml b/.github/workflows/s3tests.yml index 35ab1c177..7d481a3f6 100644 --- a/.github/workflows/s3tests.yml +++ b/.github/workflows/s3tests.yml @@ -289,6 +289,10 @@ jobs: s3tests/functional/test_s3.py::test_object_write_check_etag \ s3tests/functional/test_s3.py::test_object_write_cache_control \ s3tests/functional/test_s3.py::test_object_write_expires \ + s3tests/functional/test_s3.py::test_put_obj_enc_conflict_c_s3 \ + s3tests/functional/test_s3.py::test_put_obj_enc_conflict_c_kms \ + s3tests/functional/test_s3.py::test_put_obj_enc_conflict_s3_kms \ + s3tests/functional/test_s3.py::test_put_obj_enc_conflict_bad_enc_kms \ s3tests/functional/test_s3.py::test_object_content_encoding_aws_chunked \ s3tests/functional/test_s3.py::test_object_write_read_update_read_delete \ s3tests/functional/test_s3.py::test_object_metadata_replaced_on_put \ @@ -1151,6 +1155,10 @@ jobs: s3tests/functional/test_s3.py::test_object_write_check_etag \ s3tests/functional/test_s3.py::test_object_write_cache_control \ s3tests/functional/test_s3.py::test_object_write_expires \ + s3tests/functional/test_s3.py::test_put_obj_enc_conflict_c_s3 \ + s3tests/functional/test_s3.py::test_put_obj_enc_conflict_c_kms \ + s3tests/functional/test_s3.py::test_put_obj_enc_conflict_s3_kms \ + s3tests/functional/test_s3.py::test_put_obj_enc_conflict_bad_enc_kms \ s3tests/functional/test_s3.py::test_object_content_encoding_aws_chunked \ s3tests/functional/test_s3.py::test_object_write_read_update_read_delete \ s3tests/functional/test_s3.py::test_object_metadata_replaced_on_put \ diff --git a/weed/s3api/s3api_copy_validation.go b/weed/s3api/s3api_copy_validation.go index 9298787e0..eab3cf8d3 100644 --- a/weed/s3api/s3api_copy_validation.go +++ b/weed/s3api/s3api_copy_validation.go @@ -1,6 +1,7 @@ package s3api import ( + "errors" "fmt" "net/http" @@ -115,6 +116,26 @@ func validateSSEKMSCopyRequirements(srcMetadata map[string][]byte, headers http. // validateEncryptionCompatibility validates that encryption methods are not conflicting func validateEncryptionCompatibility(headers http.Header) error { + // A repeated header is rejected rather than deduped: the encryption paths + // apply only the first value, so extra values could hide the method the + // client actually asked for. + for _, name := range []string{ + s3_constants.AmzServerSideEncryption, + s3_constants.AmzServerSideEncryptionCustomerAlgorithm, + s3_constants.AmzServerSideEncryptionCustomerKey, + s3_constants.AmzServerSideEncryptionCustomerKeyMD5, + s3_constants.AmzServerSideEncryptionAwsKmsKeyId, + s3_constants.AmzServerSideEncryptionContext, + s3_constants.AmzServerSideEncryptionBucketKeyEnabled, + } { + if len(headers.Values(name)) > 1 { + return &CopyValidationError{ + Code: s3err.ErrInvalidRequest, + Message: fmt.Sprintf("Multiple %s headers specified - only one is allowed", name), + } + } + } + sseAlgorithm := headers.Get(s3_constants.AmzServerSideEncryption) hasSSEC := hasSSECHeaders(headers) hasSSEKMS := sseAlgorithm == s3_constants.SSEAlgorithmKMS @@ -129,19 +150,25 @@ func validateEncryptionCompatibility(headers http.Header) error { } } - // Count how many encryption methods are specified - encryptionCount := 0 - if hasSSEC { - encryptionCount++ - } - if hasSSEKMS { - encryptionCount++ - } - if hasSSES3 { - encryptionCount++ + // KMS options only apply to aws:kms; with another method or none they name + // no method at all. + if !hasSSEKMS && hasHeaderValue(headers, + s3_constants.AmzServerSideEncryptionAwsKmsKeyId, + s3_constants.AmzServerSideEncryptionContext, + s3_constants.AmzServerSideEncryptionBucketKeyEnabled) { + return &CopyValidationError{ + Code: s3err.ErrInvalidRequest, + Message: "KMS encryption options require the aws:kms encryption method", + } } // Only one encryption method should be specified + encryptionCount := 0 + for _, specified := range []bool{hasSSEC, hasSSEKMS, hasSSES3} { + if specified { + encryptionCount++ + } + } if encryptionCount > 1 { return &CopyValidationError{ Code: s3err.ErrInvalidRequest, @@ -152,6 +179,21 @@ func validateEncryptionCompatibility(headers http.Header) error { return nil } +// ValidateRequestEncryption rejects PutObject or CreateMultipartUpload +// encryption headers that no single method can honor, with the InvalidArgument +// codes S3 returns for them. +func ValidateRequestEncryption(headers http.Header) s3err.ErrorCode { + err := validateEncryptionCompatibility(headers) + if err == nil { + return s3err.ErrNone + } + var validationErr *CopyValidationError + if errors.As(err, &validationErr) && validationErr.Code == s3err.ErrInvalidEncryptionAlgorithm { + return s3err.ErrInvalidEncryptionMethod + } + return s3err.ErrIncompatibleEncryptionMethod +} + // validateSSECCopyHeaderCompleteness validates that all required SSE-C copy headers are present func validateSSECCopyHeaderCompleteness(headers http.Header) error { algorithm := headers.Get(s3_constants.AmzCopySourceServerSideEncryptionCustomerAlgorithm) @@ -229,16 +271,29 @@ func validateSSECHeaderCompleteness(headers http.Header) error { } // Helper functions for header detection +func hasHeaderValue(headers http.Header, names ...string) bool { + for _, name := range names { + for _, value := range headers.Values(name) { + if value != "" { + return true + } + } + } + return false +} + func hasSSECCopyHeaders(headers http.Header) bool { - return headers.Get(s3_constants.AmzCopySourceServerSideEncryptionCustomerAlgorithm) != "" || - headers.Get(s3_constants.AmzCopySourceServerSideEncryptionCustomerKey) != "" || - headers.Get(s3_constants.AmzCopySourceServerSideEncryptionCustomerKeyMD5) != "" + return hasHeaderValue(headers, + s3_constants.AmzCopySourceServerSideEncryptionCustomerAlgorithm, + s3_constants.AmzCopySourceServerSideEncryptionCustomerKey, + s3_constants.AmzCopySourceServerSideEncryptionCustomerKeyMD5) } func hasSSECHeaders(headers http.Header) bool { - return headers.Get(s3_constants.AmzServerSideEncryptionCustomerAlgorithm) != "" || - headers.Get(s3_constants.AmzServerSideEncryptionCustomerKey) != "" || - headers.Get(s3_constants.AmzServerSideEncryptionCustomerKeyMD5) != "" + return hasHeaderValue(headers, + s3_constants.AmzServerSideEncryptionCustomerAlgorithm, + s3_constants.AmzServerSideEncryptionCustomerKey, + s3_constants.AmzServerSideEncryptionCustomerKeyMD5) } // validateEncryptionContext validates the encryption context header format diff --git a/weed/s3api/s3api_copy_validation_test.go b/weed/s3api/s3api_copy_validation_test.go new file mode 100644 index 000000000..00edd6010 --- /dev/null +++ b/weed/s3api/s3api_copy_validation_test.go @@ -0,0 +1,107 @@ +package s3api + +import ( + "bytes" + "io" + "net/http" + "net/http/httptest" + "testing" + + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants" + "github.com/seaweedfs/seaweedfs/weed/s3api/s3err" + "github.com/stretchr/testify/assert" +) + +// TestValidateRequestEncryption verifies that a PutObject or +// CreateMultipartUpload request asking for an encryption method that cannot be +// honored is rejected instead of being stored unencrypted +func TestValidateRequestEncryption(t *testing.T) { + ssec := map[string]string{ + s3_constants.AmzServerSideEncryptionCustomerAlgorithm: "AES256", + s3_constants.AmzServerSideEncryptionCustomerKey: "a2tra2tra2tra2tra2tra2tra2tra2tra2tra2tra2s=", + s3_constants.AmzServerSideEncryptionCustomerKeyMD5: "f6OQvGsmFBq4WOqaVcuO5w==", + } + testCases := []struct { + name string + headers map[string]string + sse string + sseValues []string + dup string + want s3err.ErrorCode + }{ + {name: "no encryption", want: s3err.ErrNone}, + {name: "SSE-S3", sse: s3_constants.SSEAlgorithmAES256, want: s3err.ErrNone}, + {name: "SSE-KMS", sse: s3_constants.SSEAlgorithmKMS, want: s3err.ErrNone}, + {name: "SSE-C", headers: ssec, want: s3err.ErrNone}, + {name: "unknown algorithm", sse: "aes:kms", want: s3err.ErrInvalidEncryptionMethod}, + {name: "misspelled AES256", sse: "AES-256", want: s3err.ErrInvalidEncryptionMethod}, + {name: "SSE-C and SSE-S3", headers: ssec, sse: s3_constants.SSEAlgorithmAES256, want: s3err.ErrIncompatibleEncryptionMethod}, + {name: "SSE-C and SSE-KMS", headers: ssec, sse: s3_constants.SSEAlgorithmKMS, want: s3err.ErrIncompatibleEncryptionMethod}, + {name: "repeated algorithm header", sseValues: []string{"AES256", "aws:kms"}, want: s3err.ErrIncompatibleEncryptionMethod}, + {name: "repeated identical algorithm", sseValues: []string{"AES256", "AES256"}, want: s3err.ErrIncompatibleEncryptionMethod}, + {name: "empty value before KMS", sseValues: []string{"", "aws:kms"}, want: s3err.ErrIncompatibleEncryptionMethod}, + {name: "SSE-S3 hidden behind empty value", sseValues: []string{"", "AES256"}, headers: ssec, want: s3err.ErrIncompatibleEncryptionMethod}, + {name: "repeated SSE-C key header", headers: ssec, dup: s3_constants.AmzServerSideEncryptionCustomerKey, want: s3err.ErrIncompatibleEncryptionMethod}, + {name: "repeated KMS key id", sse: s3_constants.SSEAlgorithmKMS, dup: s3_constants.AmzServerSideEncryptionAwsKmsKeyId, want: s3err.ErrIncompatibleEncryptionMethod}, + {name: "SSE-S3 with KMS key id", sse: s3_constants.SSEAlgorithmAES256, headers: map[string]string{s3_constants.AmzServerSideEncryptionAwsKmsKeyId: "key-id"}, want: s3err.ErrIncompatibleEncryptionMethod}, + {name: "KMS key id without method", headers: map[string]string{s3_constants.AmzServerSideEncryptionAwsKmsKeyId: "key-id"}, want: s3err.ErrIncompatibleEncryptionMethod}, + {name: "KMS key id with aws:kms", sse: s3_constants.SSEAlgorithmKMS, headers: map[string]string{s3_constants.AmzServerSideEncryptionAwsKmsKeyId: "key-id"}, want: s3err.ErrNone}, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + h := http.Header{} + for k, v := range tc.headers { + h.Set(k, v) + } + if tc.sse != "" { + h.Set(s3_constants.AmzServerSideEncryption, tc.sse) + } + for _, v := range tc.sseValues { + h.Add(s3_constants.AmzServerSideEncryption, v) + } + if tc.dup != "" { + h.Add(tc.dup, "first") + h.Add(tc.dup, "second") + } + assert.Equal(t, tc.want, ValidateRequestEncryption(h)) + }) + } +} + +// TestDirectoryMarkerSSEC verifies marker content stored under SSE-C encrypts +// on write and decrypts on read through the stored entry metadata +func TestDirectoryMarkerSSEC(t *testing.T) { + ssecHeaders := func(r *http.Request) { + r.Header.Set(s3_constants.AmzServerSideEncryptionCustomerAlgorithm, "AES256") + r.Header.Set(s3_constants.AmzServerSideEncryptionCustomerKey, "a2tra2tra2tra2tra2tra2tra2tra2tra2tra2tra2s=") + r.Header.Set(s3_constants.AmzServerSideEncryptionCustomerKeyMD5, "mT2HRsMGJ5IX5C+0rreZ8Q==") + } + plaintext := []byte("directory marker content") + s3a := &S3ApiServer{} + + putReq := httptest.NewRequest(http.MethodPut, "/bucket/dir/", nil) + ssecHeaders(putReq) + sseResult, errCode := s3a.handleAllSSEEncryption(putReq, bytes.NewReader(plaintext), 0) + assert.Equal(t, s3err.ErrNone, errCode) + encrypted, err := io.ReadAll(sseResult.DataReader) + assert.NoError(t, err) + assert.NotEqual(t, plaintext, encrypted) + + entry := &filer_pb.Entry{Extended: map[string][]byte{}, Content: encrypted} + storeSSEMetadata(entry, sseResult) + + sseType := s3a.detectPrimarySSEType(entry) + assert.Equal(t, s3_constants.SSETypeC, sseType) + + getReq := httptest.NewRequest(http.MethodGet, "/bucket/dir/", nil) + ssecHeaders(getReq) + decrypted, errCode := s3a.decryptDirectoryContent(getReq, entry, sseType) + assert.Equal(t, s3err.ErrNone, errCode) + assert.Equal(t, plaintext, decrypted) + + bareReq := httptest.NewRequest(http.MethodGet, "/bucket/dir/", nil) + _, errCode = s3a.decryptDirectoryContent(bareReq, entry, sseType) + assert.Equal(t, s3err.ErrSSECustomerKeyMissing, errCode) +} diff --git a/weed/s3api/s3api_object_handlers.go b/weed/s3api/s3api_object_handlers.go index 7b45d186c..a04e846f9 100644 --- a/weed/s3api/s3api_object_handlers.go +++ b/weed/s3api/s3api_object_handlers.go @@ -465,7 +465,7 @@ func (s3a *S3ApiServer) resolveObjectEntry(bucket, object, versionId string) (*f } // serveDirectoryContent serves the content of a directory object directly -func (s3a *S3ApiServer) serveDirectoryContent(w http.ResponseWriter, r *http.Request, entry *filer_pb.Entry) { +func (s3a *S3ApiServer) serveDirectoryContent(w http.ResponseWriter, r *http.Request, entry *filer_pb.Entry, bucket, object string) { // Defensive nil checks - entry and attributes should never be nil, but guard against it if entry == nil || entry.Attributes == nil { glog.Errorf("serveDirectoryContent: entry or attributes is nil") @@ -473,10 +473,52 @@ func (s3a *S3ApiServer) serveDirectoryContent(w http.ResponseWriter, r *http.Req return } + // A directory promoted over an uploaded object still holds that object's + // chunks while entry.Content stays empty; GET streams them like a regular + // object so it delivers the bytes HEAD reports for this entry. + if r.Method != http.MethodHead && len(entry.Content) == 0 && len(entry.Chunks) > 0 { + sseType := s3a.detectPrimarySSEType(entry) + if err := s3a.streamFromVolumeServersWithSSE(w, r, entry, sseType, bucket, object, ""); err != nil { + var streamErr *StreamError + if errors.As(err, &streamErr) && streamErr.ResponseWritten { + return + } + if isCanceledStreamingError(err) { + glog.V(3).Infof("serveDirectoryContent: client disconnected while streaming %s/%s: %v", bucket, object, err) + return + } + glog.Errorf("serveDirectoryContent: failed to stream %s/%s: %v", bucket, object, err) + if errors.Is(err, util_http.ErrTooManyRequests) { + s3err.WriteErrorResponse(w, r, s3err.ErrRequestBytesExceed) + } else if shouldWriteStreamingErrorResponse(err) { + s3err.WriteErrorResponse(w, r, s3err.ErrInternalError) + } + } + return + } + // Set content type - use stored MIME type or default. A directory without a stored // mime and without data of its own answers application/x-directory, the marker type // Hadoop-style clients (e.g. flink-s3-fs-presto) require to classify the path as a // directory; defaulting to octet-stream makes them treat it as a 0-byte file. + // Marker content may be stored encrypted; decrypt before serving. HEAD + // returns no body, so it only checks access and skips decryption — and + // any KMS call — entirely. + content := entry.Content + sseType := s3a.detectPrimarySSEType(entry) + if sseType != "" && sseType != "None" { + var errCode s3err.ErrorCode + if r.Method == http.MethodHead { + errCode = s3a.checkDirectorySSEAccess(r, entry, sseType) + } else { + content, errCode = s3a.decryptDirectoryContent(r, entry, sseType) + } + if errCode != s3err.ErrNone { + s3err.WriteErrorResponse(w, r, errCode) + return + } + } + contentType := entry.Attributes.Mime if contentType == "" { if entry.IsDirectoryKeyObject() { @@ -487,8 +529,10 @@ func (s3a *S3ApiServer) serveDirectoryContent(w http.ResponseWriter, r *http.Req } w.Header().Set("Content-Type", contentType) - // Set content length - use FileSize for accuracy, especially for large files - contentLength := int64(entry.Attributes.FileSize) + contentLength := int64(len(content)) + if r.Method == http.MethodHead && len(entry.Chunks) > 0 { + contentLength = int64(entry.Attributes.FileSize) + } w.Header().Set("Content-Length", strconv.FormatInt(contentLength, 10)) // Set last modified @@ -497,6 +541,8 @@ func (s3a *S3ApiServer) serveDirectoryContent(w http.ResponseWriter, r *http.Req // Set ETag w.Header().Set("ETag", "\""+filer.ETag(entry)+"\"") + s3a.addSSEResponseHeadersFromEntry(w, r, entry, sseType) + // For HEAD requests, don't write body if r.Method == http.MethodHead { w.WriteHeader(http.StatusOK) @@ -505,13 +551,87 @@ func (s3a *S3ApiServer) serveDirectoryContent(w http.ResponseWriter, r *http.Req // Write content w.WriteHeader(http.StatusOK) - if len(entry.Content) > 0 { - if _, err := w.Write(entry.Content); err != nil { + if len(content) > 0 { + if _, err := w.Write(content); err != nil { glog.Errorf("serveDirectoryContent: failed to write response: %v", err) } } } +// decryptDirectoryContent returns the plaintext of a directory marker's inline +// entry.Content, which PutObjectHandler stores through the shared SSE path. +func (s3a *S3ApiServer) decryptDirectoryContent(r *http.Request, entry *filer_pb.Entry, sseType string) ([]byte, s3err.ErrorCode) { + var reader io.Reader + var err error + switch sseType { + case s3_constants.SSETypeC: + customerKey, parseErr := ParseSSECHeaders(r) + if parseErr != nil { + return nil, MapSSECErrorToS3Error(parseErr) + } + if customerKey == nil { + return nil, s3err.ErrSSECustomerKeyMissing + } + if storedKeyMD5 := string(entry.Extended[s3_constants.AmzServerSideEncryptionCustomerKeyMD5]); storedKeyMD5 != "" && customerKey.KeyMD5 != storedKeyMD5 { + return nil, s3err.ErrAccessDenied + } + iv, ivErr := GetSSECIVFromMetadata(entry.Extended) + if ivErr != nil { + return nil, s3err.ErrInternalError + } + reader, err = CreateSSECDecryptedReader(bytes.NewReader(entry.Content), customerKey, iv) + case s3_constants.SSETypeKMS: + sseKMSKey, deserErr := DeserializeSSEKMSMetadata(entry.Extended[s3_constants.SeaweedFSSSEKMSKey]) + if deserErr != nil { + return nil, s3err.ErrInternalError + } + reader, err = CreateSSEKMSDecryptedReader(bytes.NewReader(entry.Content), sseKMSKey) + case s3_constants.SSETypeS3: + keyManager := GetSSES3KeyManager() + sseS3Key, deserErr := DeserializeSSES3Metadata(entry.Extended[s3_constants.SeaweedFSSSES3Key], keyManager) + if deserErr != nil { + return nil, s3err.ErrInternalError + } + iv, ivErr := GetSSES3IV(entry, sseS3Key, keyManager) + if ivErr != nil { + return nil, s3err.ErrInternalError + } + reader, err = CreateSSES3DecryptedReader(bytes.NewReader(entry.Content), sseS3Key, iv) + default: + return entry.Content, s3err.ErrNone + } + if err != nil { + glog.Errorf("decryptDirectoryContent: %v", err) + return nil, s3err.ErrInternalError + } + content, readErr := io.ReadAll(reader) + if readErr != nil { + glog.Errorf("decryptDirectoryContent: %v", readErr) + return nil, s3err.ErrInternalError + } + return content, s3err.ErrNone +} + +// checkDirectorySSEAccess runs the request-level checks a HEAD of an encrypted +// marker needs: SSE-C requires the customer key to parse and to match the key +// the marker was written with; SSE-KMS and SSE-S3 need nothing from the caller. +func (s3a *S3ApiServer) checkDirectorySSEAccess(r *http.Request, entry *filer_pb.Entry, sseType string) s3err.ErrorCode { + if sseType != s3_constants.SSETypeC { + return s3err.ErrNone + } + customerKey, parseErr := ParseSSECHeaders(r) + if parseErr != nil { + return MapSSECErrorToS3Error(parseErr) + } + if customerKey == nil { + return s3err.ErrSSECustomerKeyMissing + } + if storedKeyMD5 := string(entry.Extended[s3_constants.AmzServerSideEncryptionCustomerKeyMD5]); storedKeyMD5 != "" && customerKey.KeyMD5 != storedKeyMD5 { + return s3err.ErrAccessDenied + } + return s3err.ErrNone +} + // handleDirectoryObjectRequest is a helper function that handles directory object requests // for both GET and HEAD operations, eliminating code duplication func (s3a *S3ApiServer) handleDirectoryObjectRequest(w http.ResponseWriter, r *http.Request, bucket, object, handlerName string) bool { @@ -538,7 +658,7 @@ func (s3a *S3ApiServer) handleDirectoryObjectRequest(w http.ResponseWriter, r *h return true // Request was handled (denied) } glog.V(2).Infof("%s: directory object %s/%s found, serving content", handlerName, bucket, object) - s3a.serveDirectoryContent(w, r, dirEntry) + s3a.serveDirectoryContent(w, r, dirEntry, bucket, object) return true // Request was handled successfully } else if isDirectoryObject { // Directory object but doesn't exist diff --git a/weed/s3api/s3api_object_handlers_dir_test.go b/weed/s3api/s3api_object_handlers_dir_test.go index 163d084da..4dd15f33b 100644 --- a/weed/s3api/s3api_object_handlers_dir_test.go +++ b/weed/s3api/s3api_object_handlers_dir_test.go @@ -38,7 +38,7 @@ func TestServeDirectoryContentContentType(t *testing.T) { t.Run(tt.name, func(t *testing.T) { req := httptest.NewRequest(http.MethodHead, "/bucket/dir/", nil) rec := httptest.NewRecorder() - s3a.serveDirectoryContent(rec, req, tt.entry) + s3a.serveDirectoryContent(rec, req, tt.entry, "bucket", "dir/") if rec.Code != http.StatusOK { t.Fatalf("status = %d, want 200", rec.Code) } diff --git a/weed/s3api/s3api_object_handlers_multipart.go b/weed/s3api/s3api_object_handlers_multipart.go index 4ed8c2f5c..2f68b3889 100644 --- a/weed/s3api/s3api_object_handlers_multipart.go +++ b/weed/s3api/s3api_object_handlers_multipart.go @@ -40,6 +40,12 @@ func (s3a *S3ApiServer) NewMultipartUploadHandler(w http.ResponseWriter, r *http return } + // Reject before the auto-create check so a refused upload cannot leave a bucket behind + if errCode := ValidateRequestEncryption(r.Header); errCode != s3err.ErrNone { + s3err.WriteErrorResponse(w, r, errCode) + return + } + // Check if bucket exists, and create it if it doesn't (auto-create bucket) if err := s3a.checkBucket(r, bucket); err == s3err.ErrNoSuchBucket { // Auto-create bucket if it doesn't exist (requires Admin permission) diff --git a/weed/s3api/s3api_object_handlers_put.go b/weed/s3api/s3api_object_handlers_put.go index c2ca1a8e1..bf2c2ec22 100644 --- a/weed/s3api/s3api_object_handlers_put.go +++ b/weed/s3api/s3api_object_handlers_put.go @@ -127,6 +127,11 @@ func (s3a *S3ApiServer) PutObjectHandler(w http.ResponseWriter, r *http.Request) return } + if errCode := ValidateRequestEncryption(r.Header); errCode != s3err.ErrNone { + s3err.WriteErrorResponse(w, r, errCode) + return + } + if r.Header.Get("Cache-Control") != "" { if _, err = cacheobject.ParseRequestCacheControl(r.Header.Get("Cache-Control")); err != nil { s3err.WriteErrorResponse(w, r, s3err.ErrInvalidDigest) @@ -215,6 +220,18 @@ func (s3a *S3ApiServer) PutObjectHandler(w http.ResponseWriter, r *http.Request) dirMd5 := md5.Sum(dirContent) dirEtag := fmt.Sprintf("%x", dirMd5) + // SSE applies to marker content too — it is stored in entry.Content + sseResult, sseErrCode := s3a.handleAllSSEEncryption(r, bytes.NewReader(dirContent), 0) + if sseErrCode != s3err.ErrNone { + s3err.WriteErrorResponse(w, r, sseErrCode) + return + } + if dirContent, err = io.ReadAll(sseResult.DataReader); err != nil { + glog.Errorf("PutObjectHandler: failed to encrypt directory marker content %s/%s: %v", bucket, object, err) + s3err.WriteErrorResponse(w, r, s3err.ErrInternalError) + return + } + glog.Infof("PutObjectHandler: explicit directory marker %s/%s (contentType=%q, len=%d)", bucket, object, objectContentType, r.ContentLength) // mkdir replaces this entry outright, so a lock recorded on the entry itself @@ -253,6 +270,8 @@ func (s3a *S3ApiServer) PutObjectHandler(w http.ResponseWriter, r *http.Request) } entry.Extended[s3_constants.ExtETagKey] = []byte(dirEtag) + storeSSEMetadata(entry, sseResult) + // Set object owner for directory objects (same as regular objects) s3a.setObjectOwnerFromRequest(r, bucket, entry) applyPutObjectACL(r, entry) @@ -269,6 +288,12 @@ func (s3a *S3ApiServer) PutObjectHandler(w http.ResponseWriter, r *http.Request) s3err.WriteErrorResponse(w, r, markerCode) return } + sseRespMetadata := SSEResponseMetadata{SSEType: sseResult.SSEType} + if sseResult.SSEKMSKey != nil { + sseRespMetadata.KMSKeyID = sseResult.SSEKMSKey.KeyID + sseRespMetadata.BucketKeyEnabled = sseResult.SSEKMSKey.BucketKeyEnabled + } + s3a.setSSEResponseHeaders(w, r, sseRespMetadata) setEtag(w, dirEtag) } else { // Get detailed versioning state for the bucket @@ -577,6 +602,12 @@ func (s3a *S3ApiServer) putToFiler(r *http.Request, filePath string, dataReader return "", s3err.ErrInternalError, SSEResponseMetadata{} } } + + sseResult.SSEKMSKey = sseKMSKey + sseResult.SSEKMSMetadata = sseKMSMetadata + sseResult.SSES3Key = sseS3Key + sseResult.SSES3Metadata = sseS3Metadata + sseResult.SSEType = sseType } else { glog.V(4).Infof("putToFiler: explicit encryption already applied, skipping bucket default encryption") } @@ -912,33 +943,7 @@ func (s3a *S3ApiServer) putToFiler(r *http.Request, filePath string, dataReader glog.V(3).Infof("putToFiler: stored %d tags from X-Amz-Tagging header", len(parsedTags)) } - // Set SSE-C metadata - if customerKey != nil && len(sseIV) > 0 { - // Store IV as RAW bytes (matches filer behavior - filer decodes base64 headers and stores raw bytes) - entry.Extended[s3_constants.SeaweedFSSSEIV] = sseIV - entry.Extended[s3_constants.AmzServerSideEncryptionCustomerAlgorithm] = []byte("AES256") - entry.Extended[s3_constants.AmzServerSideEncryptionCustomerKeyMD5] = []byte(customerKey.KeyMD5) - glog.V(3).Infof("putToFiler: storing SSE-C metadata - IV len=%d", len(sseIV)) - } - - // Set SSE-KMS metadata - if sseKMSKey != nil { - // Store metadata as RAW bytes (matches filer behavior - filer decodes base64 headers and stores raw bytes) - entry.Extended[s3_constants.SeaweedFSSSEKMSKey] = sseKMSMetadata - // Set standard SSE headers for detection - entry.Extended[s3_constants.AmzServerSideEncryption] = []byte("aws:kms") - entry.Extended[s3_constants.AmzServerSideEncryptionAwsKmsKeyId] = []byte(sseKMSKey.KeyID) - glog.V(3).Infof("putToFiler: storing SSE-KMS metadata - keyID=%s, raw len=%d", sseKMSKey.KeyID, len(sseKMSMetadata)) - } - - // Set SSE-S3 metadata - if sseS3Key != nil && len(sseS3Metadata) > 0 { - // Store metadata as RAW bytes (matches filer behavior - filer decodes base64 headers and stores raw bytes) - entry.Extended[s3_constants.SeaweedFSSSES3Key] = sseS3Metadata - // Set standard SSE header for detection - entry.Extended[s3_constants.AmzServerSideEncryption] = []byte("AES256") - glog.V(3).Infof("putToFiler: storing SSE-S3 metadata - keyID=%s, raw len=%d", sseS3Key.KeyID, len(sseS3Metadata)) - } + storeSSEMetadata(entry, sseResult) // Parts (object == "") stay flat: completion rebases chunk offsets, which // manifest chunks cannot express. diff --git a/weed/s3api/s3api_put_handlers.go b/weed/s3api/s3api_put_handlers.go index 825e2f904..0b3a0b5f4 100644 --- a/weed/s3api/s3api_put_handlers.go +++ b/weed/s3api/s3api_put_handlers.go @@ -7,6 +7,7 @@ import ( "strings" "github.com/seaweedfs/seaweedfs/weed/glog" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" "github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants" "github.com/seaweedfs/seaweedfs/weed/s3api/s3err" ) @@ -225,6 +226,28 @@ func (s3a *S3ApiServer) handleSSES3Encryption(r *http.Request, dataReader io.Rea return encryptedReader, sseS3Key, sseS3Metadata, s3err.ErrNone } +// storeSSEMetadata records the encryption result on the entry, in the same +// extended attributes the GET and HEAD handlers read back. +func storeSSEMetadata(entry *filer_pb.Entry, sseResult *PutToFilerEncryptionResult) { + if sseResult == nil { + return + } + if sseResult.CustomerKey != nil && len(sseResult.SSEIV) > 0 { + entry.Extended[s3_constants.SeaweedFSSSEIV] = sseResult.SSEIV + entry.Extended[s3_constants.AmzServerSideEncryptionCustomerAlgorithm] = []byte(s3_constants.SSEAlgorithmAES256) + entry.Extended[s3_constants.AmzServerSideEncryptionCustomerKeyMD5] = []byte(sseResult.CustomerKey.KeyMD5) + } + if sseResult.SSEKMSKey != nil { + entry.Extended[s3_constants.SeaweedFSSSEKMSKey] = sseResult.SSEKMSMetadata + entry.Extended[s3_constants.AmzServerSideEncryption] = []byte(s3_constants.SSEAlgorithmKMS) + entry.Extended[s3_constants.AmzServerSideEncryptionAwsKmsKeyId] = []byte(sseResult.SSEKMSKey.KeyID) + } + if sseResult.SSES3Key != nil && len(sseResult.SSES3Metadata) > 0 { + entry.Extended[s3_constants.SeaweedFSSSES3Key] = sseResult.SSES3Metadata + entry.Extended[s3_constants.AmzServerSideEncryption] = []byte(s3_constants.SSEAlgorithmAES256) + } +} + // handleAllSSEEncryption processes all SSE types in sequence and returns the final encrypted reader // This eliminates repetitive dataReader assignments and centralizes SSE processing func (s3a *S3ApiServer) handleAllSSEEncryption(r *http.Request, dataReader io.Reader, partOffset int64) (*PutToFilerEncryptionResult, s3err.ErrorCode) { diff --git a/weed/s3api/s3err/s3api_errors.go b/weed/s3api/s3err/s3api_errors.go index 4e3d4886f..f0c8bc947 100644 --- a/weed/s3api/s3err/s3api_errors.go +++ b/weed/s3api/s3err/s3api_errors.go @@ -143,6 +143,8 @@ const ( ErrSSECustomerKeyMissing ErrSSECustomerKeyNotNeeded ErrSSEEncryptionTypeMismatch + ErrInvalidEncryptionMethod + ErrIncompatibleEncryptionMethod // SSE-KMS related errors ErrKMSKeyNotFound @@ -689,6 +691,16 @@ var errorCodeResponse = map[ErrorCode]APIError{ Description: "The encryption method specified in the request does not match the encryption method used to encrypt the object.", HTTPStatusCode: http.StatusBadRequest, }, + ErrInvalidEncryptionMethod: { + Code: "InvalidArgument", + Description: "The encryption method specified is not supported", + HTTPStatusCode: http.StatusBadRequest, + }, + ErrIncompatibleEncryptionMethod: { + Code: "InvalidArgument", + Description: "Server Side Encryption with Customer provided key is incompatible with the encryption method specified", + HTTPStatusCode: http.StatusBadRequest, + }, // SSE-KMS error responses ErrKMSKeyNotFound: {