From e7e87402edd910608f77b55c1fb4d64b143d79e5 Mon Sep 17 00:00:00 2001 From: Ayush Sharma <72848455+Aeirx@users.noreply.github.com> Date: Wed, 9 Sep 2026 18:42:07 +0530 Subject: [PATCH 1/6] fix: forward the legal hold as an Object Lock header on federated CopyObject A cross-deployment CopyObject that requests `x-amz-object-lock-legal-hold: ON` answered 200 while the destination carried no hold. The resolved value reached the remote as ordinary user metadata, `X-Amz-Meta-X-Amz-Object-Lock-Legal-Hold`, so nothing applied it. Retention requested on the same copy survived, which is what made the loss easy to miss. The federation branch passes the resolved metadata map straight to `Core.PutObject` as `PutObjectOptions.UserMetadata`. minio-go's `Header()` writes the typed lock fields first, then prefixes every UserMetadata key it does not recognise with `x-amz-meta-`; `supportedHeaders` covers `x-amz-object-lock-mode` and `x-amz-object-lock-retain-until-date` but not `x-amz-object-lock-legal-hold`, and `isAmzHeader` does not match it either. Retention therefore arrives as real headers and the hold does not. The high-level `validate()` that would have rejected the key never runs, because `Core.PutObject` goes straight to the low-level PUT. Carry the hold on the typed `LegalHold` option and forward a cloned map with the raw key removed. The clone matters twice: typed fields are written before the UserMetadata loop, so a leftover raw key would add a bogus `x-amz-meta-` entry beside the correct header, and the proxy's own response and event metadata are rebuilt from the resolved values rather than the forwarding map, which no longer carries the hold. Retention stays in the map deliberately. It already passes through as a standard header, and moving it to the typed `RetainUntilDate` field would format with `time.RFC3339` and truncate a retain-until date to whole seconds. The new test asserts the wire: the remote must receive `X-Amz-Object-Lock-Legal-Hold` and never the `x-amz-meta-` spelling, and the destination version must actually store the hold. It fails without the change with "legal hold forwarded as user metadata [ON]". Fixes #166 Signed-off-by: Ayush Sharma <72848455+Aeirx@users.noreply.github.com> --- cmd/object-copy-federation-legalhold_test.go | 139 +++++++++++++++++++ cmd/object-handlers.go | 32 ++++- 2 files changed, 169 insertions(+), 2 deletions(-) create mode 100644 cmd/object-copy-federation-legalhold_test.go diff --git a/cmd/object-copy-federation-legalhold_test.go b/cmd/object-copy-federation-legalhold_test.go new file mode 100644 index 000000000..44e3e5dc5 --- /dev/null +++ b/cmd/object-copy-federation-legalhold_test.go @@ -0,0 +1,139 @@ +// Copyright (c) 2015-2026 MinIO, Inc. +// +// This file is part of MinIO Object Storage stack +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package cmd + +import ( + "net/http" + "strings" + "testing" + + "github.com/minio/minio/internal/auth" + objectlock "github.com/minio/minio/internal/bucket/object/lock" + xhttp "github.com/minio/minio/internal/http" +) + +// enableBucketObjectLock puts a lock-enabled configuration on an existing +// bucket so checkPutObjectLockAllowed accepts a legal-hold header instead of +// rejecting the request with ErrInvalidBucketObjectLockConfiguration. +func enableBucketObjectLock(t *testing.T, bucket string) { + t.Helper() + meta, err := globalBucketMetadataSys.Get(bucket) + if err != nil { + t.Fatalf("unable to read bucket metadata for %s: %v", bucket, err) + } + updated := meta + updated.ObjectLockConfigXML = enabledBucketObjectLockConfig + updated.VersioningConfigXML = enabledBucketVersioningConfig + // The XML alone is not enough: BucketMetadata keeps a parsed copy that the + // lookups actually read, and it is only populated by parseAllConfigs. + if err := updated.parseAllConfigs(t.Context(), newObjectLayerFn()); err != nil { + t.Fatalf("unable to parse bucket metadata for %s: %v", bucket, err) + } + globalBucketMetadataSys.Set(bucket, updated) +} + +// legalHoldHeaders returns the forwarded legal-hold headers the remote saw, +// separated into the real Object Lock header and the user-metadata spelling +// minio-go produces for an unrecognised UserMetadata key. +func (c *federationRemoteCapture) legalHoldHeaders() (typed, asMetadata []string) { + c.mu.Lock() + defer c.mu.Unlock() + metaKey := "X-Amz-Meta-" + xhttp.AmzObjectLockLegalHold + for _, h := range c.headers { + for k, v := range h { + switch { + case strings.EqualFold(k, xhttp.AmzObjectLockLegalHold): + typed = append(typed, strings.Join(v, ",")) + case strings.EqualFold(k, metaKey): + asMetadata = append(asMetadata, strings.Join(v, ",")) + } + } + } + return typed, asMetadata +} + +// TestAPIFederatedCopyObjectLegalHold drives the legacy etcd federation branch +// of CopyObjectHandler with an explicit legal hold on the copy. +// +// Before the fix the resolved hold was forwarded inside +// PutObjectOptions.UserMetadata. minio-go's Header() prefixes every +// UserMetadata key it does not recognise with "x-amz-meta-", and +// x-amz-object-lock-legal-hold is in neither supportedHeaders nor isAmzHeader, +// so the hold reached the remote as X-Amz-Meta-X-Amz-Object-Lock-Legal-Hold. +// The destination stored no hold and the copy still answered 200 -- a silent +// loss of a WORM control (#166). Retention requested on the same copy survived, +// which is what made it easy to miss. +func TestAPIFederatedCopyObjectLegalHold(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + objAPITest: testAPIFederatedCopyObjectLegalHold, + endpoints: []string{"CopyObject", "PutObject", "HeadObject", "GetObject"}, + }) +} + +func testAPIFederatedCopyObjectLegalHold(objectAPI ObjectLayer, instanceType, bucketName string, + apiRouter http.Handler, credentials auth.Credentials, t *testing.T, +) { + data := []byte("federated copy with a legal hold") + srcObject := "federation/legal-hold-source" + putCopyChecksumSource(t, apiRouter, credentials, bucketName, srcObject, data, nil) + + remoteBucket, capture, cleanup := setupCopyObjectFederation(t, objectAPI, apiRouter, instanceType, bucketName) + defer cleanup() + + // Both roles read bucket metadata from the shared backend in this fixture, + // so one configuration covers the proxy's own check and the remote's. + enableBucketObjectLock(t, bucketName) + enableBucketObjectLock(t, remoteBucket) + + dstObject := "federation/legal-hold-destination" + rec := federatedCopyRequest(t, apiRouter, credentials, bucketName, srcObject, remoteBucket, dstObject, + map[string]string{xhttp.AmzObjectLockLegalHold: string(objectlock.LegalHoldOn)}) + if rec.Code != http.StatusOK { + t.Fatalf("%s: federated CopyObject with a legal hold failed: %d %s", + instanceType, rec.Code, rec.Body.String()) + } + + // The wire is the point: the hold must arrive as the Object Lock header, + // never as user metadata. A 200 with the metadata spelling is exactly the + // silent loss this test exists for. + typed, asMetadata := capture.legalHoldHeaders() + if len(asMetadata) != 0 { + t.Fatalf("%s: legal hold forwarded as user metadata %v; the destination stores no hold", + instanceType, asMetadata) + } + if len(typed) == 0 { + t.Fatalf("%s: no %s header reached the remote deployment", instanceType, xhttp.AmzObjectLockLegalHold) + } + for _, got := range typed { + if !strings.EqualFold(got, string(objectlock.LegalHoldOn)) { + t.Fatalf("%s: forwarded legal hold = %q, want %q", instanceType, got, objectlock.LegalHoldOn) + } + } + + // And it must actually be stored on the destination version. + oi, err := objectAPI.GetObjectInfo(t.Context(), remoteBucket, dstObject, ObjectOptions{}) + if err != nil { + t.Fatalf("%s: unable to stat the federated copy destination: %v", instanceType, err) + } + if hold := objectlock.GetObjectLegalHoldMeta(oi.UserDefined); hold.Status != objectlock.LegalHoldOn { + t.Fatalf("%s: destination legal hold = %q, want %q (metadata: %v)", + instanceType, hold.Status, objectlock.LegalHoldOn, oi.UserDefined) + } +} diff --git a/cmd/object-handlers.go b/cmd/object-handlers.go index c8fbee472..315047035 100644 --- a/cmd/object-handlers.go +++ b/cmd/object-handlers.go @@ -1945,11 +1945,36 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re delete(srcInfo.UserDefined, k) } } + // Legal hold does not survive the metadata map. minio-go's Header() + // writes the typed lock fields first, then prefixes every UserMetadata + // key it does not recognise with "x-amz-meta-": supportedHeaders lists + // x-amz-object-lock-mode and x-amz-object-lock-retain-until-date but not + // x-amz-object-lock-legal-hold, and isAmzHeader does not match it + // either. Forwarded in the map the hold arrives as + // X-Amz-Meta-X-Amz-Object-Lock-Legal-Hold, the destination stores no + // hold, and the copy still answers 200 (#166). + // + // Carry only the hold on the typed field, and forward a clone without + // the raw key: typed fields are written before the UserMetadata loop, so + // a leftover raw key would add a bogus x-amz-meta- entry beside the + // correct header. Retention stays in the map -- it already passes + // through as a standard header, and moving it to the typed + // RetainUntilDate field would truncate the date to whole seconds. + legalHoldKey := strings.ToLower(xhttp.AmzObjectLockLegalHold) + forwardedLegalHold := srcInfo.UserDefined[legalHoldKey] + forwardedMeta := srcInfo.UserDefined + if forwardedLegalHold != "" { + forwardedMeta = cloneMSS(srcInfo.UserDefined) + delete(forwardedMeta, legalHoldKey) + } opts := miniogo.PutObjectOptions{ - UserMetadata: srcInfo.UserDefined, + UserMetadata: forwardedMeta, ServerSideEncryption: dstOpts.ServerSideEncryption, UserTags: tag.ToMap(), } + if forwardedLegalHold != "" { + opts.LegalHold = miniogo.LegalHoldStatus(forwardedLegalHold) + } // The destination must carry the same checksum the local path would // produce; the federated path has the remote compute, validate, persist // and return it. Without this the federated copy silently returns an @@ -1996,7 +2021,10 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re writeErrorResponse(ctx, w, toAPIError(ctx, rerr), r.URL) return } - objInfo.UserDefined = cloneMSS(opts.UserMetadata) + // Built from the resolved values rather than the forwarding map: the + // legal hold was moved onto the typed option above, so opts.UserMetadata + // no longer carries it and the response and event would under-report. + objInfo.UserDefined = cloneMSS(srcInfo.UserDefined) // A forwarded checksum header is a request detail, not object metadata. if checksumHeaderValue != "" { delete(objInfo.UserDefined, wantChecksumType.Key()) From 87d8b5967f7d5b1d3be5c5b6905b0e64653539ff Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Fri, 11 Sep 2026 16:24:42 +0800 Subject: [PATCH 2/6] fix(auth): align signed request and policy condition semantics Signed-off-by: Feng Ruohang --- cmd/bucket-policy.go | 33 ++- cmd/bucket-policy_test.go | 40 +++- cmd/signature-v4-boundary_test.go | 324 ++++++++++++++++++++++++++++++ cmd/signature-v4-utils.go | 24 +-- cmd/signature-v4-utils_test.go | 13 +- cmd/signature-v4.go | 2 - cmd/signature-v4_test.go | 18 +- internal/http/headers.go | 5 - 8 files changed, 412 insertions(+), 47 deletions(-) create mode 100644 cmd/signature-v4-boundary_test.go diff --git a/cmd/bucket-policy.go b/cmd/bucket-policy.go index 9657abf98..567ab0a1c 100644 --- a/cmd/bucket-policy.go +++ b/cmd/bucket-policy.go @@ -255,13 +255,32 @@ func getConditionValuesWithTags(r *http.Request, lc string, cred auth.Credential } cloneHeader := r.Header.Clone() - signatureAge := cloneHeader.Get(xhttp.AmzSignatureAge) - cloneHeader.Del(xhttp.AmzSignatureAge) - // The presigned V4 verifier overwrites this internal scratch header after - // validating the signature. Ignore a value supplied on every other request - // type, where it would otherwise synthesize s3:signatureAge. - if authType == authTypePresigned && signatureAge != "" { - args["signatureAge"] = []string{signatureAge} + + // s3:signatureAge is derived from the presigned X-Amz-Date rather than from + // anything the verifier writes back: PutObject and UploadPart authorize + // before they verify the signature, so a post-verification value is not yet + // available on the first evaluation. The date is bound by the signature + // (doesPresignedSignatureMatch rebuilds and compares it), so a forged date + // only changes the authorization outcome of a request that then fails + // verification. A date that does not parse leaves the key absent; the + // verifier rejects the request as ErrMalformedPresignedDate. + if authType == authTypePresigned { + if signedDate, err := time.Parse(iso8601Format, r.Form.Get(xhttp.AmzDate)); err == nil { + args["signatureAge"] = []string{strconv.FormatInt(currTime.Sub(signedDate).Milliseconds(), 10)} + } + } + + // s3:x-amz-content-sha256 must name the payload hash the request is actually + // verified and enforced against, and only one such value. Presence of the + // header controls whether the key exists at all (AWS documents that the + // query-string form does not populate it), but the value comes from the same + // selection getContentSha256Cksum makes for verification: the presigned query + // value takes precedence over the header, and a repeated header contributes + // only its first value. Exposing every raw header value instead let a + // request satisfy a policy with a value the verifier never checked. + if _, ok := cloneHeader[xhttp.AmzContentSha256]; ok { + args[xhttp.AmzContentSha256] = []string{getContentSha256Cksum(r, serviceS3)} + cloneHeader.Del(xhttp.AmzContentSha256) } userTags := cloneHeader.Get(xhttp.AmzObjectTagging) diff --git a/cmd/bucket-policy_test.go b/cmd/bucket-policy_test.go index 7bd5bb824..42dbb1d8a 100644 --- a/cmd/bucket-policy_test.go +++ b/cmd/bucket-policy_test.go @@ -23,8 +23,10 @@ import ( "net/url" "os" "slices" + "strconv" "strings" "testing" + "time" "github.com/minio/minio/internal/auth" "github.com/minio/minio/internal/handlers" @@ -579,8 +581,20 @@ func TestGetConditionValuesRejectsAbsentInternalKeys(t *testing.T) { } } +// s3:signatureAge is derived from the presigned X-Amz-Date, which the signature +// binds. A client header under the former scratch name must never supply it on +// any auth type, and a presign whose date is missing or malformed leaves the +// key absent (the verifier then rejects the request). func TestGetConditionValuesOnlyAcceptsPresignedSignatureAge(t *testing.T) { const signatureAgeHeader = "x-amz-signature-age" + signedDate := UTCNow().Add(-90 * time.Second) + presignQuery := func(date string) string { + q := url.Values{xhttp.AmzCredential: {"access/20260803/us-east-1/s3/aws4_request"}} + if date != "" { + q.Set(xhttp.AmzDate, date) + } + return "http://minio.local/bkt/obj?" + q.Encode() + } for _, tc := range []struct { name string @@ -602,19 +616,33 @@ func TestGetConditionValuesOnlyAcceptsPresignedSignatureAge(t *testing.T) { }, }, { - name: "presigned verifier value", - target: "http://minio.local/bkt/obj?" + url.Values{ - xhttp.AmzCredential: {"access/20260803/us-east-1/s3/aws4_request"}, - }.Encode(), + name: "presigned client header without date", + target: presignQuery(""), + headers: map[string]string{signatureAgeHeader: "250"}, + }, + { + name: "presigned malformed date", + target: presignQuery("yesterday"), + }, + { + name: "presigned signed date", + target: presignQuery(signedDate.Format(iso8601Format)), headers: map[string]string{signatureAgeHeader: "250"}, want: true, }, } { t.Run(tc.name, func(t *testing.T) { got := condValuesForRequest(t, tc.target, tc.headers) - _, ok := got["signatureAge"] + v, ok := got["signatureAge"] if ok != tc.want { - t.Fatalf("signatureAge presence: expected %v, got %v", tc.want, got["signatureAge"]) + t.Fatalf("signatureAge presence: expected %v, got %v", tc.want, v) + } + if !tc.want { + return + } + age, err := strconv.ParseInt(strings.Join(v, ""), 10, 64) + if err != nil || age < (90*time.Second).Milliseconds() || age > (2*time.Minute).Milliseconds() { + t.Fatalf("signatureAge = %v, want about 90s derived from X-Amz-Date rather than the client header", v) } }) } diff --git a/cmd/signature-v4-boundary_test.go b/cmd/signature-v4-boundary_test.go new file mode 100644 index 000000000..b513b8d71 --- /dev/null +++ b/cmd/signature-v4-boundary_test.go @@ -0,0 +1,324 @@ +// Copyright (c) 2026 PGSTY +// SPDX-License-Identifier: AGPL-3.0-only + +package cmd + +import ( + "bytes" + "encoding/xml" + "fmt" + "net/http" + "net/http/httptest" + "os" + "slices" + "strconv" + "strings" + "testing" + "time" + + "github.com/minio/madmin-go/v3" + "github.com/minio/minio/internal/auth" + xhttp "github.com/minio/minio/internal/http" + "github.com/pgsty/silo-pkg/v3/policy" +) + +// forgedSignatureAgeHeader is the name of the scratch header the presigned +// verifier once wrote s3:signatureAge through. Production code no longer knows +// it; a client that sends it is sending an ordinary unsigned x-amz-* header. +const forgedSignatureAgeHeader = "X-Amz-Signature-Age" + +// Presign at a chosen time, including exactly the supplied operation headers. +func presignBoundaryRequest(t *testing.T, r *http.Request, date time.Time, signedHeaders []string, cred auth.Credentials) { + t.Helper() + query := r.URL.Query() + query.Del(xhttp.AmzSignature) + query.Set(xhttp.AmzAlgorithm, signV4Algorithm) + query.Set(xhttp.AmzDate, date.Format(iso8601Format)) + query.Set(xhttp.AmzExpires, "3600") + query.Set(xhttp.AmzSignedHeaders, strings.Join(signedHeaders, ";")) + query.Set(xhttp.AmzCredential, cred.AccessKey+"/"+getScope(date, globalSite.Region())) + r.Form = query + headers, code := extractSignedHeaders(signedHeaders, r) + if code != ErrNone { + t.Fatal(niceError(code)) + } + canonical := getCanonicalRequest(headers, getContentSha256Cksum(r, serviceS3), query.Encode(), r.URL.Path, r.Method) + key := getSigningKey(cred.SecretKey, date, globalSite.Region(), serviceS3) + query.Set(xhttp.AmzSignature, getSignature(key, getStringToSign(canonical, date, getScope(date, globalSite.Region())))) + r.URL.RawQuery = query.Encode() + r.Form = query +} + +func setupSignatureBoundaryTest(t *testing.T) { + t.Helper() + obj, fsDir, err := prepareFS(t.Context()) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { os.RemoveAll(fsDir) }) + if err = newTestConfig(globalMinioDefaultRegion, obj); err != nil { + t.Fatal(err) + } +} + +// Canonicalization comma-joins repeated header fields, so a signature over one +// x-amz-copy-source value containing a comma also covers the same text split +// into two fields, while the copy handlers act on Header.Get alone. A single +// value may contain a literal or percent-encoded comma; a repeated header is +// rejected at the shared SigV4 boundary for both signed and presigned requests. +func TestV4CopySourceMultiplicity(t *testing.T) { + setupSignatureBoundaryTest(t) + for _, presigned := range []bool{false, true} { + for _, source := range []string{"/source/allowed,tail", "/source/allowed%2Ctail"} { + t.Run(fmt.Sprintf("presigned=%v/%s", presigned, source), func(t *testing.T) { + r := httptest.NewRequest(http.MethodPut, "http://minio.local/destination/object", nil) + r.Header.Set(xhttp.AmzContentSha256, emptySHA256) + r.Header.Set(xhttp.AmzCopySource, source) + if presigned { + presignBoundaryRequest(t, r, UTCNow(), []string{"host", "x-amz-copy-source"}, globalActiveCred) + } else { + if err := signRequestV4(r, globalActiveCred.AccessKey, globalActiveCred.SecretKey); err != nil { + t.Fatal(err) + } + r.Form = r.URL.Query() + } + if code := reqSignatureV4Verify(r, globalSite.Region(), serviceS3); code != ErrNone { + t.Fatalf("a single source key containing a comma must remain valid: %s", niceError(code)) + } + if source != "/source/allowed,tail" { + return + } + r.Header[xhttp.AmzCopySource] = []string{"/source/allowed", "tail"} + if code := reqSignatureV4Verify(r, globalSite.Region(), serviceS3); code != ErrInvalidCopySource { + t.Fatalf("split copy source: got %s, want %s; handler would copy %q", + niceError(code), niceError(ErrInvalidCopySource), r.Header.Get(xhttp.AmzCopySource)) + } + }) + } + } +} + +// PutObject and UploadPart authorize before they verify the signature, so +// s3:signatureAge must come from the signed X-Amz-Date on the first policy +// evaluation. A client header under the former scratch name must neither +// supply the value nor survive verification, and verifying the same request +// twice must give the same answer. +func TestGetConditionValuesPresignedAgeFromDate(t *testing.T) { + setupSignatureBoundaryTest(t) + for _, tc := range []struct { + name, header string + age time.Duration + wantVerify APIErrorCode + }{ + {name: "old", age: 10 * time.Minute}, + {name: "old forged 0", age: 10 * time.Minute, header: "0", wantVerify: ErrUnsignedHeaders}, + {name: "old forged 1", age: 10 * time.Minute, header: "1", wantVerify: ErrUnsignedHeaders}, + // A signer slightly ahead of the server stays within globalMaxSkewTime. + {name: "future within skew", age: -2 * time.Minute}, + } { + t.Run(tc.name, func(t *testing.T) { + r := httptest.NewRequest(http.MethodPut, "http://minio.local/bucket/object", nil) + presignBoundaryRequest(t, r, UTCNow().Add(-tc.age), []string{"host"}, globalActiveCred) + if tc.header != "" { + r.Header.Set(forgedSignatureAgeHeader, tc.header) + } + values := getConditionValues(r, "", globalActiveCred) + age, err := strconv.ParseInt(strings.Join(values["signatureAge"], ""), 10, 64) + lo, hi := (tc.age - time.Minute).Milliseconds(), (tc.age + time.Minute).Milliseconds() + if err != nil || age < lo || age > hi { + t.Errorf("pre-verification policy got age %v; want the age of the signed date (about %d ms)", values["signatureAge"], tc.age.Milliseconds()) + } + code := reqSignatureV4Verify(r, globalSite.Region(), serviceS3) + if code != tc.wantVerify { + t.Fatalf("verification: got %s, want %s", niceError(code), niceError(tc.wantVerify)) + } + if again := reqSignatureV4Verify(r, globalSite.Region(), serviceS3); again != code { + t.Fatalf("second verification changed the outcome: %s -> %s", niceError(code), niceError(again)) + } + }) + } +} + +// The s3:x-amz-content-sha256 policy value must be the single payload hash the +// request is verified and enforced against. Header presence decides whether the +// key exists at all; the value is the one getContentSha256Cksum selects, so a +// presigned query value wins over the header and a repeated header contributes +// only its first value. None of these requests is rejected at the protocol +// level; the policy simply sees what verification bound. +func TestGetConditionValuesPayloadHashMatchesVerifiedValue(t *testing.T) { + setupSignatureBoundaryTest(t) + hashA, hashB := getSHA256Hash([]byte("a")), getSHA256Hash([]byte("b")) + for _, tc := range []struct { + name string + presigned bool + query string + header []string + want []string + }{ + {name: "signed header", header: []string{hashA}, want: []string{hashA}}, + {name: "signed absent"}, + {name: "signed present empty", header: []string{""}, want: []string{""}}, + {name: "signed duplicate header", header: []string{hashA, hashB}, want: []string{hashA}}, + {name: "signed streaming with second value", header: []string{streamingContentSHA256, hashA}, want: []string{streamingContentSHA256}}, + {name: "presigned query only", presigned: true, query: hashA}, + {name: "presigned header only", presigned: true, header: []string{hashA}, want: []string{hashA}}, + {name: "presigned matching query and header", presigned: true, query: hashA, header: []string{hashA}, want: []string{hashA}}, + {name: "presigned unsigned query with forged header", presigned: true, query: unsignedPayload, header: []string{hashA}, want: []string{unsignedPayload}}, + {name: "presigned duplicate header", presigned: true, header: []string{hashA, hashB}, want: []string{hashA}}, + {name: "presigned query with empty header", presigned: true, query: hashA, header: []string{""}, want: []string{hashA}}, + } { + t.Run(tc.name, func(t *testing.T) { + target := "http://minio.local/bucket/object" + if tc.query != "" { + target += "?" + xhttp.AmzContentSha256 + "=" + tc.query + } + r := httptest.NewRequest(http.MethodPut, target, nil) + if tc.header != nil { + r.Header[xhttp.AmzContentSha256] = tc.header + } + if tc.presigned { + presignBoundaryRequest(t, r, UTCNow(), []string{"host"}, globalActiveCred) + } else { + r.Header.Set(xhttp.Authorization, signV4Algorithm+" Credential=x/20260910/us-east-1/s3/aws4_request, SignedHeaders=host, Signature=x") + r.Form = r.URL.Query() + } + got, ok := getConditionValues(r, "", globalActiveCred)[xhttp.AmzContentSha256] + if ok != (tc.want != nil) || !slices.Equal(got, tc.want) { + t.Fatalf("policy value = %v (present=%v), want %v (present=%v)", got, ok, tc.want, tc.want != nil) + } + if tc.presigned { + if code := reqSignatureV4Verify(r, globalSite.Region(), serviceS3); code != ErrNone { + t.Fatalf("the presigned request must still verify: %s", niceError(code)) + } + } + }) + } +} + +func newSignatureBoundaryUser(t *testing.T, bucket, statements string) auth.Credentials { + t.Helper() + cred, err := auth.GetNewCredentials() + if err != nil { + t.Fatal(err) + } + if _, err := globalIAMSys.CreateUser(t.Context(), cred.AccessKey, madmin.AddOrUpdateUserReq{SecretKey: cred.SecretKey, Status: madmin.AccountEnabled}); err != nil { + t.Fatal(err) + } + p, err := policy.ParseConfig(strings.NewReader(fmt.Sprintf(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Action":"s3:PutObject","Resource":"arn:aws:s3:::%s/*"},%s]}`, bucket, statements))) + if err != nil { + t.Fatal(err) + } + name := "signature-boundary-" + mustGetUUID() + if _, err := globalIAMSys.SetPolicy(t.Context(), name, *p); err != nil { + t.Fatal(err) + } + if _, err := globalIAMSys.PolicyDBSet(t.Context(), cred.AccessKey, name, regUser, false); err != nil { + t.Fatal(err) + } + return cred +} + +func TestAPIPresignedSignatureAgeBeforeAuthorization(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + endpoints: []string{"NewMultipart", "PutObjectPart", "PutObject"}, + objAPITest: func(obj ObjectLayer, instanceType, bucket string, router http.Handler, root auth.Credentials, t *testing.T) { + user := newSignatureBoundaryUser(t, bucket, fmt.Sprintf(`{"Effect":"Deny","Action":"s3:PutObject","Resource":"arn:aws:s3:::%s/*","Condition":{"NumericGreaterThan":{"s3:signatureAge":"60000"}}}`, bucket)) + initReq, err := newTestSignedRequestV4(http.MethodPost, getNewMultipartURL("", bucket, "multipart"), 0, nil, root.AccessKey, root.SecretKey, nil) + if err != nil { + t.Fatal(err) + } + initRec := httptest.NewRecorder() + router.ServeHTTP(initRec, initReq) + var upload InitiateMultipartUploadResponse + if initRec.Code != http.StatusOK || xml.Unmarshal(initRec.Body.Bytes(), &upload) != nil { + t.Fatalf("multipart initiation: %d %s", initRec.Code, initRec.Body.String()) + } + for _, operation := range []string{"put", "part"} { + for _, tc := range []struct { + name, header string + age time.Duration + want int + }{ + {name: "fresh", want: http.StatusOK}, + {name: "old without header", age: 10 * time.Minute, want: http.StatusForbidden}, + {name: "old forged header", age: 10 * time.Minute, header: "0", want: http.StatusForbidden}, + // Policy allows a fresh signature; the unsigned header then fails + // verification with the existing ErrUnsignedHeaders (HTTP 400). + {name: "fresh with unsigned header", header: "0", want: http.StatusBadRequest}, + } { + t.Run(instanceType+"/"+operation+"/"+tc.name, func(t *testing.T) { + target := getPutObjectURL("", bucket, "put-"+strings.ReplaceAll(tc.name, " ", "-")) + if operation == "part" { + target = getPutObjectPartURL("", bucket, "multipart", upload.UploadID, "1") + } + r, err := newTestRequest(http.MethodPut, target, 4, bytes.NewReader([]byte("body"))) + if err != nil { + t.Fatal(err) + } + r.Header.Del(xhttp.AmzContentSha256) + presignBoundaryRequest(t, r, UTCNow().Add(-tc.age), []string{"host"}, user) + if tc.header != "" { + r.Header.Set(forgedSignatureAgeHeader, tc.header) + } + rec := httptest.NewRecorder() + router.ServeHTTP(rec, r) + if rec.Code != tc.want { + t.Errorf("got %d %s, want %d", rec.Code, rec.Body.String(), tc.want) + } + }) + } + } + }, + }) +} + +func TestAPIPayloadHashPolicyMatchesVerifiedValue(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + endpoints: []string{"PutObject"}, + objAPITest: func(_ ObjectLayer, instanceType, bucket string, router http.Handler, _ auth.Credentials, t *testing.T) { + allowedHash := getSHA256Hash([]byte("expected")) + user := newSignatureBoundaryUser(t, bucket, fmt.Sprintf(`{"Effect":"Deny","Action":"s3:PutObject","Resource":"arn:aws:s3:::%s/*","Condition":{"StringNotEquals":{"s3:x-amz-content-sha256":"%s"}}}`, bucket, allowedHash)) + for _, kind := range []string{"signed control", "presigned control", "unsigned query forged header", "signed duplicate header", "presigned duplicate header"} { + t.Run(instanceType+"/"+kind, func(t *testing.T) { + control := strings.HasSuffix(kind, "control") + body := "modified" + if control { + body = "expected" + } + r, err := newTestRequest(http.MethodPut, getPutObjectURL("", bucket, strings.ReplaceAll(kind, " ", "-")), int64(len(body)), strings.NewReader(body)) + if err != nil { + t.Fatal(err) + } + if strings.Contains(kind, "duplicate") { + r.Header.Add(xhttp.AmzContentSha256, allowedHash) + } + if kind == "unsigned query forged header" { + q := r.URL.Query() + q.Set(xhttp.AmzContentSha256, unsignedPayload) + r.URL.RawQuery = q.Encode() + r.Header.Set(xhttp.AmzContentSha256, allowedHash) + } + if strings.HasPrefix(kind, "signed") { + if err := signRequestV4(r, user.AccessKey, user.SecretKey); err != nil { + t.Fatal(err) + } + } else { + presignBoundaryRequest(t, r, UTCNow(), []string{"host"}, user) + } + rec := httptest.NewRecorder() + router.ServeHTTP(rec, r) + if control && rec.Code != http.StatusOK { + t.Errorf("control rejected: %d %s", rec.Code, rec.Body.String()) + } + if !control && rec.Code < http.StatusBadRequest { + t.Errorf("payload-hash policy bypass returned %d %s", rec.Code, rec.Body.String()) + } + }) + } + }, + }) +} diff --git a/cmd/signature-v4-utils.go b/cmd/signature-v4-utils.go index 88ed7e9cd..f3f6815c6 100644 --- a/cmd/signature-v4-utils.go +++ b/cmd/signature-v4-utils.go @@ -211,7 +211,16 @@ func extractSignedHeaders(signedHeaders []string, r *http.Request) (http.Header, // `host` will not be found in the headers, can be found in r.Host. // but its always necessary that the list of signed headers containing host in it. val, ok := reqHeaders[http.CanonicalHeaderKey(header)] - if !ok { + if ok { + // Canonicalization comma-joins repeated header fields, so a signature + // over the single value "/src/a,tail" also verifies a request carrying + // ["/src/a", "tail"]. The copy handlers read only Header.Get, so that + // rewrite would copy a different source than the one signed. Reject a + // repeated x-amz-copy-source; a single value may still contain commas. + if len(val) > 1 && strings.EqualFold(header, strings.ToLower(xhttp.AmzCopySource)) { + return nil, ErrInvalidCopySource + } + } else { // try to set headers from Query String val, ok = reqQueries[header] } @@ -274,9 +283,8 @@ func signV4TrimAll(input string) string { // object the signing key can reach. // // Only headers actually sent by the client are inspected. Server-synthesized -// x-amz-* headers (e.g. x-amz-tagging derived from a request body, or the -// post-verification x-amz-signature-age scratch header) are set after signature -// verification and therefore never reach this walk. +// x-amz-* headers (e.g. x-amz-tagging derived from a request body) are set +// after signature verification and therefore never reach this walk. func checkUnsignedHeaders(signedHeadersMap http.Header, r *http.Request) APIErrorCode { // check headers that arrived on the request for k := range r.Header { @@ -294,14 +302,6 @@ func checkUnsignedHeaders(signedHeadersMap http.Header, r *http.Request) APIErro if strings.EqualFold(k, xhttp.AmzContentSha256) { continue } - // X-Amz-Signature-Age is an internal scratch header written by the - // presigned verifier itself, after this check, purely so bucket-policy - // evaluation can expose s3:signatureAge. It is never sent or signed by a - // client, and exempting it keeps signature verification idempotent when - // the same request is verified more than once. - if strings.EqualFold(k, xhttp.AmzSignatureAge) { - continue - } // The header must be a member of the signed-headers list. Testing // membership (not value equality) is essential: an unsigned header whose // first value is empty would otherwise compare equal to the empty string diff --git a/cmd/signature-v4-utils_test.go b/cmd/signature-v4-utils_test.go index d957908b9..8cccf6116 100644 --- a/cmd/signature-v4-utils_test.go +++ b/cmd/signature-v4-utils_test.go @@ -469,15 +469,16 @@ func TestCheckUnsignedHeaders(t *testing.T) { t.Fatalf("unsigned x-amz-content-sha256 must be exempt: expected %d, got %d", ErrNone, errCode) } - // X-Amz-Signature-Age is the presigned verifier's own scratch header, - // written after this check. Exempting it keeps verification idempotent when - // the same request object is verified more than once. + // X-Amz-Signature-Age was once a scratch header the presigned verifier wrote + // back and this check exempted. s3:signatureAge is now derived from the + // signed date, so a client that sends the header is sending an ordinary + // unsigned x-amz-* header and must be rejected like any other. r, err = http.NewRequest(http.MethodPut, "http://play.min.io:9000", nil) if err != nil { t.Fatal("Unable to create http.Request :", err) } - r.Header.Set(xhttp.AmzSignatureAge, "1234") - if errCode = checkUnsignedHeaders(signedHeadersMap, r); errCode != ErrNone { - t.Fatalf("internal x-amz-signature-age must be exempt: expected %d, got %d", ErrNone, errCode) + r.Header.Set("X-Amz-Signature-Age", "1234") + if errCode = checkUnsignedHeaders(signedHeadersMap, r); errCode != ErrUnsignedHeaders { + t.Fatalf("unsigned x-amz-signature-age must be rejected: expected %d, got %d", ErrUnsignedHeaders, errCode) } } diff --git a/cmd/signature-v4.go b/cmd/signature-v4.go index f204cd33c..d0125fb1c 100644 --- a/cmd/signature-v4.go +++ b/cmd/signature-v4.go @@ -336,8 +336,6 @@ func doesPresignedSignatureMatch(hashedPayload string, r *http.Request, region s return ErrSignatureDoesNotMatch } - r.Header.Set(xhttp.AmzSignatureAge, strconv.FormatInt(UTCNow().Sub(pSignValues.Date).Milliseconds(), 10)) - return ErrNone } diff --git a/cmd/signature-v4_test.go b/cmd/signature-v4_test.go index 72129940d..4f8b6dc07 100644 --- a/cmd/signature-v4_test.go +++ b/cmd/signature-v4_test.go @@ -25,8 +25,6 @@ import ( "os" "testing" "time" - - xhttp "github.com/minio/minio/internal/http" ) func niceError(code APIErrorCode) string { @@ -317,11 +315,11 @@ func TestDoesPresignedSignatureMatch(t *testing.T) { } // TestPresignedVerifyIdempotent guards against a regression where verifying the -// same presigned request twice began to fail. doesPresignedSignatureMatch -// writes an internal x-amz-signature-age header after validating the signature; -// the unsigned-header check must exempt that scratch header (and an unsigned -// x-amz-content-sha256 the client may carry) so a second verification of the -// same *http.Request still succeeds. +// same presigned request twice began to fail. The verifier must not write +// anything back to the request (it once recorded an internal +// x-amz-signature-age header that the unsigned-header check then had to +// exempt), and an unsigned x-amz-content-sha256 the client may carry stays +// exempt, so a second verification of the same *http.Request still succeeds. func TestPresignedVerifyIdempotent(t *testing.T) { ctx, cancel := context.WithCancel(t.Context()) defer cancel() @@ -350,7 +348,9 @@ func TestPresignedVerifyIdempotent(t *testing.T) { t.Fatalf("first verification: expected ErrNone, got %s", niceError(got)) } if got := reqSignatureV4Verify(req, globalSite.Region(), serviceS3); got != ErrNone { - t.Fatalf("second verification of the same request: expected ErrNone, got %s (x-amz-signature-age=%q)", - niceError(got), req.Header.Get(xhttp.AmzSignatureAge)) + t.Fatalf("second verification of the same request: expected ErrNone, got %s", niceError(got)) + } + if _, ok := req.Header["X-Amz-Signature-Age"]; ok { + t.Fatal("the verifier wrote a scratch header back to the request") } } diff --git a/internal/http/headers.go b/internal/http/headers.go index 93e55264f..9195e122e 100644 --- a/internal/http/headers.go +++ b/internal/http/headers.go @@ -132,11 +132,6 @@ const ( AmzMaxParts = "X-Amz-Max-Parts" AmzPartNumberMarker = "X-Amz-Part-Number-Marker" - // AmzSignatureAge is an internal scratch header the presigned verifier - // writes after validating the signature so that bucket-policy evaluation can - // expose s3:signatureAge. It is never sent or signed by a client. - AmzSignatureAge = "X-Amz-Signature-Age" - // Constants used for GetObjectAttributes and GetObjectVersionAttributes AmzObjectAttributes = "X-Amz-Object-Attributes" From 9a303f5096c572c74efaca605f50d3c511e9885b Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Fri, 11 Sep 2026 16:24:42 +0800 Subject: [PATCH 3/6] fix(object): preserve retention and federated copy destination state Signed-off-by: Feng Ruohang --- cmd/bucket-object-lock.go | 9 +- cmd/object-copy-federation-legalhold_test.go | 106 ++++++++++++- cmd/object-copy-federation-result_test.go | 158 +++++++++++++++++++ cmd/object-handlers.go | 90 ++++++----- cmd/object-lock-default-retention_test.go | 140 ++++++++++++++++ 5 files changed, 450 insertions(+), 53 deletions(-) create mode 100644 cmd/object-copy-federation-result_test.go create mode 100644 cmd/object-lock-default-retention_test.go diff --git a/cmd/bucket-object-lock.go b/cmd/bucket-object-lock.go index e5828e2dd..36211c45d 100644 --- a/cmd/bucket-object-lock.go +++ b/cmd/bucket-object-lock.go @@ -334,11 +334,10 @@ func checkPutObjectLockAllowed(ctx context.Context, rq *http.Request, bucket, ob return mode, retainDate, legalHold, ErrObjectLocked } - if !legalHoldRequested && retentionCfg.LockEnabled { - // inherit retention from bucket configuration - return retentionCfg.Mode, objectlock.RetentionDate{Time: t.Add(retentionCfg.Validity)}, legalHold, ErrNone - } - return "", objectlock.RetentionDate{}, legalHold, ErrNone + // Inherit retention from the bucket configuration. A legal-hold header + // on the same request, ON or OFF, is independent of retention and must + // not suppress the default (#165). + return retentionCfg.Mode, objectlock.RetentionDate{Time: t.Add(retentionCfg.Validity)}, legalHold, ErrNone } return mode, retainDate, legalHold, ErrNone } diff --git a/cmd/object-copy-federation-legalhold_test.go b/cmd/object-copy-federation-legalhold_test.go index 44e3e5dc5..941719699 100644 --- a/cmd/object-copy-federation-legalhold_test.go +++ b/cmd/object-copy-federation-legalhold_test.go @@ -18,13 +18,16 @@ package cmd import ( + "bytes" "net/http" "strings" "testing" + "time" "github.com/minio/minio/internal/auth" objectlock "github.com/minio/minio/internal/bucket/object/lock" xhttp "github.com/minio/minio/internal/http" + "github.com/minio/minio/internal/kms" ) // enableBucketObjectLock puts a lock-enabled configuration on an existing @@ -47,9 +50,108 @@ func enableBucketObjectLock(t *testing.T, bucket string) { globalBucketMetadataSys.Set(bucket, updated) } +// Exercise #165 and #166 together, including hold OFF, default retention and +// explicit subsecond retention. Ordinary copies must not inherit source locks. +func TestAPIFederatedCopyObjectLockParity(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + endpoints: []string{"CopyObject", "PutObject", "HeadObject", "GetObject"}, + objAPITest: func(obj ObjectLayer, instanceType, bucket string, router http.Handler, cred auth.Credentials, t *testing.T) { + testKMS, err := kms.NewBuiltin(federationTestKMSKeyID, bytes.Repeat([]byte{0x58}, 32)) + if err != nil { + t.Fatal(err) + } + previousKMS := GlobalKMS + GlobalKMS = testKMS + defer func() { GlobalKMS = previousKMS }() + remoteBucket, capture, cleanup := setupCopyObjectFederation(t, obj, router, instanceType, bucket) + defer cleanup() + enableBucketObjectLock(t, bucket) + until := UTCNow().Add(7 * 24 * time.Hour).Truncate(time.Second).Add(789 * time.Millisecond).Format(time.RFC3339Nano) + for _, sourceType := range []string{"plain", "s3"} { + source := sourceType + "-held-source" + headers := federationSSEHeaders(sourceType, 0, false) + headers[xhttp.AmzObjectLockLegalHold] = "ON" + putCopyChecksumSource(t, router, cred, bucket, source, []byte("held source"), headers) + for _, tc := range []struct { + name, hold, defaultMode, explicitMode string + replace bool + }{ + {name: "on", hold: "ON"}, + {name: "off", hold: "OFF"}, + {name: "no source inheritance"}, + {name: "on with default", hold: "ON", defaultMode: "GOVERNANCE"}, + {name: "off with default", hold: "OFF", defaultMode: "COMPLIANCE"}, + {name: "explicit retention", explicitMode: "GOVERNANCE"}, + {name: "hold and explicit retention", hold: "ON", explicitMode: "COMPLIANCE"}, + {name: "replace metadata", hold: "ON", explicitMode: "GOVERNANCE", replace: true}, + } { + t.Run(instanceType+"/"+sourceType+"/"+tc.name, func(t *testing.T) { + setTestBucketDefaultRetention(t, bucket, tc.defaultMode) + setTestBucketDefaultRetention(t, remoteBucket, tc.defaultMode) + headers := map[string]string{} + if tc.hold != "" { + headers[xhttp.AmzObjectLockLegalHold] = tc.hold + } + if tc.explicitMode != "" { + headers[xhttp.AmzObjectLockMode] = tc.explicitMode + headers[xhttp.AmzObjectLockRetainUntilDate] = until + } + if tc.replace { + headers[xhttp.AmzMetadataDirective] = "REPLACE" + headers["X-Amz-Meta-Origin"] = "replacement" + } + capture.mu.Lock() + capture.headers = nil + capture.mu.Unlock() + for _, destinationBucket := range []string{bucket, remoteBucket} { + destination := sourceType + "-copy-" + strings.ReplaceAll(tc.name, " ", "-") + rec := federatedCopyRequest(t, router, cred, bucket, source, destinationBucket, destination, headers) + if rec.Code != http.StatusOK { + t.Fatalf("copy to %s: %d %s", destinationBucket, rec.Code, rec.Body.String()) + } + info, err := obj.GetObjectInfo(t.Context(), destinationBucket, destination, ObjectOptions{}) + if err != nil { + t.Fatal(err) + } + if got := objectlock.GetObjectLegalHoldMeta(info.UserDefined).Status; string(got) != tc.hold { + t.Errorf("stored hold = %q, want %q", got, tc.hold) + } + retention := objectlock.GetObjectRetentionMeta(info.UserDefined) + mode := tc.explicitMode + if mode == "" { + mode = tc.defaultMode + } + if string(retention.Mode) != mode { + t.Errorf("stored retention = %q, want %q", retention.Mode, mode) + } + if tc.explicitMode != "" && retention.RetainUntilDate.Format(time.RFC3339Nano) != until { + t.Errorf("retention date lost precision: %s, want %s", retention.RetainUntilDate, until) + } + for key := range info.UserDefined { + if stringsHasPrefixFold(key, "X-Amz-Meta-X-Amz-Object-Lock-") { + t.Errorf("lock state became user metadata: %s", key) + } + } + if tc.replace && info.UserDefined["X-Amz-Meta-Origin"] != "replacement" { + t.Errorf("replacement metadata was lost: %v", info.UserDefined) + } + } + typed, asMetadata := capture.legalHoldHeaders() + if len(asMetadata) != 0 || (tc.hold != "" && strings.Join(typed, "") != tc.hold) || (tc.hold == "" && len(typed) != 0) { + t.Errorf("forwarded hold headers = %v, metadata = %v; want %q", typed, asMetadata, tc.hold) + } + }) + } + } + }, + }) +} + // legalHoldHeaders returns the forwarded legal-hold headers the remote saw, // separated into the real Object Lock header and the user-metadata spelling -// minio-go produces for an unrecognised UserMetadata key. +// minio-go produces for an unrecognized UserMetadata key. func (c *federationRemoteCapture) legalHoldHeaders() (typed, asMetadata []string) { c.mu.Lock() defer c.mu.Unlock() @@ -72,7 +174,7 @@ func (c *federationRemoteCapture) legalHoldHeaders() (typed, asMetadata []string // // Before the fix the resolved hold was forwarded inside // PutObjectOptions.UserMetadata. minio-go's Header() prefixes every -// UserMetadata key it does not recognise with "x-amz-meta-", and +// UserMetadata key it does not recognize with "x-amz-meta-", and // x-amz-object-lock-legal-hold is in neither supportedHeaders nor isAmzHeader, // so the hold reached the remote as X-Amz-Meta-X-Amz-Object-Lock-Legal-Hold. // The destination stored no hold and the copy still answered 200 -- a silent diff --git a/cmd/object-copy-federation-result_test.go b/cmd/object-copy-federation-result_test.go new file mode 100644 index 000000000..20d70f10b --- /dev/null +++ b/cmd/object-copy-federation-result_test.go @@ -0,0 +1,158 @@ +// Copyright (c) 2026 PGSTY +// SPDX-License-Identifier: AGPL-3.0-only + +package cmd + +import ( + "bytes" + "net/http" + "net/url" + "strings" + "testing" + "time" + + "github.com/minio/minio/internal/auth" + sse "github.com/minio/minio/internal/bucket/encryption" + "github.com/minio/minio/internal/event" + xhttp "github.com/minio/minio/internal/http" + "github.com/minio/minio/internal/kms" + "github.com/minio/minio/internal/pubsub" +) + +func TestAPIFederatedCopyObjectVersionAndEvent(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + endpoints: []string{"CopyObject", "PutObject", "GetObject", "HeadObject"}, + objAPITest: func(obj ObjectLayer, instanceType, bucket string, router http.Handler, cred auth.Credentials, t *testing.T) { + testKMS, err := kms.NewBuiltin(federationTestKMSKeyID, bytes.Repeat([]byte{0x58}, 32)) + if err != nil { + t.Fatal(err) + } + previousKMS := GlobalKMS + GlobalKMS = testKMS + defer func() { GlobalKMS = previousKMS }() + restore := setCopyChecksumCompression(true) + defer restore() + remoteBucket, _, cleanup := setupCopyObjectFederation(t, obj, router, instanceType, bucket) + defer cleanup() + enableBucketObjectLock(t, remoteBucket) + events := make(chan event.Event, 8) + done := make(chan struct{}) + defer close(done) + if err := globalHTTPListen.Subscribe(pubsub.MaskFromMaskable(event.ObjectCreatedCopy), events, done, nil); err != nil { + t.Fatal(err) + } + for _, kind := range []string{"plain", "compressed", "encrypted"} { + t.Run(instanceType+"/"+kind, func(t *testing.T) { + data := []byte("logical object size") + source := kind + "-source.bin" + var headers map[string]string + if kind == "compressed" { + data = bytes.Repeat(data, 8192) + source = kind + "-source.txt" + } + if kind == "encrypted" { + headers = federationSSEHeaders("s3", 0, false) + } + putCopyChecksumSource(t, router, cred, bucket, source, data, headers) + destination := "result/" + kind + " with space.bin" + rec := federatedCopyRequest(t, router, cred, bucket, source, remoteBucket, destination, nil) + if rec.Code != http.StatusOK { + t.Fatalf("copy: %d %s", rec.Code, rec.Body.String()) + } + versionID := strings.Join(rec.Header()[xhttp.AmzVersionID], "") + if versionID == "" { + t.Error("copy response omitted destination version ID") + } else { + info, err := obj.GetObjectInfo(t.Context(), remoteBucket, destination, ObjectOptions{VersionID: versionID}) + if err != nil || info.VersionID != versionID { + t.Fatalf("response does not name the written version: %v, %q", err, info.VersionID) + } + } + select { + case evt := <-events: + key, err := url.QueryUnescape(evt.S3.Object.Key) + if err != nil || key != destination || evt.S3.Bucket.Name != remoteBucket || evt.S3.Object.Size != int64(len(data)) || evt.S3.Object.VersionID == "" || evt.S3.Object.VersionID != versionID { + t.Errorf("copy event does not describe the written object: %+v", evt.S3) + } + case <-time.After(5 * time.Second): + t.Fatal("copy event was not emitted") + } + }) + } + }, + }) +} + +func TestAPIFederatedCopyObjectDestinationSSEDefaults(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + endpoints: []string{"CopyObject", "PutObject", "GetObject", "HeadObject"}, + objAPITest: func(obj ObjectLayer, instanceType, bucket string, router http.Handler, cred auth.Credentials, t *testing.T) { + data := []byte("destination chooses default encryption") + putCopyChecksumSource(t, router, cred, bucket, "source", data, nil) + testKMS, err := kms.NewBuiltin(federationTestKMSKeyID, bytes.Repeat([]byte{0x58}, 32)) + if err != nil { + t.Fatal(err) + } + previousKMS, previousAuto := GlobalKMS, globalAutoEncryption + GlobalKMS, globalAutoEncryption = testKMS, true + defer func() { GlobalKMS, globalAutoEncryption = previousKMS, previousAuto }() + remoteBucket, capture, cleanup := setupCopyObjectFederationRemote(t, obj, router, instanceType, bucket, true) + defer cleanup() + destinationConfig, err := sse.ParseBucketSSEConfig(strings.NewReader(`aws:kmsdestination-bucket-key`)) + if err != nil { + t.Fatal(err) + } + for _, kind := range []string{"plain", "s3", "kms", "kms-default", "c"} { + t.Run(instanceType+"/"+kind, func(t *testing.T) { + var headers map[string]string + if kind == "kms-default" { + headers = map[string]string{xhttp.AmzServerSideEncryption: xhttp.AmzEncryptionKMS} + } else { + headers = federationSSEHeaders(kind, 0x22, false) + } + rec := federatedCopyRequest(t, router, cred, bucket, "source", remoteBucket, kind, headers) + if rec.Code != http.StatusOK { + t.Fatalf("copy: %d %s", rec.Code, rec.Body.String()) + } + capture.mu.Lock() + forwarded := capture.headers[len(capture.headers)-1].Clone() + capture.mu.Unlock() + if kind == "plain" { + if forwarded.Get(xhttp.AmzServerSideEncryption) != "" { + t.Errorf("proxy injected SSE %q into a request with no client SSE", forwarded.Get(xhttp.AmzServerSideEncryption)) + } + // With nothing injected, the KMS encryption the destination + // stores can only be the remote applying its own defaults. + info, err := obj.GetObjectInfo(t.Context(), remoteBucket, kind, ObjectOptions{}) + if err != nil || federationStoredSSE(info.UserDefined) != "kms" { + t.Errorf("remote destination did not apply its own auto-encryption: %v, %v", err, info.UserDefined) + } + } + // Replay the actual inbound headers against an independent bucket + // configuration. No handler flips shared globals while serving. + destinationConfig.Apply(forwarded, sse.ApplyOptions{}) + wantKey := headers[xhttp.AmzServerSideEncryptionKmsID] + if kind == "plain" { + wantKey = "destination-bucket-key" + } + if got := forwarded.Get(xhttp.AmzServerSideEncryptionKmsID); got != wantKey { + t.Errorf("destination selected key %q, want %q", got, wantKey) + } + }) + } + // A local destination still inherits the server's auto-encryption. + rec := federatedCopyRequest(t, router, cred, bucket, "source", bucket, "local-copy", nil) + if rec.Code != http.StatusOK { + t.Fatalf("local copy: %d %s", rec.Code, rec.Body.String()) + } + info, err := obj.GetObjectInfo(t.Context(), bucket, "local-copy", ObjectOptions{}) + if err != nil || federationStoredSSE(info.UserDefined) != "kms" { + t.Errorf("local copy lost auto-encryption: %v, %v", err, info.UserDefined) + } + }, + }) +} diff --git a/cmd/object-handlers.go b/cmd/object-handlers.go index 315047035..a4dfab00b 100644 --- a/cmd/object-handlers.go +++ b/cmd/object-handlers.go @@ -1404,11 +1404,27 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re } allowReplicationMetadata := replicaTrusted - // Check if bucket encryption is enabled - sseConfig, _ := globalBucketSSEConfigSys.Get(dstBucket) - sseConfig.Apply(r.Header, sse.ApplyOptions{ - AutoEncrypt: globalAutoEncryption, - }) + // Federation only: the destination bucket lives on another deployment and + // the copy is forwarded to it as a PutObject. That remote write owns the + // destination's storage transformations, so this handler hands it the + // logical (decompressed, decrypted) bytes at their logical size and lets + // the remote compress and encrypt once. Encrypting here as well would + // forward ciphertext under the destination's own SSE option, which either + // fails the length check or, for SSE to SSE, has the remote encrypt the + // ciphertext a second time and store an unreadable object (#158). + remoteCallRequired := isRemoteCopyRequired(ctx, srcBucket, dstBucket, objectAPI) + + // Apply the destination bucket's default encryption only when this + // deployment writes the destination. For a remote destination the proxy + // has no authority over that bucket's defaults: applying its own here + // would forward an explicit SSE header that the remote then honors in + // place of the destination's configuration (#167). + if !remoteCallRequired { + sseConfig, _ := globalBucketSSEConfigSys.Get(dstBucket) + sseConfig.Apply(r.Header, sse.ApplyOptions{ + AutoEncrypt: globalAutoEncryption, + }) + } var srcOpts, dstOpts ObjectOptions srcOpts, err = copySrcOpts(ctx, r, srcBucket, srcObject) if err != nil { @@ -1518,16 +1534,6 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re } } - // Federation only: the destination bucket lives on another deployment and - // the copy is forwarded to it as a PutObject. That remote write owns the - // destination's storage transformations, so this handler hands it the - // logical (decompressed, decrypted) bytes at their logical size and lets - // the remote compress and encrypt once. Encrypting here as well would - // forward ciphertext under the destination's own SSE option, which either - // fails the length check or, for SSE to SSE, has the remote encrypt the - // ciphertext a second time and store an unreadable object (#158). - remoteCallRequired := isRemoteCopyRequired(ctx, srcBucket, dstBucket, objectAPI) - var compressMetadata map[string]string // No need to compress for remote etcd calls // Pass the decompressed stream to such calls. @@ -1945,35 +1951,23 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re delete(srcInfo.UserDefined, k) } } - // Legal hold does not survive the metadata map. minio-go's Header() - // writes the typed lock fields first, then prefixes every UserMetadata - // key it does not recognise with "x-amz-meta-": supportedHeaders lists - // x-amz-object-lock-mode and x-amz-object-lock-retain-until-date but not - // x-amz-object-lock-legal-hold, and isAmzHeader does not match it - // either. Forwarded in the map the hold arrives as - // X-Amz-Meta-X-Amz-Object-Lock-Legal-Hold, the destination stores no - // hold, and the copy still answers 200 (#166). - // - // Carry only the hold on the typed field, and forward a clone without - // the raw key: typed fields are written before the UserMetadata loop, so - // a leftover raw key would add a bogus x-amz-meta- entry beside the - // correct header. Retention stays in the map -- it already passes - // through as a standard header, and moving it to the typed - // RetainUntilDate field would truncate the date to whole seconds. + // Forward a clone: request-only keys are removed or added below and + // srcInfo.UserDefined stays the resolved record for the response. + forwardedMeta := cloneMSS(srcInfo.UserDefined) + // minio-go prefixes any UserMetadata key it does not recognize with + // x-amz-meta-, and it recognizes the retention headers but not + // x-amz-object-lock-legal-hold, so a hold left in the map reaches the + // remote as user metadata and is silently dropped (#166). Carry it on + // the typed option instead. Retention stays in the map: the typed + // RetainUntilDate would truncate the date to whole seconds. legalHoldKey := strings.ToLower(xhttp.AmzObjectLockLegalHold) - forwardedLegalHold := srcInfo.UserDefined[legalHoldKey] - forwardedMeta := srcInfo.UserDefined - if forwardedLegalHold != "" { - forwardedMeta = cloneMSS(srcInfo.UserDefined) - delete(forwardedMeta, legalHoldKey) - } + forwardedLegalHold := forwardedMeta[legalHoldKey] + delete(forwardedMeta, legalHoldKey) opts := miniogo.PutObjectOptions{ UserMetadata: forwardedMeta, ServerSideEncryption: dstOpts.ServerSideEncryption, UserTags: tag.ToMap(), - } - if forwardedLegalHold != "" { - opts.LegalHold = miniogo.LegalHoldStatus(forwardedLegalHold) + LegalHold: miniogo.LegalHoldStatus(forwardedLegalHold), } // The destination must carry the same checksum the local path would // produce; the federated path has the remote compute, validate, persist @@ -2021,14 +2015,18 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re writeErrorResponse(ctx, w, toAPIError(ctx, rerr), r.URL) return } - // Built from the resolved values rather than the forwarding map: the - // legal hold was moved onto the typed option above, so opts.UserMetadata - // no longer carries it and the response and event would under-report. + // Keep the resolved legal hold in response and event metadata; the + // request-only checksum header exists only in the forwarding map. objInfo.UserDefined = cloneMSS(srcInfo.UserDefined) - // A forwarded checksum header is a request detail, not object metadata. - if checksumHeaderValue != "" { - delete(objInfo.UserDefined, wantChecksumType.Key()) - } + // The response headers and the ObjectCreated:Copy event describe the + // object this handler wrote, so name it: without these the response + // carries no x-amz-version-id and the event has an empty key and zero + // size (#170). Size is the logical size the remote was handed, which is + // what it reports back. + objInfo.Bucket = dstBucket + objInfo.Name = dstObject + objInfo.Size = actualSize + objInfo.VersionID = remoteObjInfo.VersionID objInfo.ETag = remoteObjInfo.ETag objInfo.ModTime = remoteObjInfo.LastModified // Bind the checksum the remote computed for this exact write. A single diff --git a/cmd/object-lock-default-retention_test.go b/cmd/object-lock-default-retention_test.go new file mode 100644 index 000000000..325657fc6 --- /dev/null +++ b/cmd/object-lock-default-retention_test.go @@ -0,0 +1,140 @@ +// Copyright (c) 2026 PGSTY +// SPDX-License-Identifier: AGPL-3.0-only + +package cmd + +import ( + "fmt" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/minio/minio/internal/auth" + objectlock "github.com/minio/minio/internal/bucket/object/lock" + xhttp "github.com/minio/minio/internal/http" +) + +func setTestBucketDefaultRetention(t *testing.T, bucket, mode string) { + t.Helper() + enableBucketObjectLock(t, bucket) + if mode == "" { + return + } + meta, err := globalBucketMetadataSys.Get(bucket) + if err != nil { + t.Fatal(err) + } + meta.ObjectLockConfigXML = fmt.Appendf(nil, `Enabled%s2`, mode) + if err := meta.parseAllConfigs(t.Context(), newObjectLayerFn()); err != nil { + t.Fatal(err) + } + globalBucketMetadataSys.Set(bucket, meta) +} + +func TestAPIObjectLockDefaultRetentionWithLegalHold(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + endpoints: []string{"NewMultipart", "PutObjectPart", "CompleteMultipart", "CopyObject", "PutObject", "DeleteObject"}, + objAPITest: func(obj ObjectLayer, instanceType, bucket string, router http.Handler, cred auth.Credentials, t *testing.T) { + data := []byte("default retention and legal hold are independent") + putCopyChecksumSource(t, router, cred, bucket, "source", data, nil) + for _, mode := range []string{"GOVERNANCE", "COMPLIANCE", ""} { + setTestBucketDefaultRetention(t, bucket, mode) + for _, hold := range []string{"ON", "OFF"} { + for _, operation := range []string{"put", "copy", "multipart"} { + t.Run(instanceType+"/"+mode+"/"+hold+"/"+operation, func(t *testing.T) { + object := mode + "-" + hold + "-" + operation + headers := map[string]string{xhttp.AmzObjectLockLegalHold: hold} + // The stored retention date carries millisecond precision. + before := UTCNow().Truncate(time.Millisecond) + switch operation { + case "put": + putCopyChecksumSource(t, router, cred, bucket, object, data, headers) + case "copy": + rec := federatedCopyRequest(t, router, cred, bucket, "source", bucket, object, headers) + if rec.Code != http.StatusOK { + t.Fatalf("copy: %d %s", rec.Code, rec.Body.String()) + } + case "multipart": + putFederationMultipartSource(t, router, cred, bucket, object, [][]byte{data}, headers) + } + info, err := obj.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{}) + if err != nil { + t.Fatal(err) + } + ret := objectlock.GetObjectRetentionMeta(info.UserDefined) + if string(ret.Mode) != mode { + t.Errorf("stored retention mode = %q, want %q", ret.Mode, mode) + } + if mode != "" && (ret.RetainUntilDate.Before(before.Add(48*time.Hour)) || ret.RetainUntilDate.After(UTCNow().Add(48*time.Hour))) { + t.Errorf("default retention date = %s, want write time + 2 days", ret.RetainUntilDate) + } + if got := objectlock.GetObjectLegalHoldMeta(info.UserDefined).Status; string(got) != hold { + t.Errorf("stored legal hold = %q, want %q", got, hold) + } + if mode == "COMPLIANCE" && hold == "OFF" { + req, err := newTestSignedRequestV4(http.MethodDelete, getPutObjectURL("", bucket, object)+"?versionId="+info.VersionID, 0, nil, cred.AccessKey, cred.SecretKey, nil) + if err != nil { + t.Fatal(err) + } + rec := httptest.NewRecorder() + router.ServeHTTP(rec, req) + if rec.Code != http.StatusBadRequest || !strings.Contains(rec.Body.String(), "InvalidRequest") { + t.Errorf("version DELETE must enforce retention: %d %s", rec.Code, rec.Body.String()) + } + } + }) + } + } + } + }, + }) +} + +func TestObjectLockDefaultRetentionBoundaries(t *testing.T) { + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + objAPITest: func(obj ObjectLayer, instanceType, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) { + setTestBucketDefaultRetention(t, bucket, "COMPLIANCE") + for _, tc := range []struct { + name string + replica, marker, explicit bool + retentionErr, holdErr APIErrorCode + wantMode objectlock.RetMode + wantErr APIErrorCode + }{ + {name: "trusted replica", replica: true}, + {name: "marker-only ordinary write", marker: true, wantMode: objectlock.RetCompliance}, + {name: "explicit retention", explicit: true, wantMode: objectlock.RetGovernance}, + {name: "retention permission denied", retentionErr: ErrAccessDenied, wantErr: ErrAccessDenied}, + {name: "legal hold permission denied", holdErr: ErrAccessDenied, wantErr: ErrAccessDenied}, + } { + t.Run(instanceType+"/"+tc.name, func(t *testing.T) { + r := httptest.NewRequest(http.MethodPut, "http://minio.local/"+bucket+"/object", nil) + r.Header.Set(xhttp.AmzObjectLockLegalHold, "OFF") + if tc.marker { + r.Header.Set(xhttp.MinIOSourceReplicationRequest, "true") + } + until := UTCNow().Add(24 * time.Hour).Truncate(time.Second).Add(789 * time.Millisecond) + if tc.explicit { + r.Header.Set(xhttp.AmzObjectLockMode, "GOVERNANCE") + r.Header.Set(xhttp.AmzObjectLockRetainUntilDate, until.Format(time.RFC3339Nano)) + } + mode, date, _, code := checkPutObjectLockAllowed(t.Context(), r, bucket, "object", obj.GetObjectInfo, tc.retentionErr, tc.holdErr, tc.replica) + if code != tc.wantErr || mode != tc.wantMode { + t.Fatalf("got mode %s error %s, want mode %s error %s", mode, niceError(code), tc.wantMode, niceError(tc.wantErr)) + } + if tc.explicit && !date.Equal(until) { + t.Errorf("explicit retention lost precision: %s != %s", date, until) + } + if tc.replica && !date.IsZero() { + t.Errorf("replica acquired a destination default: %s", date) + } + }) + } + }, + }) +} From c4b5e1cb45265220604a4dea16aeabb3d8c06c40 Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Fri, 11 Sep 2026 16:24:42 +0800 Subject: [PATCH 4/6] fix(auth): enforce header-only presigned payload checksums Signed-off-by: Feng Ruohang --- cmd/auth-handler.go | 15 ++--- cmd/presigned-payload_test.go | 114 ++++++++++++++++++++++++++++++++++ 2 files changed, 118 insertions(+), 11 deletions(-) create mode 100644 cmd/presigned-payload_test.go diff --git a/cmd/auth-handler.go b/cmd/auth-handler.go index 884fa9c83..fa35d9058 100644 --- a/cmd/auth-handler.go +++ b/cmd/auth-handler.go @@ -632,18 +632,11 @@ func isReqAuthenticated(ctx context.Context, r *http.Request, region string, sty return ErrInvalidDigest } - // Extract either 'X-Amz-Content-Sha256' header or 'X-Amz-Content-Sha256' query parameter (if V4 presigned) - // Do not verify 'X-Amz-Content-Sha256' if skipSHA256. + // Honor the selected header/query checksum, including the header fallback + // for a presigned request. STS separately hashes its body for its signature. var contentSHA256 []byte - if skipSHA256 := skipContentSha256Cksum(r); !skipSHA256 && isRequestPresignedSignatureV4(r) { - if sha256Sum, ok := r.Form[xhttp.AmzContentSha256]; ok && len(sha256Sum) > 0 { - contentSHA256, err = hex.DecodeString(sha256Sum[0]) - if err != nil { - return ErrContentSHA256Mismatch - } - } - } else if _, ok := r.Header[xhttp.AmzContentSha256]; !skipSHA256 && ok { - contentSHA256, err = hex.DecodeString(r.Header.Get(xhttp.AmzContentSha256)) + if !skipContentSha256Cksum(r) { + contentSHA256, err = hex.DecodeString(getContentSha256Cksum(r, serviceS3)) if err != nil || len(contentSHA256) == 0 { return ErrContentSHA256Mismatch } diff --git a/cmd/presigned-payload_test.go b/cmd/presigned-payload_test.go new file mode 100644 index 000000000..ef39a0a84 --- /dev/null +++ b/cmd/presigned-payload_test.go @@ -0,0 +1,114 @@ +// Copyright (c) 2026 PGSTY +// SPDX-License-Identifier: AGPL-3.0-only + +package cmd + +import ( + "fmt" + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/minio/minio/internal/auth" + xhttp "github.com/minio/minio/internal/http" +) + +func TestPresignedHeaderPayloadHashIsEnforced(t *testing.T) { + setupSignatureBoundaryTest(t) + digest := getSHA256Hash([]byte("expected")) + for _, tc := range []struct { + name, query, header string + unsigned bool + }{ + {name: "header only", header: digest}, + {name: "query only", query: digest}, + {name: "query wins", query: digest, header: getSHA256Hash([]byte("different"))}, + {name: "unsigned query wins", query: unsignedPayload, header: digest, unsigned: true}, + {name: "unsigned header", header: unsignedPayload, unsigned: true}, + {name: "absent", unsigned: true}, + } { + for _, body := range []string{"expected", "modified"} { + t.Run(tc.name+"/"+body, func(t *testing.T) { + r := httptest.NewRequest(http.MethodPut, "http://minio.local/bucket/object", strings.NewReader(body)) + if tc.header != "" { + r.Header.Set(xhttp.AmzContentSha256, tc.header) + } + if tc.query != "" { + query := r.URL.Query() + query.Set(xhttp.AmzContentSha256, tc.query) + r.URL.RawQuery = query.Encode() + } + presignBoundaryRequest(t, r, UTCNow(), []string{"host"}, globalActiveCred) + if code := isReqAuthenticated(t.Context(), r, globalSite.Region(), serviceS3); code != ErrNone { + t.Fatal(niceError(code)) + } + _, err := io.ReadAll(r.Body) + wantMismatch := body == "modified" && !tc.unsigned + if wantMismatch { + if toAPIErrorCode(t.Context(), err) != ErrContentSHA256Mismatch { + t.Errorf("modified signed payload: got %v, want SHA-256 mismatch", err) + } + } else if err != nil { + t.Errorf("valid payload: %v", err) + } + }) + } + } +} + +func TestAPIPresignedBucketPolicyPayloadHash(t *testing.T) { + defer DetectTestLeak(t)() + ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{ + t: t, + objAPITest: func(objectAPI ObjectLayer, instanceType, bucket string, apiRouter http.Handler, credentials auth.Credentials, t *testing.T) { + server := httptest.NewServer(apiRouter) + defer server.Close() + client := server.Client() + client.Timeout = 10 * time.Second + expected := fmt.Sprintf(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":{"AWS":["*"]},"Action":["s3:GetObject"],"Resource":["arn:aws:s3:::%s/*"]}]}`, bucket) + modified := strings.Replace(expected, "s3:GetObject", "s3:PutObject", 1) + for _, tc := range []struct { + name, body string + status int + }{ + {"original", expected, http.StatusNoContent}, + {"modified", modified, http.StatusBadRequest}, + } { + t.Run(instanceType+"/"+tc.name, func(t *testing.T) { + r, err := http.NewRequestWithContext(t.Context(), http.MethodPut, server.URL+"/"+bucket+"?policy=", strings.NewReader(tc.body)) + if err != nil { + t.Fatal(err) + } + r.Header.Set(xhttp.AmzContentSha256, getSHA256Hash([]byte(expected))) + presignBoundaryRequest(t, r, UTCNow(), []string{"host"}, credentials) + response, err := client.Do(r) + if err != nil { + t.Fatal(err) + } + defer response.Body.Close() + body, err := io.ReadAll(response.Body) + if err != nil { + t.Fatal(err) + } + if response.StatusCode != tc.status { + t.Fatalf("PutBucketPolicy: got %d %s, want %d", response.StatusCode, body, tc.status) + } + if tc.status == http.StatusBadRequest && !strings.Contains(string(body), "XAmzContentSHA256Mismatch") { + t.Fatalf("expected payload checksum rejection, got %s", body) + } + }) + } + meta, err := globalBucketMetadataSys.Get(bucket) + if err != nil { + t.Fatal(err) + } + if strings.Contains(string(meta.PolicyConfigJSON), "s3:PutObject") { + t.Fatal("tampered presigned request replaced the bucket policy") + } + }, + endpoints: []string{"PutBucketPolicy"}, + }) +} From 68127c5a632bf1bcfd4c5de9ac57f78c1a46af3a Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Fri, 11 Sep 2026 16:40:42 +0800 Subject: [PATCH 5/6] build(deps): embed the consolidated Console source Signed-off-by: Feng Ruohang --- go.mod | 2 +- go.sum | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/go.mod b/go.mod index a9ff1166e..c43a937ca 100644 --- a/go.mod +++ b/go.mod @@ -4,7 +4,7 @@ go 1.27.1 // Console and MC retain their historical module paths for best-effort upstream // compatibility. Pin the maintained PGSTY implementations used by SILO. -replace github.com/minio/console => github.com/pgsty/silo-console v0.0.0-20260910093545-6a0b31b5ade2 +replace github.com/minio/console => github.com/pgsty/silo-console v0.0.0-20260911081814-153f48fb5f80 replace github.com/minio/mc => github.com/pgsty/mc v0.0.0-20260910093317-e6a60edf0952 diff --git a/go.sum b/go.sum index d94781d25..dc9ec046b 100644 --- a/go.sum +++ b/go.sum @@ -547,8 +547,8 @@ github.com/pborman/getopt v0.0.0-20170112200414-7148bc3a4c30/go.mod h1:85jBQOZwp github.com/pelletier/go-toml v1.2.0/go.mod h1:5z9KED0ma1S8pY6P1sdut58dfprrGBbd/94hg7ilaic= github.com/pgsty/mc v0.0.0-20260910093317-e6a60edf0952 h1:b6jpBuZWUiZhBox64eF1cmNL4HWGwdSk/LVXNc3DBQ0= github.com/pgsty/mc v0.0.0-20260910093317-e6a60edf0952/go.mod h1:VsGNEmditljwBgmYKy7X+//dXLN8gEfSPR2m4lcMPmo= -github.com/pgsty/silo-console v0.0.0-20260910093545-6a0b31b5ade2 h1:v/AXKa/UZkDheLLGrpKmyu8bVHo11dVrDHdbDvjX1/M= -github.com/pgsty/silo-console v0.0.0-20260910093545-6a0b31b5ade2/go.mod h1:bCyOnZactQajLOqlgCn0lHwQWMNylvP6MUOqqLako00= +github.com/pgsty/silo-console v0.0.0-20260911081814-153f48fb5f80 h1:bEpCn6B+aEOLriH9EjJwsGY9tarR5P/9Npv54S7DMJQ= +github.com/pgsty/silo-console v0.0.0-20260911081814-153f48fb5f80/go.mod h1:bCyOnZactQajLOqlgCn0lHwQWMNylvP6MUOqqLako00= github.com/pgsty/silo-pkg/v3 v3.13.4-0.20260910091716-2d8fd3cbbf07 h1:IKm2AyPsvuL4NyniK3ixqqQz1CIargmhyCkL2eZvlhA= github.com/pgsty/silo-pkg/v3 v3.13.4-0.20260910091716-2d8fd3cbbf07/go.mod h1:1JdcUcM+TRXObb9QrC8FWhav0RKk6/u8RT2RgimfxRQ= github.com/philhofer/fwd v1.2.0 h1:e6DnBTl7vGY+Gz322/ASL4Gyp1FspeMvx1RNDoToZuM= From d63c92e3930bbcd2a392313852483db3ed657fc6 Mon Sep 17 00:00:00 2001 From: Feng Ruohang Date: Fri, 11 Sep 2026 16:56:33 +0800 Subject: [PATCH 6/6] build(deps): secure AMQP frames and select merged Console Signed-off-by: Feng Ruohang --- go.mod | 4 +- go.sum | 8 ++-- internal/event/target/amqp_test.go | 77 ++++++++++++++++++++++++++++++ 3 files changed, 83 insertions(+), 6 deletions(-) create mode 100644 internal/event/target/amqp_test.go diff --git a/go.mod b/go.mod index c43a937ca..9f3618526 100644 --- a/go.mod +++ b/go.mod @@ -4,7 +4,7 @@ go 1.27.1 // Console and MC retain their historical module paths for best-effort upstream // compatibility. Pin the maintained PGSTY implementations used by SILO. -replace github.com/minio/console => github.com/pgsty/silo-console v0.0.0-20260911081814-153f48fb5f80 +replace github.com/minio/console => github.com/pgsty/silo-console v0.0.0-20260911085047-638eefd7aece replace github.com/minio/mc => github.com/pgsty/mc v0.0.0-20260910093317-e6a60edf0952 @@ -92,7 +92,7 @@ require ( github.com/prometheus/common v0.71.0 github.com/prometheus/procfs v0.22.0 github.com/puzpuzpuz/xsync/v3 v3.5.1 - github.com/rabbitmq/amqp091-go v1.10.0 + github.com/rabbitmq/amqp091-go v1.14.0 github.com/rcrowley/go-metrics v0.0.0-20250401214520-65e299d6c5c9 github.com/rs/cors v1.11.1 github.com/secure-io/sio-go v0.3.1 diff --git a/go.sum b/go.sum index dc9ec046b..bbf0396f1 100644 --- a/go.sum +++ b/go.sum @@ -547,8 +547,8 @@ github.com/pborman/getopt v0.0.0-20170112200414-7148bc3a4c30/go.mod h1:85jBQOZwp github.com/pelletier/go-toml v1.2.0/go.mod h1:5z9KED0ma1S8pY6P1sdut58dfprrGBbd/94hg7ilaic= github.com/pgsty/mc v0.0.0-20260910093317-e6a60edf0952 h1:b6jpBuZWUiZhBox64eF1cmNL4HWGwdSk/LVXNc3DBQ0= github.com/pgsty/mc v0.0.0-20260910093317-e6a60edf0952/go.mod h1:VsGNEmditljwBgmYKy7X+//dXLN8gEfSPR2m4lcMPmo= -github.com/pgsty/silo-console v0.0.0-20260911081814-153f48fb5f80 h1:bEpCn6B+aEOLriH9EjJwsGY9tarR5P/9Npv54S7DMJQ= -github.com/pgsty/silo-console v0.0.0-20260911081814-153f48fb5f80/go.mod h1:bCyOnZactQajLOqlgCn0lHwQWMNylvP6MUOqqLako00= +github.com/pgsty/silo-console v0.0.0-20260911085047-638eefd7aece h1:1QGCeVCRV4TzDLww9IFfRT70YOI8r/X88pRo+KILvZ4= +github.com/pgsty/silo-console v0.0.0-20260911085047-638eefd7aece/go.mod h1:bCyOnZactQajLOqlgCn0lHwQWMNylvP6MUOqqLako00= github.com/pgsty/silo-pkg/v3 v3.13.4-0.20260910091716-2d8fd3cbbf07 h1:IKm2AyPsvuL4NyniK3ixqqQz1CIargmhyCkL2eZvlhA= github.com/pgsty/silo-pkg/v3 v3.13.4-0.20260910091716-2d8fd3cbbf07/go.mod h1:1JdcUcM+TRXObb9QrC8FWhav0RKk6/u8RT2RgimfxRQ= github.com/philhofer/fwd v1.2.0 h1:e6DnBTl7vGY+Gz322/ASL4Gyp1FspeMvx1RNDoToZuM= @@ -590,8 +590,8 @@ github.com/prometheus/prometheus v0.314.0 h1:YjsimqsIi6/mOtzZcrPEYUALO6zpfaht9O5 github.com/prometheus/prometheus v0.314.0/go.mod h1:zjg3pMTAkY0/JG8jy/h8/YgSQUVB+aCXMhUqN6l64jg= github.com/puzpuzpuz/xsync/v3 v3.5.1 h1:GJYJZwO6IdxN/IKbneznS6yPkVC+c3zyY/j19c++5Fg= github.com/puzpuzpuz/xsync/v3 v3.5.1/go.mod h1:VjzYrABPabuM4KyBh1Ftq6u8nhwY5tBPKP9jpmh0nnA= -github.com/rabbitmq/amqp091-go v1.10.0 h1:STpn5XsHlHGcecLmMFCtg7mqq0RnD+zFr4uzukfVhBw= -github.com/rabbitmq/amqp091-go v1.10.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= +github.com/rabbitmq/amqp091-go v1.14.0 h1:RSaT7aOKt/OrkVUyswPDW29lnRz9psuGmfZFBmLqLek= +github.com/rabbitmq/amqp091-go v1.14.0/go.mod h1:Hy4jKW5kQART1u+JkDTF9YYOQUHXqMuhrgxOEeS7G4o= github.com/rcrowley/go-metrics v0.0.0-20250401214520-65e299d6c5c9 h1:bsUq1dX0N8AOIL7EB/X911+m4EHsnWEHeJ0c+3TTBrg= github.com/rcrowley/go-metrics v0.0.0-20250401214520-65e299d6c5c9/go.mod h1:bCqnVzQkZxMG4s8nGwiZ5l3QUCyqpo9Y+/ZMZ9VjZe4= github.com/rivo/uniseg v0.1.0/go.mod h1:J6wj4VEh+S6ZtnVlnTBMWIodfgj8LQOQFoIToxlJtxc= diff --git a/internal/event/target/amqp_test.go b/internal/event/target/amqp_test.go new file mode 100644 index 000000000..596de201f --- /dev/null +++ b/internal/event/target/amqp_test.go @@ -0,0 +1,77 @@ +// Copyright (c) 2026 PGSTY +// SPDX-License-Identifier: AGPL-3.0-only + +package target + +import ( + "context" + "errors" + "fmt" + "io" + "net" + "testing" + "time" + + "github.com/rabbitmq/amqp091-go" +) + +// A broker must not make the notification client allocate or read an oversized +// frame, including before connection.tune negotiates the frame size limit. +func TestAMQPRejectsOversizedHandshakeFrame(t *testing.T) { + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { listener.Close() }) + + uri, err := amqp091.ParseURI("amqp://" + listener.Addr().String()) + if err != nil { + t.Fatal(err) + } + target, err := NewAMQPTarget("oversized-frame", AMQPArgs{Enable: true, URL: uri}, + func(context.Context, error, string, ...any) {}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { target.Close() }) + + brokerDone := make(chan error, 1) + go func() { + brokerDone <- func() error { + conn, err := listener.Accept() + if err != nil { + return err + } + defer conn.Close() + if err := conn.SetDeadline(time.Now().Add(5 * time.Second)); err != nil { + return err + } + var protocol [8]byte + if _, err := io.ReadFull(conn, protocol[:]); err != nil { + return err + } + if protocol != [8]byte{'A', 'M', 'Q', 'P', 0, 0, 9, 1} { + return fmt.Errorf("unexpected AMQP protocol header: %x", protocol) + } + // A method frame declaring an 8 KiB payload exceeds the 4 KiB + // pre-negotiation limit. Send only its header: rejection must + // happen without allocating or waiting for the declared body. + if _, err := conn.Write([]byte{1, 0, 0, 0, 0, 0x20, 0}); err != nil { + return err + } + var response [1]byte + if n, err := conn.Read(response[:]); n != 0 || !errors.Is(err, io.EOF) { + return fmt.Errorf("client did not close after oversized frame header: read %d bytes, error %v", n, err) + } + return nil + }() + }() + + active, err := target.IsActive() + if active || err == nil { + t.Errorf("oversized handshake frame accepted: active=%v, error=%v", active, err) + } + if err := <-brokerDone; err != nil { + t.Fatal(err) + } +}