Merge pull request #101 from pgsty/codex/fix-cors-replication-trust

fix: harden CORS and replication request trust
This commit is contained in:
Feng Ruohang
2026-09-02 07:07:49 +08:00
committed by GitHub
22 changed files with 1783 additions and 196 deletions
@@ -693,6 +693,7 @@
"/%s/us-east-1/s3/aws4_request", "/%s/us-east-1/s3/aws4_request",
"/*", "/*",
"/../../etc", "/../../etc",
"/../obj",
"/./abc/def", "/./abc/def",
"/.dockerenv", "/.dockerenv",
"/.trash", "/.trash",
@@ -705,6 +706,7 @@
"//contains/double-forwardslash-prefix", "//contains/double-forwardslash-prefix",
"/?", "/?",
"/?X-Amz-Algorithm=AWS4-HMAC-SHA256\u0026X-Amz-Credential=USWUXHGYZQYFYFFIT3RE%2F20170529%2Fus-east-1%2Fs3%2Faws4_request\u0026X-Amz-Date=20170529T190139Z\u0026X-Amz-Expires=600\u0026X-Amz-Signature=19b58080999df54b446fc97304eb8dda60d3df1812ae97f3e8783351bfd9781d\u0026X-Amz-SignedHeaders=host\u0026prefix=Hello%2AWorld%2A", "/?X-Amz-Algorithm=AWS4-HMAC-SHA256\u0026X-Amz-Credential=USWUXHGYZQYFYFFIT3RE%2F20170529%2Fus-east-1%2Fs3%2Faws4_request\u0026X-Amz-Date=20170529T190139Z\u0026X-Amz-Expires=600\u0026X-Amz-Signature=19b58080999df54b446fc97304eb8dda60d3df1812ae97f3e8783351bfd9781d\u0026X-Amz-SignedHeaders=host\u0026prefix=Hello%2AWorld%2A",
"/A/obj",
"/a", "/a",
"/a/b/c", "/a/b/c",
"/a/b/c/d/e/f/g", "/a/b/c/d/e/f/g",
@@ -721,6 +723,7 @@
"/admin", "/admin",
"/afile", "/afile",
"/api/requests", "/api/requests",
"/api/v1/login",
"/apis", "/apis",
"/audit", "/audit",
"/background-heal/status", "/background-heal/status",
@@ -852,6 +855,7 @@
"/ls", "/ls",
"/metrics", "/metrics",
"/metrics/v3", "/metrics/v3",
"/minio/admin/v3/info",
"/minio/grid/", "/minio/grid/",
"/minio/grid/lock/", "/minio/grid/lock/",
"/minio/health/cluster", "/minio/health/cluster",
@@ -970,6 +974,7 @@
"/speedtest/site", "/speedtest/site",
"/start-job", "/start-job",
"/startprofiling", "/startprofiling",
"/startup-missing/object",
"/status", "/status",
"/status-job", "/status-job",
"/storage", "/storage",
@@ -1020,6 +1025,7 @@
"/version", "/version",
"/vfile", "/vfile",
"/wall", "/wall",
"/x/obj",
"/xl.meta", "/xl.meta",
"/{bucket}", "/{bucket}",
"/{object:.+}" "/{object:.+}"
@@ -3139,6 +3145,7 @@
"cmd:cmd:method:BucketMetadataSys.GetPolicyConfig", "cmd:cmd:method:BucketMetadataSys.GetPolicyConfig",
"cmd:cmd:method:BucketMetadataSys.GetQuotaConfig", "cmd:cmd:method:BucketMetadataSys.GetQuotaConfig",
"cmd:cmd:method:BucketMetadataSys.GetReplicationConfig", "cmd:cmd:method:BucketMetadataSys.GetReplicationConfig",
"cmd:cmd:method:BucketMetadataSys.GetResidentCorsConfig",
"cmd:cmd:method:BucketMetadataSys.GetSSEConfig", "cmd:cmd:method:BucketMetadataSys.GetSSEConfig",
"cmd:cmd:method:BucketMetadataSys.GetTaggingConfig", "cmd:cmd:method:BucketMetadataSys.GetTaggingConfig",
"cmd:cmd:method:BucketMetadataSys.GetVersioningConfig", "cmd:cmd:method:BucketMetadataSys.GetVersioningConfig",
@@ -10258,6 +10265,10 @@
"cmd/object-multipart-handlers.go=\"X-Minio-Replication-Server-Side-Encryption-Iv\"", "cmd/object-multipart-handlers.go=\"X-Minio-Replication-Server-Side-Encryption-Iv\"",
"cmd/object-multipart-handlers.go=\"X-Minio-Replication-Server-Side-Encryption-Seal-Algorithm\"", "cmd/object-multipart-handlers.go=\"X-Minio-Replication-Server-Side-Encryption-Seal-Algorithm\"",
"cmd/object-multipart-handlers.go=\"X-Minio-Replication-Server-Side-Encryption-Sealed-Key\"", "cmd/object-multipart-handlers.go=\"X-Minio-Replication-Server-Side-Encryption-Sealed-Key\"",
"cmd/replication-trust.go=\"X-Minio-Replication-Encrypted-Multipart\"",
"cmd/replication-trust.go=\"X-Minio-Replication-Server-Side-Encryption-Iv\"",
"cmd/replication-trust.go=\"X-Minio-Replication-Server-Side-Encryption-Seal-Algorithm\"",
"cmd/replication-trust.go=\"X-Minio-Replication-Server-Side-Encryption-Sealed-Key\"",
"cmd/s3-zip-handlers.go=\"x-minio-extract\"", "cmd/s3-zip-handlers.go=\"x-minio-extract\"",
"cmd/server-startup-msg.go=\"https://silo.pgsty.com/reference/minio-mc/#quickstart\"", "cmd/server-startup-msg.go=\"https://silo.pgsty.com/reference/minio-mc/#quickstart\"",
"cmd/storage-rest-server.go=\"X-Minio-Time\"", "cmd/storage-rest-server.go=\"X-Minio-Time\"",
+10 -1
View File
@@ -787,7 +787,16 @@ func corsHandler(handler http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Header.Get("Origin") != "" { if r.Header.Get("Origin") != "" {
if bucket, _ := request2BucketObjectName(r); bucket != "" && globalBucketMetadataSys != nil { if bucket, _ := request2BucketObjectName(r); bucket != "" && globalBucketMetadataSys != nil {
cfg, _, err := globalBucketMetadataSys.GetCorsConfig(bucket) // Resident-only lookup: this runs pre-auth for every
// Origin-bearing request using a client-supplied path segment as
// the bucket name. It must never load or cache metadata for
// arbitrary names (see GetResidentCorsConfig). GetResidentCorsConfig
// is the single decision point: it returns errInvalidArgument for
// the internal .minio.sys namespace (fail closed), a config for a
// resident bucket, errBucketMetadataNotInitialized for a real but
// unloaded bucket (fail closed), and errConfigNotFound otherwise
// (fall back to the global policy below).
cfg, _, err := globalBucketMetadataSys.GetResidentCorsConfig(bucket)
if err == nil && cfg != nil { if err == nil && cfg != nil {
if applyBucketCors(w, r, cfg) { if applyBucketCors(w, r, cfg) {
return return
+13 -3
View File
@@ -786,10 +786,20 @@ func isPutActionAllowedWithRequestTags(ctx context.Context, atype authType, buck
return s3Err return s3Err
} }
logger.GetReqInfo(ctx).Cred = cred reqInfo := logger.GetReqInfo(ctx)
logger.GetReqInfo(ctx).Owner = owner if reqInfo == nil {
logger.GetReqInfo(ctx).Region = region return ErrAccessDenied
}
reqInfo.Lock()
reqInfo.Cred = cred
reqInfo.Owner = owner
reqInfo.Region = region
reqInfo.Unlock()
return isPutActionAllowedWithCred(bucketName, objectName, r, action, requestTags, cred, owner)
}
func isPutActionAllowedWithCred(bucketName, objectName string, r *http.Request, action policy.Action, requestTags *string, cred auth.Credentials, owner bool) APIErrorCode {
// Do not check for PutObjectRetentionAction permission, // Do not check for PutObjectRetentionAction permission,
// if mode and retain until date are not set. // if mode and retain until date are not set.
// Can happen when bucket has default lock config set // Can happen when bucket has default lock config set
+268 -1
View File
@@ -19,6 +19,7 @@ package cmd
import ( import (
"context" "context"
"fmt"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
"strings" "strings"
@@ -355,7 +356,20 @@ func TestBucketCorsMissingBucketUsesGlobalFallback(t *testing.T) {
}) })
} }
func testBucketCorsMissingBucketUsesGlobalFallback(_ ObjectLayer, _ string, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) { func testBucketCorsMissingBucketUsesGlobalFallback(obj ObjectLayer, _ string, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) {
// Model a fully started server: bucket metadata loading has completed, so
// a name that is not resident is genuinely not a CORS-bearing bucket.
restore := markBucketMetadataInitialized(t)
defer restore()
// A non-resident bucket name must not cause any bucket-metadata disk read.
oldObjectAPI := newObjectLayerFn()
counting := &corsLookupCountingObjectLayer{ObjectLayer: obj}
setObjectLayer(counting)
defer setObjectLayer(oldObjectAPI)
before := bucketMetadataMapLen()
wrapped := corsHandler(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { wrapped := corsHandler(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNotFound) w.WriteHeader(http.StatusNotFound)
})) }))
@@ -373,6 +387,14 @@ func testBucketCorsMissingBucketUsesGlobalFallback(_ ObjectLayer, _ string, buck
if got := rec.Header().Get("Access-Control-Allow-Credentials"); got != "true" { if got := rec.Header().Get("Access-Control-Allow-Credentials"); got != "true" {
t.Fatalf("allow-credentials = %q", got) t.Fatalf("allow-credentials = %q", got)
} }
// Regression guard: the pre-auth CORS lookup for a non-existent bucket must
// neither read bucket metadata from disk nor cache a synthetic entry.
if got := counting.getObjectNInfoCalls.Load(); got != 0 {
t.Fatalf("missing-bucket CORS lookup performed %d bucket metadata reads", got)
}
if after := bucketMetadataMapLen(); after != before {
t.Fatalf("missing-bucket CORS lookup grew metadataMap from %d to %d", before, after)
}
} }
func testBucketCorsNoConfigUsesGlobalFallback(_ ObjectLayer, _ string, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) { func testBucketCorsNoConfigUsesGlobalFallback(_ ObjectLayer, _ string, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) {
@@ -471,3 +493,248 @@ func requireCorsOriginVary(t *testing.T, header http.Header) {
t.Fatalf("Vary = %q, missing Origin", values) t.Fatalf("Vary = %q, missing Origin", values)
} }
} }
// markBucketMetadataInitialized marks the global bucket-metadata subsystem as
// fully loaded, modeling a running server (the API test harness sets up the
// subsystem but does not run Init). It returns a function that restores the
// previous state.
func markBucketMetadataInitialized(t *testing.T) func() {
t.Helper()
sys := globalBucketMetadataSys
if sys == nil {
t.Fatal("globalBucketMetadataSys is nil")
}
sys.Lock()
prev := sys.initialized
sys.initialized = true
sys.Unlock()
return func() {
sys.Lock()
sys.initialized = prev
sys.Unlock()
}
}
// bucketMetadataMapLen returns the number of resident bucket-metadata entries.
func bucketMetadataMapLen() int {
sys := globalBucketMetadataSys
if sys == nil {
return 0
}
sys.RLock()
defer sys.RUnlock()
return len(sys.metadataMap)
}
// TestBucketCorsUnknownBucketDoesNotGrowMetadata is the regression guard for
// the pre-auth resource-exhaustion path: an unauthenticated, Origin-bearing
// request whose first path segment is not a real bucket must fall back to the
// global CORS policy without loading bucket metadata from disk and without
// caching a synthetic entry. Before the resident-only lookup, each distinct
// name grew metadataMap by one and issued an erasure metadata probe.
func TestBucketCorsUnknownBucketDoesNotGrowMetadata(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testBucketCorsUnknownBucketDoesNotGrowMetadata,
endpoints: []string{"GetBucketCors"},
})
}
func testBucketCorsUnknownBucketDoesNotGrowMetadata(obj ObjectLayer, _ string, _ string, _ http.Handler, _ auth.Credentials, t *testing.T) {
restore := markBucketMetadataInitialized(t)
defer restore()
oldObjectAPI := newObjectLayerFn()
counting := &corsLookupCountingObjectLayer{ObjectLayer: obj}
setObjectLayer(counting)
defer setObjectLayer(oldObjectAPI)
wrapped := corsHandler(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNotFound)
}))
before := bucketMetadataMapLen()
// Console/admin routes plus enough distinct valid names to make accidental
// cache growth or one metadata probe per name unambiguous.
names := []string{"minio", "api"}
for i := 0; i < 500; i++ {
names = append(names, fmt.Sprintf("cors-missing-%03d", i))
}
for _, name := range names {
for _, method := range []string{http.MethodGet, http.MethodOptions} {
rec := httptest.NewRecorder()
req := httptest.NewRequest(method, getGetObjectURL("", name, "obj"), nil)
req.Header.Set("Origin", "https://app.example.com")
if method == http.MethodOptions {
req.Header.Set("Access-Control-Request-Method", http.MethodGet)
}
wrapped.ServeHTTP(rec, req)
if got := rec.Header().Get("Access-Control-Allow-Origin"); got != "https://app.example.com" {
t.Fatalf("%s/%s: allow-origin = %q, want global fallback", name, method, got)
}
}
}
for _, path := range []string{"/../obj", "/A/obj", "/x/obj", "/minio/admin/v3/info", "/api/v1/login"} {
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, path, nil)
req.Header.Set("Origin", "https://app.example.com")
wrapped.ServeHTTP(rec, req)
if got := rec.Header().Get("Access-Control-Allow-Origin"); got != "https://app.example.com" {
t.Fatalf("%s: allow-origin = %q, want global fallback", path, got)
}
}
if got := counting.getObjectNInfoCalls.Load(); got != 0 {
t.Fatalf("unknown-bucket CORS lookups performed %d bucket metadata reads", got)
}
if after := bucketMetadataMapLen(); after != before {
t.Fatalf("unknown-bucket CORS lookups grew metadataMap from %d to %d", before, after)
}
}
func TestBucketCorsStartupMissFailsClosedWithoutIO(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testBucketCorsStartupMissFailsClosedWithoutIO,
endpoints: []string{"GetBucketCors"},
})
}
func testBucketCorsStartupMissFailsClosedWithoutIO(obj ObjectLayer, _ string, _ string, _ http.Handler, _ auth.Credentials, t *testing.T) {
oldObjectAPI := newObjectLayerFn()
oldMetadataSys := globalBucketMetadataSys
counting := &corsLookupCountingObjectLayer{ObjectLayer: obj}
setObjectLayer(counting)
globalBucketMetadataSys = NewBucketMetadataSys()
defer func() {
setObjectLayer(oldObjectAPI)
globalBucketMetadataSys = oldMetadataSys
}()
innerCalled := false
wrapped := corsHandler(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
innerCalled = true
w.WriteHeader(http.StatusNoContent)
}))
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, "/startup-missing/object", nil)
req.Header.Set("Origin", "https://app.example.com")
wrapped.ServeHTTP(rec, req)
if !innerCalled || rec.Code != http.StatusNoContent {
t.Fatalf("startup miss did not reach inner handler: called=%v status=%d", innerCalled, rec.Code)
}
if got := rec.Header().Get("Access-Control-Allow-Origin"); got != "" {
t.Fatalf("startup miss used permissive global CORS: %q", got)
}
if got := counting.getObjectNInfoCalls.Load(); got != 0 {
t.Fatalf("startup miss performed %d metadata reads", got)
}
if got := globalBucketMetadataSys.Count(); got != 0 {
t.Fatalf("startup miss grew metadataMap to %d", got)
}
}
// markBucketMetadataLoadFailed records a bucket as one whose metadata failed to
// load at startup while the subsystem is Initialized, modeling the degraded
// state where a real bucket is not resident. Returns a restore function.
func markBucketMetadataLoadFailed(t *testing.T, bucket string) func() {
t.Helper()
sys := globalBucketMetadataSys
if sys == nil {
t.Fatal("globalBucketMetadataSys is nil")
}
sys.Lock()
_, had := sys.loadFailed[bucket]
sys.loadFailed[bucket] = struct{}{}
sys.Unlock()
return func() {
sys.Lock()
if !had {
delete(sys.loadFailed, bucket)
}
sys.Unlock()
}
}
// TestBucketCorsLoadFailedBucketFailsClosed guards P1: a real bucket whose
// metadata could not be loaded at startup (present in loadFailed, subsystem
// Initialized) must NOT be answered with the permissive global CORS policy. We
// cannot rule out a restrictive per-bucket config for it, so it must fail
// closed — without a synchronous disk read.
func TestBucketCorsLoadFailedBucketFailsClosed(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testBucketCorsLoadFailedBucketFailsClosed,
endpoints: []string{"GetBucketCors"},
})
}
func testBucketCorsLoadFailedBucketFailsClosed(obj ObjectLayer, _ string, _ string, _ http.Handler, _ auth.Credentials, t *testing.T) {
restoreInit := markBucketMetadataInitialized(t)
defer restoreInit()
restoreFail := markBucketMetadataLoadFailed(t, "strict-cors-bucket")
defer restoreFail()
oldObjectAPI := newObjectLayerFn()
counting := &corsLookupCountingObjectLayer{ObjectLayer: obj}
setObjectLayer(counting)
defer setObjectLayer(oldObjectAPI)
wrapped := corsHandler(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNotFound)
}))
for _, method := range []string{http.MethodGet, http.MethodOptions} {
rec := httptest.NewRecorder()
req := httptest.NewRequest(method, getGetObjectURL("", "strict-cors-bucket", "object"), nil)
req.Header.Set("Origin", "https://app.example.com")
if method == http.MethodOptions {
req.Header.Set("Access-Control-Request-Method", http.MethodGet)
}
wrapped.ServeHTTP(rec, req)
if got := rec.Header().Get("Access-Control-Allow-Origin"); got != "" {
t.Fatalf("%s: load-failed bucket fell back to global allow-origin %q", method, got)
}
if got := rec.Header().Get("Access-Control-Allow-Credentials"); got != "" {
t.Fatalf("%s: load-failed bucket fell back to global credentials %q", method, got)
}
}
if got := counting.getObjectNInfoCalls.Load(); got != 0 {
t.Fatalf("load-failed CORS lookup performed %d synchronous bucket metadata reads", got)
}
}
// TestBucketCorsInternalBucketFailsClosed guards P2: an Origin-bearing request
// whose first path segment is the reserved .minio.sys namespace must preserve
// GetConfig's errInvalidArgument semantics and fail closed, not fall back to
// the permissive global CORS policy.
func TestBucketCorsInternalBucketFailsClosed(t *testing.T) {
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testBucketCorsInternalBucketFailsClosed,
endpoints: []string{"GetBucketCors"},
})
}
func testBucketCorsInternalBucketFailsClosed(_ ObjectLayer, _ string, _ string, _ http.Handler, _ auth.Credentials, t *testing.T) {
restoreInit := markBucketMetadataInitialized(t)
defer restoreInit()
wrapped := corsHandler(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusNotFound)
}))
rec := httptest.NewRecorder()
req := httptest.NewRequest(http.MethodGet, getGetObjectURL("", minioMetaBucket, "object"), nil)
req.Header.Set("Origin", "https://app.example.com")
wrapped.ServeHTTP(rec, req)
if got := rec.Header().Get("Access-Control-Allow-Origin"); got != "" {
t.Fatalf("internal bucket fell back to global allow-origin %q", got)
}
if got := rec.Header().Get("Access-Control-Allow-Credentials"); got != "" {
t.Fatalf("internal bucket fell back to global credentials %q", got)
}
}
+33 -5
View File
@@ -634,7 +634,7 @@ func testSiteReplicationStatusCountsCorsPerSite(obj ObjectLayer, _ string, bucke
globalSiteReplicationSys.Unlock() globalSiteReplicationSys.Unlock()
}() }()
check := func(name string, wantLocal, wantRemote int) { check := func(name string, wantLocal, wantRemote int, wantMismatch, wantReplicated bool) {
t.Helper() t.Helper()
status, err := globalSiteReplicationSys.siteReplicationStatus(ctx, obj, madmin.SRStatusOptions{Buckets: true}) status, err := globalSiteReplicationSys.siteReplicationStatus(ctx, obj, madmin.SRStatusOptions{Buckets: true})
if err != nil { if err != nil {
@@ -646,23 +646,51 @@ func testSiteReplicationStatusCountsCorsPerSite(obj ObjectLayer, _ string, bucke
if got := status.StatsSummary[remoteID].TotalCorsConfigCount; got != wantRemote { if got := status.StatsSummary[remoteID].TotalCorsConfigCount; got != wantRemote {
t.Fatalf("%s: remote TotalCorsConfigCount = %d, want %d", name, got, wantRemote) t.Fatalf("%s: remote TotalCorsConfigCount = %d, want %d", name, got, wantRemote)
} }
for _, id := range []string{localID, remoteID} {
bucketStatus := status.BucketStats[bucket][id]
wantSet := wantLocal != 0
if id == remoteID {
wantSet = wantRemote != 0
}
if bucketStatus.HasCorsCfgSet != wantSet {
t.Fatalf("%s: %s HasCorsCfgSet = %v, want %v", name, id, bucketStatus.HasCorsCfgSet, wantSet)
}
if bucketStatus.CorsCfgMismatch != wantMismatch {
t.Fatalf("%s: %s CorsCfgMismatch = %v, want %v", name, id, bucketStatus.CorsCfgMismatch, wantMismatch)
}
gotReplicated := status.StatsSummary[id].ReplicatedCorsConfig != 0
if gotReplicated != wantReplicated {
t.Fatalf("%s: %s ReplicatedCorsConfig = %d, want replicated %v", name, id, status.StatsSummary[id].ReplicatedCorsConfig, wantReplicated)
}
}
} }
check("neither site", 0, 0) check("neither site", 0, 0, false, false)
t1 := meta.Created.Add(time.Second) t1 := meta.Created.Add(time.Second)
if err = globalSiteReplicationSys.PeerBucketCorsConfigHandler(ctx, bucket, &encoded, t1); err != nil { if err = globalSiteReplicationSys.PeerBucketCorsConfigHandler(ctx, bucket, &encoded, t1); err != nil {
t.Fatal(err) t.Fatal(err)
} }
check("local site only", 1, 0) check("local site only", 1, 0, true, false)
t2 := t1.Add(time.Second) t2 := t1.Add(time.Second)
if err = globalSiteReplicationSys.PeerBucketCorsConfigHandler(ctx, bucket, nil, t2); err != nil { if err = globalSiteReplicationSys.PeerBucketCorsConfigHandler(ctx, bucket, nil, t2); err != nil {
t.Fatal(err) t.Fatal(err)
} }
remoteInfo.Buckets[bucket] = madmin.SRBucketInfo{
Bucket: bucket, CreatedAt: meta.Created, CorsConfig: &encoded,
}
check("live remote without timestamp", 0, 0, true, false)
invalidXML := base64.StdEncoding.EncodeToString([]byte(`not xml`))
remoteInfo.Buckets[bucket] = madmin.SRBucketInfo{
Bucket: bucket, CreatedAt: meta.Created, CorsConfig: &invalidXML, CorsConfigUpdatedAt: t2,
}
check("invalid remote XML", 0, 0, true, false)
remoteInfo.Buckets[bucket] = madmin.SRBucketInfo{ remoteInfo.Buckets[bucket] = madmin.SRBucketInfo{
Bucket: bucket, CreatedAt: meta.Created, CorsConfig: &encoded, CorsConfigUpdatedAt: t2, Bucket: bucket, CreatedAt: meta.Created, CorsConfig: &encoded, CorsConfigUpdatedAt: t2,
} }
check("remote site only", 0, 1) check("remote site only", 0, 1, true, false)
t3 := t2.Add(time.Second) t3 := t2.Add(time.Second)
if err = globalSiteReplicationSys.PeerBucketCorsConfigHandler(ctx, bucket, &encoded, t3); err != nil { if err = globalSiteReplicationSys.PeerBucketCorsConfigHandler(ctx, bucket, &encoded, t3); err != nil {
@@ -671,7 +699,7 @@ func testSiteReplicationStatusCountsCorsPerSite(obj ObjectLayer, _ string, bucke
remoteInfo.Buckets[bucket] = madmin.SRBucketInfo{ remoteInfo.Buckets[bucket] = madmin.SRBucketInfo{
Bucket: bucket, CreatedAt: meta.Created, CorsConfig: &encoded, CorsConfigUpdatedAt: t3, Bucket: bucket, CreatedAt: meta.Created, CorsConfig: &encoded, CorsConfigUpdatedAt: t3,
} }
check("both sites", 1, 1) check("both sites", 1, 1, false, true)
} }
func TestCORSReplicationStateOrdering(t *testing.T) { func TestCORSReplicationStateOrdering(t *testing.T) {
+86
View File
@@ -51,6 +51,14 @@ type BucketMetadataSys struct {
initialized bool initialized bool
group *singleflight.Group group *singleflight.Group
metadataMap map[string]BucketMetadata metadataMap map[string]BucketMetadata
// loadFailed tracks real buckets whose metadata could not be loaded at
// startup (concurrentLoad) or during a refresh. Such buckets are NOT
// resident in metadataMap even though the subsystem is Initialized, so a
// plain map miss cannot distinguish "not a bucket" from "known bucket whose
// config we could not read". Callers that must fail closed for a real but
// unreadable bucket (e.g. per-bucket CORS) consult this set. It is bounded
// by the number of load failures and is empty in normal operation.
loadFailed map[string]struct{}
} }
// Count returns number of bucket metadata map entries. // Count returns number of bucket metadata map entries.
@@ -67,6 +75,7 @@ func (sys *BucketMetadataSys) Remove(buckets ...string) {
for _, bucket := range buckets { for _, bucket := range buckets {
sys.group.Forget(bucket) sys.group.Forget(bucket)
delete(sys.metadataMap, bucket) delete(sys.metadataMap, bucket)
delete(sys.loadFailed, bucket)
globalBucketMonitor.DeleteBucket(bucket) globalBucketMonitor.DeleteBucket(bucket)
} }
sys.Unlock() sys.Unlock()
@@ -84,6 +93,11 @@ func (sys *BucketMetadataSys) RemoveStaleBuckets(diskBuckets set.StringSet) {
delete(sys.metadataMap, bucket) delete(sys.metadataMap, bucket)
globalBucketMonitor.DeleteBucket(bucket) globalBucketMonitor.DeleteBucket(bucket)
} }
for bucket := range sys.loadFailed {
if !diskBuckets.Contains(bucket) {
delete(sys.loadFailed, bucket)
}
}
} }
// Set - sets a new metadata in-memory. // Set - sets a new metadata in-memory.
@@ -95,6 +109,7 @@ func (sys *BucketMetadataSys) Set(bucket string, meta BucketMetadata) {
if !isMinioMetaBucketName(bucket) { if !isMinioMetaBucketName(bucket) {
sys.Lock() sys.Lock()
sys.metadataMap[bucket] = meta sys.metadataMap[bucket] = meta
delete(sys.loadFailed, bucket)
sys.Unlock() sys.Unlock()
} }
} }
@@ -379,6 +394,65 @@ func (sys *BucketMetadataSys) GetCorsConfig(bucket string) (*cors.Config, time.T
return meta.corsConfig, meta.CorsConfigUpdatedAt, nil return meta.corsConfig, meta.CorsConfigUpdatedAt, nil
} }
// GetResidentCorsConfig returns the CORS configuration for the given bucket
// using only bucket metadata that is already resident in memory. Unlike
// GetCorsConfig it never loads metadata from disk and never caches a new
// entry.
//
// The per-request CORS middleware runs before authentication, for every
// Origin-bearing request, using the validated first path segment as the bucket
// name.
// Routing that through GetCorsConfig (which loads and caches) let an
// unauthenticated client grow metadataMap without bound and trigger an
// erasure metadata probe for every distinct, attacker-controlled,
// non-existent name it sent with an Origin header (e.g. /minio/... , /api/... ,
// or random buckets). Every bucket that can carry a CORS document is made
// resident when the document is written (Set) and when metadata is loaded at
// startup (Init/concurrentLoad), so a resident-only read is complete for real
// buckets while costing only an in-memory map lookup for everything else.
//
// While bucket metadata is still loading (not yet Initialized) a non-resident
// bucket returns errBucketMetadataNotInitialized so the caller fails closed
// rather than answering with the permissive global policy for a bucket whose
// restrictive CORS document may simply not be loaded yet.
func (sys *BucketMetadataSys) GetResidentCorsConfig(bucket string) (*cors.Config, time.Time, error) {
if isMinioMetaBucketName(bucket) {
// Preserve GetConfig's semantics for the internal namespace: this is
// not a real bucket, and returning a non-errConfigNotFound error makes
// the CORS middleware fail closed rather than answer for .minio.sys
// with the permissive global policy.
return nil, time.Time{}, errInvalidArgument
}
if isReservedOrInvalidBucket(bucket, true) {
return nil, time.Time{}, errConfigNotFound
}
sys.RLock()
meta, ok := sys.metadataMap[bucket]
_, failed := sys.loadFailed[bucket]
initialized := sys.initialized
sys.RUnlock()
if ok {
if meta.corsConfigErr != nil {
return nil, meta.CorsConfigUpdatedAt, meta.corsConfigErr
}
if meta.corsConfig == nil {
return nil, time.Time{}, errConfigNotFound
}
return meta.corsConfig, meta.CorsConfigUpdatedAt, nil
}
// Not resident. Two cases must not be conflated:
// - metadata is still loading (!initialized), or this is a real bucket
// whose metadata failed to load: we cannot rule out a restrictive CORS
// config, so fail closed rather than answer with the global policy.
// - a fully initialized subsystem with no record of the name: it is not a
// bucket that can carry CORS, so fall back to the global policy without
// loading or caching metadata for an arbitrary, client-supplied name.
if !initialized || failed {
return nil, time.Time{}, errBucketMetadataNotInitialized
}
return nil, time.Time{}, errConfigNotFound
}
// GetCorsConfigXML returns the raw stored CORS configuration XML for the // GetCorsConfigXML returns the raw stored CORS configuration XML for the
// given bucket, preserving the document exactly as it was PUT (including // given bucket, preserving the document exactly as it was PUT (including
// the S3 xmlns and any unmodeled elements). // the S3 xmlns and any unmodeled elements).
@@ -576,8 +650,13 @@ func (sys *BucketMetadataSys) concurrentLoad(ctx context.Context, buckets []stri
sys.Lock() sys.Lock()
for i, meta := range bucketMetas { for i, meta := range bucketMetas {
if errs[i] != nil { if errs[i] != nil {
// Real bucket whose metadata could not be loaded: record it so
// consumers that must fail closed (per-bucket CORS) can tell it
// apart from a name that is not a bucket at all.
sys.loadFailed[buckets[i]] = struct{}{}
continue continue
} }
delete(sys.loadFailed, buckets[i])
sys.metadataMap[buckets[i]] = meta sys.metadataMap[buckets[i]] = meta
} }
sys.Unlock() sys.Unlock()
@@ -627,6 +706,9 @@ func (sys *BucketMetadataSys) refreshBucketsMetadataLoop(ctx context.Context) {
meta, err := loadBucketMetadata(ctx, sys.objAPI, bucket) meta, err := loadBucketMetadata(ctx, sys.objAPI, bucket)
if err != nil { if err != nil {
internalLogIf(ctx, err, logger.WarningKind) internalLogIf(ctx, err, logger.WarningKind)
sys.Lock()
sys.loadFailed[bucket] = struct{}{}
sys.Unlock()
wait() // wait to proceed to next entry. wait() // wait to proceed to next entry.
continue continue
} }
@@ -637,6 +719,8 @@ func (sys *BucketMetadataSys) refreshBucketsMetadataLoop(ctx context.Context) {
updated = true updated = true
sys.metadataMap[bucket] = meta sys.metadataMap[bucket] = meta
} }
// A successful (re)load clears any earlier load failure.
delete(sys.loadFailed, bucket)
sys.Unlock() sys.Unlock()
if updated { if updated {
@@ -684,6 +768,7 @@ func (sys *BucketMetadataSys) init(ctx context.Context, buckets []string) {
func (sys *BucketMetadataSys) Reset() { func (sys *BucketMetadataSys) Reset() {
sys.Lock() sys.Lock()
clear(sys.metadataMap) clear(sys.metadataMap)
clear(sys.loadFailed)
sys.Unlock() sys.Unlock()
} }
@@ -691,6 +776,7 @@ func (sys *BucketMetadataSys) Reset() {
func NewBucketMetadataSys() *BucketMetadataSys { func NewBucketMetadataSys() *BucketMetadataSys {
return &BucketMetadataSys{ return &BucketMetadataSys{
metadataMap: make(map[string]BucketMetadata), metadataMap: make(map[string]BucketMetadata),
loadFailed: make(map[string]struct{}),
group: &singleflight.Group{}, group: &singleflight.Group{},
} }
} }
+10 -10
View File
@@ -25,8 +25,6 @@ import (
"github.com/minio/minio/internal/auth" "github.com/minio/minio/internal/auth"
objectlock "github.com/minio/minio/internal/bucket/object/lock" objectlock "github.com/minio/minio/internal/bucket/object/lock"
"github.com/minio/minio/internal/bucket/replication"
xhttp "github.com/minio/minio/internal/http"
"github.com/minio/minio/internal/logger" "github.com/minio/minio/internal/logger"
"github.com/minio/pkg/v3/policy" "github.com/minio/pkg/v3/policy"
) )
@@ -150,7 +148,11 @@ func enforceRetentionBypassForDelete(ctx context.Context, r *http.Request, bucke
} }
// https://docs.aws.amazon.com/AmazonS3/latest/dev/object-lock-overview.html#object-lock-retention-modes // https://docs.aws.amazon.com/AmazonS3/latest/dev/object-lock-overview.html#object-lock-retention-modes
// If you try to delete objects protected by governance mode and have s3:BypassGovernanceRetention, the operation will succeed. // If you try to delete objects protected by governance mode and have s3:BypassGovernanceRetention, the operation will succeed.
if checkRequestAuthType(ctx, r, policy.BypassGovernanceRetentionAction, bucket, object.ObjectName) != ErrNone { if reqInfo := logger.GetReqInfo(ctx); reqInfo != nil {
reqInfo.BucketName = bucket
reqInfo.ObjectName = object.ObjectName
}
if authorizeRequest(ctx, r, policy.BypassGovernanceRetentionAction) != ErrNone {
return errAuthentication return errAuthentication
} }
} }
@@ -242,7 +244,7 @@ func enforceRetentionBypassForPut(ctx context.Context, r *http.Request, oi Objec
// For objects in "Compliance" mode, retention date cannot be shortened, and mode cannot be altered. // For objects in "Compliance" mode, retention date cannot be shortened, and mode cannot be altered.
// For objects with legal hold header set, the s3:PutObjectLegalHold permission is expected to be set // For objects with legal hold header set, the s3:PutObjectLegalHold permission is expected to be set
// Both legal hold and retention can be applied independently on an object // Both legal hold and retention can be applied independently on an object
func checkPutObjectLockAllowed(ctx context.Context, rq *http.Request, bucket, object string, getObjectInfoFn GetObjectInfoFn, retentionPermErr, legalHoldPermErr APIErrorCode) (objectlock.RetMode, objectlock.RetentionDate, objectlock.ObjectLegalHold, APIErrorCode) { func checkPutObjectLockAllowed(ctx context.Context, rq *http.Request, bucket, object string, getObjectInfoFn GetObjectInfoFn, retentionPermErr, legalHoldPermErr APIErrorCode, replicaTrusted bool) (objectlock.RetMode, objectlock.RetentionDate, objectlock.ObjectLegalHold, APIErrorCode) {
var mode objectlock.RetMode var mode objectlock.RetMode
var retainDate objectlock.RetentionDate var retainDate objectlock.RetentionDate
var legalHold objectlock.ObjectLegalHold var legalHold objectlock.ObjectLegalHold
@@ -269,9 +271,7 @@ func checkPutObjectLockAllowed(ctx context.Context, rq *http.Request, bucket, ob
return mode, retainDate, legalHold, toAPIErrorCode(ctx, err) return mode, retainDate, legalHold, toAPIErrorCode(ctx, err)
} }
replica := rq.Header.Get(xhttp.AmzBucketReplicationStatus) == replication.Replica.String() if opts.VersionID != "" && !replicaTrusted {
if opts.VersionID != "" && !replica {
if objInfo, err := getObjectInfoFn(ctx, bucket, object, opts); err == nil { if objInfo, err := getObjectInfoFn(ctx, bucket, object, opts); err == nil {
r := objectlock.GetObjectRetentionMeta(objInfo.UserDefined) r := objectlock.GetObjectRetentionMeta(objInfo.UserDefined)
t, err := objectlock.UTCNowNTP() t, err := objectlock.UTCNowNTP()
@@ -307,8 +307,8 @@ func checkPutObjectLockAllowed(ctx context.Context, rq *http.Request, bucket, ob
if err != nil { if err != nil {
return mode, retainDate, legalHold, toAPIErrorCode(ctx, err) return mode, retainDate, legalHold, toAPIErrorCode(ctx, err)
} }
rMode, rDate, err := objectlock.ParseObjectLockRetentionHeaders(rq.Header) rMode, rDate, err := objectlock.ParseObjectLockRetentionHeaders(rq.Header, replicaTrusted)
if err != nil && (!replica || rMode != "" || !rDate.IsZero()) { if err != nil && (!replicaTrusted || rMode != "" || !rDate.IsZero()) {
return mode, retainDate, legalHold, toAPIErrorCode(ctx, err) return mode, retainDate, legalHold, toAPIErrorCode(ctx, err)
} }
if retentionPermErr != ErrNone { if retentionPermErr != ErrNone {
@@ -316,7 +316,7 @@ func checkPutObjectLockAllowed(ctx context.Context, rq *http.Request, bucket, ob
} }
return rMode, rDate, legalHold, ErrNone return rMode, rDate, legalHold, ErrNone
} }
if replica { // replica inherits retention metadata only from source if replicaTrusted { // replica inherits retention metadata only from source
return "", objectlock.RetentionDate{}, legalHold, ErrNone return "", objectlock.RetentionDate{}, legalHold, ErrNone
} }
if !retentionRequested && retentionCfg.Validity > 0 { if !retentionRequested && retentionCfg.Validity > 0 {
+5 -1
View File
@@ -617,7 +617,7 @@ func (api objectAPIHandlers) ValidateBucketReplicationCredsHandler(w http.Respon
ReplicationValidityCheck: true, // set this to validate the replication config ReplicationValidityCheck: true, // set this to validate the replication config
}, },
} }
obj := path.Join(minioReservedBucket, globalLocalNodeNameHex, "deleteme") obj := replicationValidationObject(rule)
ui, err := c.PutObject(ctx, clnt.Bucket, obj, reader, int64(len(buf)), "", "", putOpts) ui, err := c.PutObject(ctx, clnt.Bucket, obj, reader, int64(len(buf)), "", "", putOpts)
if err != nil && !isReplicationPermissionCheck(ErrorRespToObjectError(err, bucket, obj)) { if err != nil && !isReplicationPermissionCheck(ErrorRespToObjectError(err, bucket, obj)) {
writeErrorResponse(ctx, w, errorCodes.ToAPIErrWithErr(ErrReplicationValidationError, fmt.Errorf("s3:ReplicateObject permissions missing for replication user: %w", err)), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErrWithErr(ErrReplicationValidationError, fmt.Errorf("s3:ReplicateObject permissions missing for replication user: %w", err)), r.URL)
@@ -658,3 +658,7 @@ func (api objectAPIHandlers) ValidateBucketReplicationCredsHandler(w http.Respon
// Write success response. // Write success response.
writeSuccessResponseHeadersOnly(w) writeSuccessResponseHeadersOnly(w)
} }
func replicationValidationObject(rule replication.Rule) string {
return path.Join(rule.Prefix(), minioReservedBucket, globalLocalNodeNameHex, "deleteme")
}
+20
View File
@@ -20,6 +20,7 @@ package cmd
import ( import (
"fmt" "fmt"
"net/http" "net/http"
"path"
"testing" "testing"
"time" "time"
@@ -287,3 +288,22 @@ func TestReplicationResyncwrapper(t *testing.T) {
} }
} }
} }
func TestReplicationValidationObjectUsesRulePrefix(t *testing.T) {
tests := []struct {
name string
rule replication.Rule
want string
}{
{name: "empty prefix", rule: replication.Rule{}, want: path.Join(minioReservedBucket, globalLocalNodeNameHex, "deleteme")},
{name: "filter prefix", rule: replication.Rule{Filter: replication.Filter{Prefix: "data/"}}, want: path.Join("data", minioReservedBucket, globalLocalNodeNameHex, "deleteme")},
{name: "and prefix", rule: replication.Rule{Filter: replication.Filter{And: replication.And{Prefix: "archive/"}}}, want: path.Join("archive", minioReservedBucket, globalLocalNodeNameHex, "deleteme")},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
if got := replicationValidationObject(test.rule); got != test.want {
t.Fatalf("replicationValidationObject() = %q, want %q", got, test.want)
}
})
}
}
+1 -1
View File
@@ -1049,7 +1049,7 @@ func DecryptObjectInfo(info *ObjectInfo, r *http.Request) (encrypted bool, err e
if encrypted { if encrypted {
if crypto.SSEC.IsEncrypted(info.UserDefined) { if crypto.SSEC.IsEncrypted(info.UserDefined) {
if !crypto.SSEC.IsRequested(headers) && !crypto.SSECopy.IsRequested(headers) { if !crypto.SSEC.IsRequested(headers) && !crypto.SSECopy.IsRequested(headers) {
if r.Header.Get(xhttp.MinIOSourceReplicationRequest) != "true" { if !isReplicaTrusted(r.Context()) {
return encrypted, errEncryptedObject return encrypted, errEncryptedObject
} }
} }
+1 -2
View File
@@ -83,7 +83,6 @@ var supportedHeaders = []string{
xhttp.AmzStorageClass, xhttp.AmzStorageClass,
xhttp.AmzObjectTagging, xhttp.AmzObjectTagging,
"expires", "expires",
xhttp.AmzBucketReplicationStatus,
"X-Minio-Replication-Server-Side-Encryption-Sealed-Key", "X-Minio-Replication-Server-Side-Encryption-Sealed-Key",
"X-Minio-Replication-Server-Side-Encryption-Seal-Algorithm", "X-Minio-Replication-Server-Side-Encryption-Seal-Algorithm",
"X-Minio-Replication-Server-Side-Encryption-Iv", "X-Minio-Replication-Server-Side-Encryption-Iv",
@@ -332,7 +331,7 @@ func extractReqParams(r *http.Request) map[string]string {
m["range"] = rangeField m["range"] = rangeField
} }
if _, ok := r.Header[xhttp.MinIOSourceReplicationRequest]; ok { if isTrustedReplication(r.Context()) {
m[xhttp.MinIOSourceReplicationRequest] = "" m[xhttp.MinIOSourceReplicationRequest] = ""
} }
return m return m
+8 -2
View File
@@ -307,7 +307,7 @@ func TestGetCopyObjectMetadataFromHeaderReplication(t *testing.T) {
} }
} }
func TestCloneRequestWithoutCopyReplicationHeaders(t *testing.T) { func TestCloneRequestWithoutReplicationHeaders(t *testing.T) {
req, err := http.NewRequest(http.MethodPut, "http://localhost/test", nil) req, err := http.NewRequest(http.MethodPut, "http://localhost/test", nil)
if err != nil { if err != nil {
t.Fatal(err) t.Fatal(err)
@@ -320,9 +320,12 @@ func TestCloneRequestWithoutCopyReplicationHeaders(t *testing.T) {
req.Header.Set(xhttp.MinIOSourceObjectLegalHoldTimestamp, "2026-04-15T10:00:00Z") req.Header.Set(xhttp.MinIOSourceObjectLegalHoldTimestamp, "2026-04-15T10:00:00Z")
req.Header.Set(xhttp.MinIOReplicationActualObjectSize, "123") req.Header.Set(xhttp.MinIOReplicationActualObjectSize, "123")
req.Header.Set(ReplicationSsecChecksumHeader, "checksum") req.Header.Set(ReplicationSsecChecksumHeader, "checksum")
req.Header.Set(xhttp.AmzBucketReplicationStatus, "REPLICA")
req.Header.Set(xhttp.MinIOSourceDeleteMarker, "true")
req.Header.Set("X-Minio-Replication-Server-Side-Encryption-Sealed-Key", "sealed")
req.Header.Set("Content-Type", "application/octet-stream") req.Header.Set("Content-Type", "application/octet-stream")
clone := cloneRequestWithoutCopyReplicationHeaders(req) clone := cloneRequestWithoutReplicationHeaders(t.Context(), req)
if clone == req { if clone == req {
t.Fatal("expected cloned request") t.Fatal("expected cloned request")
} }
@@ -336,6 +339,9 @@ func TestCloneRequestWithoutCopyReplicationHeaders(t *testing.T) {
xhttp.MinIOSourceObjectLegalHoldTimestamp, xhttp.MinIOSourceObjectLegalHoldTimestamp,
xhttp.MinIOReplicationActualObjectSize, xhttp.MinIOReplicationActualObjectSize,
ReplicationSsecChecksumHeader, ReplicationSsecChecksumHeader,
xhttp.AmzBucketReplicationStatus,
xhttp.MinIOSourceDeleteMarker,
"X-Minio-Replication-Server-Side-Encryption-Sealed-Key",
} { } {
if got := clone.Header.Get(header); got != "" { if got := clone.Header.Get(header); got != "" {
t.Fatalf("expected %s to be stripped, got %q", header, got) t.Fatalf("expected %s to be stripped, got %q", header, got)
+46 -34
View File
@@ -41,9 +41,6 @@ func getDefaultOpts(header http.Header, copySource bool, metadata map[string]str
opts.ProxyHeaderSet = true opts.ProxyHeaderSet = true
opts.ProxyRequest = strings.Join(v, "") == "true" opts.ProxyRequest = strings.Join(v, "") == "true"
} }
if _, ok := header[xhttp.MinIOSourceReplicationRequest]; ok {
opts.ReplicationRequest = true
}
opts.Speedtest = header.Get(globalObjectPerfUserMetadata) != "" opts.Speedtest = header.Get(globalObjectPerfUserMetadata) != ""
if copySource { if copySource {
@@ -116,12 +113,15 @@ func getOpts(ctx context.Context, r *http.Request, bucket, object string) (Objec
} }
opts.PartNumber = partNumber opts.PartNumber = partNumber
opts.VersionID = vid opts.VersionID = vid
opts.ReplicationRequest = isTrustedReplication(ctx)
delMarker, err := parseBoolHeader(bucket, object, r.Header, xhttp.MinIOSourceDeleteMarker) if opts.ReplicationRequest {
if err != nil { delMarker, err := parseBoolHeader(bucket, object, r.Header, xhttp.MinIOSourceDeleteMarker)
return opts, err if err != nil {
return opts, err
}
opts.DeleteMarker = delMarker
} }
opts.DeleteMarker = delMarker
replReadyCheck, err := parseBoolHeader(bucket, object, r.Header, xhttp.MinIOCheckDMReplicationReady) replReadyCheck, err := parseBoolHeader(bucket, object, r.Header, xhttp.MinIOCheckDMReplicationReady)
if err != nil { if err != nil {
@@ -297,20 +297,22 @@ func delOpts(ctx context.Context, r *http.Request, bucket, object string) (opts
opts.VersionID = nullVersionID opts.VersionID = nullVersionID
} }
delMarker, err := parseBoolHeader(bucket, object, r.Header, xhttp.MinIOSourceDeleteMarker) if isTrustedReplication(ctx) {
if err != nil { delMarker, err := parseBoolHeader(bucket, object, r.Header, xhttp.MinIOSourceDeleteMarker)
return opts, err
}
opts.DeleteMarker = delMarker
mtime := strings.TrimSpace(r.Header.Get(xhttp.MinIOSourceMTime))
if mtime != "" {
opts.MTime, err = time.Parse(time.RFC3339Nano, mtime)
if err != nil { if err != nil {
return opts, InvalidArgument{ return opts, err
Bucket: bucket, }
Object: object, opts.DeleteMarker = delMarker
Err: fmt.Errorf("Unable to parse %s, failed with %w", xhttp.MinIOSourceMTime, err),
mtime := strings.TrimSpace(r.Header.Get(xhttp.MinIOSourceMTime))
if mtime != "" {
opts.MTime, err = time.Parse(time.RFC3339Nano, mtime)
if err != nil {
return opts, InvalidArgument{
Bucket: bucket,
Object: object,
Err: fmt.Errorf("Unable to parse %s, failed with %w", xhttp.MinIOSourceMTime, err),
}
} }
} }
} }
@@ -319,10 +321,10 @@ func delOpts(ctx context.Context, r *http.Request, bucket, object string) (opts
// get ObjectOptions for PUT calls from encryption headers and metadata // get ObjectOptions for PUT calls from encryption headers and metadata
func putOptsFromReq(ctx context.Context, r *http.Request, bucket, object string, metadata map[string]string) (opts ObjectOptions, err error) { func putOptsFromReq(ctx context.Context, r *http.Request, bucket, object string, metadata map[string]string) (opts ObjectOptions, err error) {
return putOpts(ctx, bucket, object, r.Form.Get(xhttp.VersionID), r.Header, metadata) return putOpts(ctx, bucket, object, r.Form.Get(xhttp.VersionID), r.Header, metadata, isTrustedReplication(ctx))
} }
func putOpts(ctx context.Context, bucket, object, vid string, hdrs http.Header, metadata map[string]string) (opts ObjectOptions, err error) { func putOpts(ctx context.Context, bucket, object, vid string, hdrs http.Header, metadata map[string]string, trustedReplication bool) (opts ObjectOptions, err error) {
versioned := globalBucketVersioningSys.PrefixEnabled(bucket, object) versioned := globalBucketVersioningSys.PrefixEnabled(bucket, object)
versionSuspended := globalBucketVersioningSys.PrefixSuspended(bucket, object) versionSuspended := globalBucketVersioningSys.PrefixSuspended(bucket, object)
@@ -344,7 +346,7 @@ func putOpts(ctx context.Context, bucket, object, vid string, hdrs http.Header,
} }
} }
} }
opts, err = putOptsFromHeaders(ctx, hdrs, metadata) opts, err = putOptsFromHeaders(ctx, hdrs, metadata, trustedReplication)
if err != nil { if err != nil {
return opts, InvalidArgument{ return opts, InvalidArgument{
Bucket: bucket, Bucket: bucket,
@@ -365,8 +367,15 @@ func putOpts(ctx context.Context, bucket, object, vid string, hdrs http.Header,
return opts, nil return opts, nil
} }
func putOptsFromHeaders(ctx context.Context, hdr http.Header, metadata map[string]string) (opts ObjectOptions, err error) { func putOptsFromHeaders(ctx context.Context, hdr http.Header, metadata map[string]string, trustedReplication bool) (opts ObjectOptions, err error) {
mtimeStr := strings.TrimSpace(hdr.Get(xhttp.MinIOSourceMTime)) var mtimeStr, retaintimeStr, lholdtimeStr, tagtimeStr, etag string
if trustedReplication {
mtimeStr = strings.TrimSpace(hdr.Get(xhttp.MinIOSourceMTime))
retaintimeStr = strings.TrimSpace(hdr.Get(xhttp.MinIOSourceObjectRetentionTimestamp))
lholdtimeStr = strings.TrimSpace(hdr.Get(xhttp.MinIOSourceObjectLegalHoldTimestamp))
tagtimeStr = strings.TrimSpace(hdr.Get(xhttp.MinIOSourceTaggingTimestamp))
etag = strings.TrimSpace(hdr.Get(xhttp.MinIOSourceETag))
}
var mtime time.Time var mtime time.Time
if mtimeStr != "" { if mtimeStr != "" {
mtime, err = time.Parse(time.RFC3339Nano, mtimeStr) mtime, err = time.Parse(time.RFC3339Nano, mtimeStr)
@@ -374,7 +383,6 @@ func putOptsFromHeaders(ctx context.Context, hdr http.Header, metadata map[strin
return opts, fmt.Errorf("Unable to parse %s, failed with %w", xhttp.MinIOSourceMTime, err) return opts, fmt.Errorf("Unable to parse %s, failed with %w", xhttp.MinIOSourceMTime, err)
} }
} }
retaintimeStr := strings.TrimSpace(hdr.Get(xhttp.MinIOSourceObjectRetentionTimestamp))
var retaintimestmp time.Time var retaintimestmp time.Time
if retaintimeStr != "" { if retaintimeStr != "" {
retaintimestmp, err = time.Parse(time.RFC3339, retaintimeStr) retaintimestmp, err = time.Parse(time.RFC3339, retaintimeStr)
@@ -383,7 +391,6 @@ func putOptsFromHeaders(ctx context.Context, hdr http.Header, metadata map[strin
} }
} }
lholdtimeStr := strings.TrimSpace(hdr.Get(xhttp.MinIOSourceObjectLegalHoldTimestamp))
var lholdtimestmp time.Time var lholdtimestmp time.Time
if lholdtimeStr != "" { if lholdtimeStr != "" {
lholdtimestmp, err = time.Parse(time.RFC3339, lholdtimeStr) lholdtimestmp, err = time.Parse(time.RFC3339, lholdtimeStr)
@@ -391,7 +398,6 @@ func putOptsFromHeaders(ctx context.Context, hdr http.Header, metadata map[strin
return opts, fmt.Errorf("Unable to parse %s, failed with %w", xhttp.MinIOSourceObjectLegalHoldTimestamp, err) return opts, fmt.Errorf("Unable to parse %s, failed with %w", xhttp.MinIOSourceObjectLegalHoldTimestamp, err)
} }
} }
tagtimeStr := strings.TrimSpace(hdr.Get(xhttp.MinIOSourceTaggingTimestamp))
var taggingtimestmp time.Time var taggingtimestmp time.Time
if tagtimeStr != "" { if tagtimeStr != "" {
taggingtimestmp, err = time.Parse(time.RFC3339, tagtimeStr) taggingtimestmp, err = time.Parse(time.RFC3339, tagtimeStr)
@@ -404,7 +410,6 @@ func putOptsFromHeaders(ctx context.Context, hdr http.Header, metadata map[strin
metadata = make(map[string]string) metadata = make(map[string]string)
} }
etag := strings.TrimSpace(hdr.Get(xhttp.MinIOSourceETag))
if crypto.S3KMS.IsRequested(hdr) { if crypto.S3KMS.IsRequested(hdr) {
keyID, context, err := crypto.S3KMS.ParseHTTP(hdr) keyID, context, err := crypto.S3KMS.ParseHTTP(hdr)
if err != nil { if err != nil {
@@ -419,6 +424,7 @@ func putOptsFromHeaders(ctx context.Context, hdr http.Header, metadata map[strin
UserDefined: metadata, UserDefined: metadata,
MTime: mtime, MTime: mtime,
PreserveETag: etag, PreserveETag: etag,
ReplicationRequest: trustedReplication,
} }
return op, nil return op, nil
} }
@@ -429,6 +435,7 @@ func putOptsFromHeaders(ctx context.Context, hdr http.Header, metadata map[strin
} }
opts.MTime = mtime opts.MTime = mtime
opts.ReplicationRequest = trustedReplication
opts.ReplicationSourceLegalholdTimestamp = lholdtimestmp opts.ReplicationSourceLegalholdTimestamp = lholdtimestmp
opts.ReplicationSourceRetentionTimestamp = retaintimestmp opts.ReplicationSourceRetentionTimestamp = retaintimestmp
opts.ReplicationSourceTaggingTimestamp = taggingtimestmp opts.ReplicationSourceTaggingTimestamp = taggingtimestmp
@@ -454,12 +461,17 @@ func copySrcOpts(ctx context.Context, r *http.Request, bucket, object string) (O
if err != nil { if err != nil {
return opts, err return opts, err
} }
opts.ReplicationRequest = isReplicaTrusted(ctx)
return opts, nil return opts, nil
} }
// get ObjectOptions for CompleteMultipart calls // get ObjectOptions for CompleteMultipart calls
func completeMultipartOpts(ctx context.Context, r *http.Request, bucket, object string) (opts ObjectOptions, err error) { func completeMultipartOpts(ctx context.Context, r *http.Request, bucket, object string) (opts ObjectOptions, err error) {
mtimeStr := strings.TrimSpace(r.Header.Get(xhttp.MinIOSourceMTime)) trustedReplication := isTrustedReplication(ctx)
var mtimeStr string
if trustedReplication {
mtimeStr = strings.TrimSpace(r.Header.Get(xhttp.MinIOSourceMTime))
}
var mtime time.Time var mtime time.Time
if mtimeStr != "" { if mtimeStr != "" {
mtime, err = time.Parse(time.RFC3339Nano, mtimeStr) mtime, err = time.Parse(time.RFC3339Nano, mtimeStr)
@@ -495,12 +507,12 @@ func completeMultipartOpts(ctx context.Context, r *http.Request, bucket, object
} }
} }
} }
if _, ok := r.Header[xhttp.MinIOSourceReplicationRequest]; ok { if trustedReplication {
opts.ReplicationRequest = true opts.ReplicationRequest = true
opts.UserDefined[ReservedMetadataPrefix+"Actual-Object-Size"] = r.Header.Get(xhttp.MinIOReplicationActualObjectSize) opts.UserDefined[ReservedMetadataPrefix+"Actual-Object-Size"] = r.Header.Get(xhttp.MinIOReplicationActualObjectSize)
} if r.Header.Get(ReplicationSsecChecksumHeader) != "" {
if r.Header.Get(ReplicationSsecChecksumHeader) != "" { opts.UserDefined[ReplicationSsecChecksumHeader] = r.Header.Get(ReplicationSsecChecksumHeader)
opts.UserDefined[ReplicationSsecChecksumHeader] = r.Header.Get(ReplicationSsecChecksumHeader) }
} }
return opts, nil return opts, nil
} }
+179 -98
View File
@@ -356,6 +356,18 @@ func (api objectAPIHandlers) getObjectHandler(ctx context.Context, objectAPI Obj
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL)
return return
} }
if hasReplicationMarkerHeader(r.Header) {
trusted := hasReplicationMarker(r.Header) &&
replicationPermissionAllowed(ctx, r, bucket, object, policy.ReplicateObjectAction)
ctx, r = applyReplicationTrust(ctx, r, trusted, trusted)
if trusted {
opts, err = getOpts(ctx, r, bucket, object)
if err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return
}
}
}
getObjectNInfo := objectAPI.GetObjectNInfo getObjectNInfo := objectAPI.GetObjectNInfo
@@ -494,8 +506,8 @@ func (api objectAPIHandlers) getObjectHandler(ctx context.Context, objectAPI Obj
} }
// filter object lock metadata if permission does not permit // filter object lock metadata if permission does not permit
getRetPerms := checkRequestAuthType(ctx, r, policy.GetObjectRetentionAction, bucket, object) getRetPerms := authorizeRequest(ctx, r, policy.GetObjectRetentionAction)
legalHoldPerms := checkRequestAuthType(ctx, r, policy.GetObjectLegalHoldAction, bucket, object) legalHoldPerms := authorizeRequest(ctx, r, policy.GetObjectLegalHoldAction)
// filter object lock metadata if permission does not permit // filter object lock metadata if permission does not permit
objInfo.UserDefined = objectlock.FilterObjectLockMetadata(objInfo.UserDefined, getRetPerms != ErrNone, legalHoldPerms != ErrNone) objInfo.UserDefined = objectlock.FilterObjectLockMetadata(objInfo.UserDefined, getRetPerms != ErrNone, legalHoldPerms != ErrNone)
@@ -599,10 +611,16 @@ func (api objectAPIHandlers) getObjectAttributesHandler(ctx context.Context, obj
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL)
return return
} }
if hasReplicationMarkerHeader(r.Header) {
trusted := hasReplicationMarker(r.Header) &&
replicationPermissionAllowed(ctx, r, bucket, object, policy.ReplicateObjectAction)
ctx, r = applyReplicationTrust(ctx, r, trusted, trusted)
opts.ReplicationRequest = trusted
}
objInfo, err := objectAPI.GetObjectInfo(ctx, bucket, object, opts) objInfo, err := objectAPI.GetObjectInfo(ctx, bucket, object, opts)
if err != nil { if err != nil {
s3Error = checkRequestAuthType(ctx, r, policy.ListBucketAction, bucket, object) s3Error = authorizeRequest(ctx, r, policy.ListBucketAction)
if s3Error == ErrNone { if s3Error == ErrNone {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return return
@@ -622,9 +640,7 @@ func (api objectAPIHandlers) getObjectAttributesHandler(ctx context.Context, obj
// Only a caller authorized to replicate this object may read SSE-C // Only a caller authorized to replicate this object may read SSE-C
// attributes without presenting the customer key. The header alone is // attributes without presenting the customer key. The header alone is
// client controlled, so it cannot stand in for that authorization. // client controlled, so it cannot stand in for that authorization.
trustedReplicationRequest := r.Header.Get(xhttp.MinIOSourceReplicationRequest) == "true" && if crypto.SSEC.IsEncrypted(objInfo.UserDefined) && !isReplicaTrusted(ctx) {
checkRequestAuthType(ctx, r, policy.ReplicateObjectAction, bucket, object) == ErrNone
if crypto.SSEC.IsEncrypted(objInfo.UserDefined) && !trustedReplicationRequest {
if _, err = crypto.SSEC.UnsealObjectKey(r.Header, objInfo.UserDefined, bucket, object); err != nil { if _, err = crypto.SSEC.UnsealObjectKey(r.Header, objInfo.UserDefined, bucket, object); err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return return
@@ -800,6 +816,18 @@ func (api objectAPIHandlers) headObjectHandler(ctx context.Context, objectAPI Ob
writeErrorResponseHeadersOnly(w, errorCodes.ToAPIErr(s3Error)) writeErrorResponseHeadersOnly(w, errorCodes.ToAPIErr(s3Error))
return return
} }
if hasReplicationMarkerHeader(r.Header) {
trusted := hasReplicationMarker(r.Header) &&
replicationPermissionAllowed(ctx, r, bucket, object, policy.ReplicateObjectAction)
ctx, r = applyReplicationTrust(ctx, r, trusted, trusted)
if trusted {
opts, err = getOpts(ctx, r, bucket, object)
if err != nil {
writeErrorResponseHeadersOnly(w, toAPIError(ctx, err))
return
}
}
}
// Get request range. // Get request range.
var rs *HTTPRangeSpec var rs *HTTPRangeSpec
@@ -911,8 +939,8 @@ func (api objectAPIHandlers) headObjectHandler(ctx context.Context, objectAPI Ob
} }
// filter object lock metadata if permission does not permit // filter object lock metadata if permission does not permit
getRetPerms := checkRequestAuthType(ctx, r, policy.GetObjectRetentionAction, bucket, object) getRetPerms := authorizeRequest(ctx, r, policy.GetObjectRetentionAction)
legalHoldPerms := checkRequestAuthType(ctx, r, policy.GetObjectLegalHoldAction, bucket, object) legalHoldPerms := authorizeRequest(ctx, r, policy.GetObjectLegalHoldAction)
// filter object lock metadata if permission does not permit // filter object lock metadata if permission does not permit
objInfo.UserDefined = objectlock.FilterObjectLockMetadata(objInfo.UserDefined, getRetPerms != ErrNone, legalHoldPerms != ErrNone) objInfo.UserDefined = objectlock.FilterObjectLockMetadata(objInfo.UserDefined, getRetPerms != ErrNone, legalHoldPerms != ErrNone)
@@ -938,10 +966,12 @@ func (api objectAPIHandlers) headObjectHandler(ctx context.Context, objectAPI Ob
w.Header().Set(xhttp.AmzServerSideEncryptionKmsContext, kmsCtx) w.Header().Set(xhttp.AmzServerSideEncryptionKmsContext, kmsCtx)
} }
case crypto.SSEC: case crypto.SSEC:
// Validate the SSE-C Key set in the header. if !isReplicaTrusted(ctx) {
if _, err = crypto.SSEC.UnsealObjectKey(r.Header, objInfo.UserDefined, bucket, object); err != nil { // Validate the SSE-C Key set in the header for ordinary reads.
writeErrorResponseHeadersOnly(w, toAPIError(ctx, err)) if _, err = crypto.SSEC.UnsealObjectKey(r.Header, objInfo.UserDefined, bucket, object); err != nil {
return writeErrorResponseHeadersOnly(w, toAPIError(ctx, err))
return
}
} }
w.Header().Set(xhttp.AmzServerSideEncryptionCustomerAlgorithm, r.Header.Get(xhttp.AmzServerSideEncryptionCustomerAlgorithm)) w.Header().Set(xhttp.AmzServerSideEncryptionCustomerAlgorithm, r.Header.Get(xhttp.AmzServerSideEncryptionCustomerAlgorithm))
w.Header().Set(xhttp.AmzServerSideEncryptionCustomerKeyMD5, r.Header.Get(xhttp.AmzServerSideEncryptionCustomerKeyMD5)) w.Header().Set(xhttp.AmzServerSideEncryptionCustomerKeyMD5, r.Header.Get(xhttp.AmzServerSideEncryptionCustomerKeyMD5))
@@ -1093,31 +1123,6 @@ func getCpObjMetadataFromHeader(ctx context.Context, r *http.Request, userMeta m
return defaultMeta, nil return defaultMeta, nil
} }
func cloneRequestWithoutCopyReplicationHeaders(r *http.Request) *http.Request {
if r == nil {
return nil
}
clone := new(http.Request)
*clone = *r
clone.Header = r.Header.Clone()
for _, header := range []string{
xhttp.MinIOSourceReplicationRequest,
xhttp.MinIOSourceETag,
xhttp.MinIOSourceMTime,
xhttp.MinIOSourceTaggingTimestamp,
xhttp.MinIOSourceObjectRetentionTimestamp,
xhttp.MinIOSourceObjectLegalHoldTimestamp,
xhttp.MinIOReplicationActualObjectSize,
ReplicationSsecChecksumHeader,
} {
clone.Header.Del(header)
}
return clone
}
func copyDestinationSSEHeaders(h http.Header) http.Header { func copyDestinationSSEHeaders(h http.Header) http.Header {
dst := h.Clone() dst := h.Clone()
dst.Del(xhttp.AmzServerSideEncryptionCopyCustomerAlgorithm) dst.Del(xhttp.AmzServerSideEncryptionCopyCustomerAlgorithm)
@@ -1303,28 +1308,30 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrInvalidStorageClass), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrInvalidStorageClass), r.URL)
return return
} }
allowReplicationMetadata := false rawReplica := hasReplicaStatus(r.Header)
if r.Header.Get(xhttp.AmzBucketReplicationStatus) == replication.Replica.String() { markerExact := hasReplicationMarker(r.Header)
if s3Error := checkRequestAuthType(ctx, r, policy.ReplicateObjectAction, dstBucket, dstObject); s3Error != ErrNone { replicationPermitted := false
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL) if rawReplica || markerExact {
return replicationPermitted = replicationPermissionAllowed(ctx, r, dstBucket, dstObject, policy.ReplicateObjectAction)
}
allowReplicationMetadata = true
} }
trustedReplicationRequest := allowReplicationMetadata && r.Header.Get(xhttp.MinIOSourceReplicationRequest) == "true" if rawReplica && !replicationPermitted {
optsReq := r writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrAccessDenied), r.URL)
if !trustedReplicationRequest { return
optsReq = cloneRequestWithoutCopyReplicationHeaders(r)
} }
trustedReplication := markerExact && replicationPermitted
replicaTrusted := trustedReplication && rawReplica
if hasReplicationRequestHeaders(r.Header) {
ctx, r = applyReplicationTrust(ctx, r, trustedReplication, replicaTrusted)
}
allowReplicationMetadata := replicaTrusted
// Check if bucket encryption is enabled // Check if bucket encryption is enabled
sseConfig, _ := globalBucketSSEConfigSys.Get(dstBucket) sseConfig, _ := globalBucketSSEConfigSys.Get(dstBucket)
sseConfig.Apply(r.Header, sse.ApplyOptions{ sseConfig.Apply(r.Header, sse.ApplyOptions{
AutoEncrypt: globalAutoEncryption, AutoEncrypt: globalAutoEncryption,
}) })
var srcOpts, dstOpts ObjectOptions var srcOpts, dstOpts ObjectOptions
srcOpts, err = copySrcOpts(ctx, optsReq, srcBucket, srcObject) srcOpts, err = copySrcOpts(ctx, r, srcBucket, srcObject)
if err != nil { if err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return return
@@ -1336,14 +1343,14 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
VersionID: srcOpts.VersionID, VersionID: srcOpts.VersionID,
Versioned: srcOpts.Versioned, Versioned: srcOpts.Versioned,
VersionSuspended: srcOpts.VersionSuspended, VersionSuspended: srcOpts.VersionSuspended,
ReplicationRequest: trustedReplicationRequest, ReplicationRequest: replicaTrusted,
} }
getSSE := encrypt.SSE(srcOpts.ServerSideEncryption) getSSE := encrypt.SSE(srcOpts.ServerSideEncryption)
if getSSE != srcOpts.ServerSideEncryption { if getSSE != srcOpts.ServerSideEncryption {
getOpts.ServerSideEncryption = getSSE getOpts.ServerSideEncryption = getSSE
} }
dstOpts, err = copyDstOpts(ctx, optsReq, dstBucket, dstObject, nil) dstOpts, err = copyDstOpts(ctx, r, dstBucket, dstObject, nil)
if err != nil { if err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return return
@@ -1353,7 +1360,7 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
getObjectNInfo := objectAPI.GetObjectNInfo getObjectNInfo := objectAPI.GetObjectNInfo
checkCopyPrecondFn := func(o ObjectInfo) bool { checkCopyPrecondFn := func(o ObjectInfo) bool {
if _, err := DecryptObjectInfo(&o, optsReq); err != nil { if _, err := DecryptObjectInfo(&o, r); err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return true return true
} }
@@ -1465,7 +1472,7 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
return return
} }
// Encryption parameters not present for this object. // Encryption parameters not present for this object.
if crypto.SSEC.IsEncrypted(srcInfo.UserDefined) && !crypto.SSECopy.IsRequested(r.Header) && !trustedReplicationRequest { if crypto.SSEC.IsEncrypted(srcInfo.UserDefined) && !crypto.SSECopy.IsRequested(r.Header) && !replicaTrusted {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrInvalidSSECustomerAlgorithm), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrInvalidSSECustomerAlgorithm), r.URL)
return return
} }
@@ -1716,7 +1723,7 @@ func (api objectAPIHandlers) CopyObjectHandler(w http.ResponseWriter, r *http.Re
getObjectInfo := objectAPI.GetObjectInfo getObjectInfo := objectAPI.GetObjectInfo
// apply default bucket configuration/governance headers for dest side. // apply default bucket configuration/governance headers for dest side.
retentionMode, retentionDate, legalHold, s3Err := checkPutObjectLockAllowed(ctx, r, dstBucket, dstObject, getObjectInfo, retPerms, holdPerms) retentionMode, retentionDate, legalHold, s3Err := checkPutObjectLockAllowed(ctx, r, dstBucket, dstObject, getObjectInfo, retPerms, holdPerms, replicaTrusted)
if s3Err == ErrNone && retentionMode.Valid() { if s3Err == ErrNone && retentionMode.Valid() {
lastretentionTimestamp := srcInfo.UserDefined[ReservedMetadataPrefixLower+ObjectLockRetentionTimestamp] lastretentionTimestamp := srcInfo.UserDefined[ReservedMetadataPrefixLower+ObjectLockRetentionTimestamp]
if dstOpts.ReplicationRequest { if dstOpts.ReplicationRequest {
@@ -2066,6 +2073,16 @@ func (api objectAPIHandlers) PutObjectHandler(w http.ResponseWriter, r *http.Req
} }
} }
rawReplica := hasReplicaStatus(r.Header)
markerExact := hasReplicationMarker(r.Header)
replicationPermitted := false
if rawReplica || markerExact {
replicationPermitted = replicationPermissionAllowed(ctx, r, bucket, object, policy.ReplicateObjectAction)
}
if rawReplica && !replicationPermitted {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrAccessDenied), r.URL)
return
}
if _, ok := r.Header[xhttp.MinIOSourceReplicationCheck]; ok { if _, ok := r.Header[xhttp.MinIOSourceReplicationCheck]; ok {
// requests to just validate replication settings and permissions are not allowed to write data // requests to just validate replication settings and permissions are not allowed to write data
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrReplicationPermissionCheckError), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrReplicationPermissionCheckError), r.URL)
@@ -2076,18 +2093,22 @@ func (api objectAPIHandlers) PutObjectHandler(w http.ResponseWriter, r *http.Req
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return return
} }
if r.Header.Get(xhttp.AmzBucketReplicationStatus) == replication.Replica.String() { trustedReplication := markerExact && replicationPermitted
if s3Err = isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.ReplicateObjectAction); s3Err != ErrNone { replicaTrusted := trustedReplication && rawReplica
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Err), r.URL) if hasReplicationRequestHeaders(r.Header) {
return ctx, r = applyReplicationTrust(ctx, r, trustedReplication, replicaTrusted)
} }
if replicaTrusted {
if err = extractReplicationMetadataFromMime(ctx, textproto.MIMEHeader(r.Header), metadata); err != nil { if err = extractReplicationMetadataFromMime(ctx, textproto.MIMEHeader(r.Header), metadata); err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return return
} }
metadata[xhttp.AmzBucketReplicationStatus] = replication.Replica.String()
metadata[ReservedMetadataPrefixLower+ReplicaStatus] = replication.Replica.String() metadata[ReservedMetadataPrefixLower+ReplicaStatus] = replication.Replica.String()
metadata[ReservedMetadataPrefixLower+ReplicaTimestamp] = UTCNow().Format(time.RFC3339Nano) metadata[ReservedMetadataPrefixLower+ReplicaTimestamp] = UTCNow().Format(time.RFC3339Nano)
defer globalReplicationStats.Load().UpdateReplicaStat(bucket, size) defer globalReplicationStats.Load().UpdateReplicaStat(bucket, size)
} else {
delete(metadata, xhttp.AmzBucketReplicationStatus)
} }
// Check if bucket encryption is enabled // Check if bucket encryption is enabled
@@ -2186,7 +2207,7 @@ func (api objectAPIHandlers) PutObjectHandler(w http.ResponseWriter, r *http.Req
getObjectInfo := objectAPI.GetObjectInfo getObjectInfo := objectAPI.GetObjectInfo
retentionMode, retentionDate, legalHold, s3Err := checkPutObjectLockAllowed(ctx, r, bucket, object, getObjectInfo, retPerms, holdPerms) retentionMode, retentionDate, legalHold, s3Err := checkPutObjectLockAllowed(ctx, r, bucket, object, getObjectInfo, retPerms, holdPerms, isReplicaTrusted(ctx))
if s3Err == ErrNone && retentionMode.Valid() { if s3Err == ErrNone && retentionMode.Valid() {
metadata[strings.ToLower(xhttp.AmzObjectLockMode)] = string(retentionMode) metadata[strings.ToLower(xhttp.AmzObjectLockMode)] = string(retentionMode)
metadata[strings.ToLower(xhttp.AmzObjectLockRetainUntilDate)] = amztime.ISO8601Format(retentionDate.UTC()) metadata[strings.ToLower(xhttp.AmzObjectLockRetainUntilDate)] = amztime.ISO8601Format(retentionDate.UTC())
@@ -2470,7 +2491,6 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
sha256hex = getContentSha256Cksum(r, serviceS3) sha256hex = getContentSha256Cksum(r, serviceS3)
} }
} }
hreader, err := hash.NewReader(ctx, reader, size, md5hex, sha256hex, size) hreader, err := hash.NewReader(ctx, reader, size, md5hex, sha256hex, size)
if err != nil { if err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
@@ -2491,14 +2511,26 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
sseConfig.Apply(r.Header, sse.ApplyOptions{ sseConfig.Apply(r.Header, sse.ApplyOptions{
AutoEncrypt: globalAutoEncryption, AutoEncrypt: globalAutoEncryption,
}) })
entryRequestBase := r.Clone(ctx)
// The streaming reader fills r.Trailer while untar writes small entries in
// parallel. Entry authorization never consumes trailers, so keep them out
// of the immutable request template cloned by those goroutines. Snapshot
// after applying bucket defaults so extracted objects retain encryption.
entryRequestBase.Trailer = nil
rawReplica := hasReplicaStatus(r.Header)
markerExact := hasReplicationMarker(r.Header)
trustedRequestCtx := withReplicationTrust(ctx, true, rawReplica)
trustedRequest := entryRequestBase.WithContext(trustedRequestCtx)
cleanRequestCtx := withReplicationTrust(ctx, false, false)
cleanRequest := cloneRequestWithoutReplicationHeaders(cleanRequestCtx, entryRequestBase)
trustedReqParams := extractReqParams(trustedRequest)
cleanReqParams := extractReqParams(cleanRequest)
retPerms := isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.PutObjectRetentionAction) retPerms := isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.PutObjectRetentionAction)
holdPerms := isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.PutObjectLegalHoldAction) holdPerms := isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.PutObjectLegalHoldAction)
getObjectInfo := objectAPI.GetObjectInfo getObjectInfo := objectAPI.GetObjectInfo
// These are static for all objects extracted.
reqParams := extractReqParams(r)
respElements := map[string]string{ respElements := map[string]string{
"requestId": w.Header().Get(xhttp.AmzRequestID), "requestId": w.Header().Get(xhttp.AmzRequestID),
"nodeId": w.Header().Get(xhttp.AmzRequestHostID), "nodeId": w.Header().Get(xhttp.AmzRequestHostID),
@@ -2506,12 +2538,46 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
if sc == "" { if sc == "" {
sc = storageclass.STANDARD sc = storageclass.STANDARD
} }
reqInfo := logger.GetReqInfo(ctx)
if reqInfo == nil {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrAccessDenied), r.URL)
return
}
reqInfo.RLock()
tarCred := reqInfo.Cred
tarOwner := reqInfo.Owner
reqInfo.RUnlock()
var tarS3Err atomic.Int32
setTarS3Err := func(code APIErrorCode) {
tarS3Err.CompareAndSwap(int32(ErrNone), int32(code))
}
ignoreEntryErrors := opts.ignoreErrs
putObjectTar := func(reader io.Reader, info os.FileInfo, object string) error { putObjectTar := func(reader io.Reader, info os.FileInfo, object string) error {
size := info.Size() size := info.Size()
if s3Err = isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.PutObjectAction); s3Err != ErrNone { entryAuthReq := entryRequestBase.Clone(ctx)
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Err), r.URL) entryS3Err := isPutActionAllowedWithCred(bucket, object, entryAuthReq, policy.PutObjectAction, nil, tarCred, tarOwner)
return errors.New(errorCodes.ToAPIErr(s3Err).Code) if entryS3Err != ErrNone {
setTarS3Err(entryS3Err)
return errors.New(errorCodes.ToAPIErr(entryS3Err).Code)
}
replicationPermitted := false
if tarCred.AccessKey != "" && (rawReplica || markerExact) {
replicationPermitted = isPutActionAllowedWithCred(bucket, object, entryAuthReq, policy.ReplicateObjectAction, nil, tarCred, tarOwner) == ErrNone
}
if rawReplica && !replicationPermitted {
setTarS3Err(ErrAccessDenied)
return errors.New(errorCodes.ToAPIErr(ErrAccessDenied).Code)
}
entryTrusted := markerExact && replicationPermitted
replicaTrusted := entryTrusted && rawReplica
entryCtx := cleanRequestCtx
entryReq := cloneRequestWithoutReplicationHeaders(cleanRequestCtx, entryAuthReq)
reqParams := cleanReqParams
if entryTrusted {
entryCtx = trustedRequestCtx
entryReq = entryAuthReq.WithContext(trustedRequestCtx)
reqParams = trustedReqParams
} }
metadata := map[string]string{ metadata := map[string]string{
xhttp.AmzStorageClass: sc, // save same storage-class as incoming stream. xhttp.AmzStorageClass: sc, // save same storage-class as incoming stream.
@@ -2519,7 +2585,7 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
actualSize := size actualSize := size
var idxCb func() []byte var idxCb func() []byte
if isCompressible(r.Header, object) && size > minCompressibleSize { if isCompressible(entryReq.Header, object) && size > minCompressibleSize {
// Storing the compression metadata. // Storing the compression metadata.
metadata[ReservedMetadataPrefix+"compression"] = compressionAlgorithmV2 metadata[ReservedMetadataPrefix+"compression"] = compressionAlgorithmV2
metadata[ReservedMetadataPrefix+"actual-size"] = strconv.FormatInt(size, 10) metadata[ReservedMetadataPrefix+"actual-size"] = strconv.FormatInt(size, 10)
@@ -2530,7 +2596,7 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
} }
// Set compression metrics. // Set compression metrics.
wantEncryption := crypto.Requested(r.Header) wantEncryption := crypto.Requested(entryReq.Header)
s2c, cb := newS2CompressReader(actualReader, actualSize, wantEncryption) s2c, cb := newS2CompressReader(actualReader, actualSize, wantEncryption)
defer s2c.Close() defer s2c.Close()
idxCb = cb idxCb = cb
@@ -2546,15 +2612,11 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
rawReader := hashReader rawReader := hashReader
pReader := NewPutObjReader(rawReader) pReader := NewPutObjReader(rawReader)
allowReplicationMetadata := false if replicaTrusted {
if r.Header.Get(xhttp.AmzBucketReplicationStatus) == replication.Replica.String() { if err := extractReplicationMetadataFromMime(entryCtx, textproto.MIMEHeader(entryReq.Header), metadata); err != nil {
if s3Err = isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.ReplicateObjectAction); s3Err != ErrNone {
return errors.New(errorCodes.ToAPIErr(s3Err).Code)
}
allowReplicationMetadata = true
if err = extractReplicationMetadataFromMime(ctx, textproto.MIMEHeader(r.Header), metadata); err != nil {
return err return err
} }
metadata[xhttp.AmzBucketReplicationStatus] = replication.Replica.String()
metadata[ReservedMetadataPrefixLower+ReplicaStatus] = replication.Replica.String() metadata[ReservedMetadataPrefixLower+ReplicaStatus] = replication.Replica.String()
metadata[ReservedMetadataPrefixLower+ReplicaTimestamp] = UTCNow().Format(time.RFC3339Nano) metadata[ReservedMetadataPrefixLower+ReplicaTimestamp] = UTCNow().Format(time.RFC3339Nano)
} }
@@ -2576,22 +2638,25 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
hdrs.Set(k, v) hdrs.Set(k, v)
} }
} }
m, err := extractMetadata(ctx, textproto.MIMEHeader(hdrs)) if !entryTrusted {
stripReplicationRequestHeaders(hdrs)
}
m, err := extractMetadata(entryCtx, textproto.MIMEHeader(hdrs))
if err != nil { if err != nil {
return err return err
} }
if allowReplicationMetadata { if replicaTrusted {
if err = extractReplicationMetadataFromMime(ctx, textproto.MIMEHeader(hdrs), m); err != nil { if err = extractReplicationMetadataFromMime(entryCtx, textproto.MIMEHeader(hdrs), m); err != nil {
return err return err
} }
} }
maps.Copy(metadata, m) maps.Copy(metadata, m)
} else { } else {
versionID = r.Form.Get(xhttp.VersionID) versionID = entryReq.Form.Get(xhttp.VersionID)
hdrs = r.Header hdrs = entryReq.Header
} }
opts, err := putOpts(ctx, bucket, object, versionID, hdrs, metadata) opts, err := putOpts(entryCtx, bucket, object, versionID, hdrs, metadata, entryTrusted)
if err != nil { if err != nil {
return err return err
} }
@@ -2602,7 +2667,7 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
} }
opts.IndexCB = idxCb opts.IndexCB = idxCb
retentionMode, retentionDate, legalHold, s3err := checkPutObjectLockAllowed(ctx, r, bucket, object, getObjectInfo, retPerms, holdPerms) retentionMode, retentionDate, legalHold, s3err := checkPutObjectLockAllowed(entryCtx, entryReq, bucket, object, getObjectInfo, retPerms, holdPerms, replicaTrusted)
if s3err == ErrNone && retentionMode.Valid() { if s3err == ErrNone && retentionMode.Valid() {
metadata[strings.ToLower(xhttp.AmzObjectLockMode)] = string(retentionMode) metadata[strings.ToLower(xhttp.AmzObjectLockMode)] = string(retentionMode)
metadata[strings.ToLower(xhttp.AmzObjectLockRetainUntilDate)] = amztime.ISO8601Format(retentionDate.UTC()) metadata[strings.ToLower(xhttp.AmzObjectLockRetainUntilDate)] = amztime.ISO8601Format(retentionDate.UTC())
@@ -2613,7 +2678,9 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
} }
if s3err != ErrNone { if s3err != ErrNone {
s3Err = s3err if !ignoreEntryErrors {
setTarS3Err(s3err)
}
return ObjectLocked{} return ObjectLocked{}
} }
@@ -2623,12 +2690,12 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
} }
var objectEncryptionKey crypto.ObjectKey var objectEncryptionKey crypto.ObjectKey
if crypto.Requested(r.Header) { if crypto.Requested(entryReq.Header) {
if crypto.SSECopy.IsRequested(r.Header) { if crypto.SSECopy.IsRequested(entryReq.Header) {
return errInvalidEncryptionParameters return errInvalidEncryptionParameters
} }
reader, objectEncryptionKey, err = EncryptRequest(hashReader, r, bucket, object, metadata) reader, objectEncryptionKey, err = EncryptRequest(hashReader, entryReq, bucket, object, metadata)
if err != nil { if err != nil {
return err return err
} }
@@ -2673,7 +2740,7 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
} }
origETag := objInfo.ETag origETag := objInfo.ETag
objInfo.ETag = getDecryptedETag(r.Header, objInfo, false) objInfo.ETag = getDecryptedETag(entryReq.Header, objInfo, false)
if dsc := mustReplicate(ctx, bucket, object, getMustReplicateOptions(metadata, "", "", replication.ObjectReplicationType, opts)); dsc.ReplicateAny() { if dsc := mustReplicate(ctx, bucket, object, getMustReplicateOptions(metadata, "", "", replication.ObjectReplicationType, opts)); dsc.ReplicateAny() {
scheduleReplication(ctx, objInfo, objectAPI, dsc, replication.ObjectReplicationType) scheduleReplication(ctx, objInfo, objectAPI, dsc, replication.ObjectReplicationType)
@@ -2686,8 +2753,8 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
Object: objInfo, Object: objInfo,
ReqParams: reqParams, ReqParams: reqParams,
RespElements: respElements, RespElements: respElements,
UserAgent: r.UserAgent(), UserAgent: entryReq.UserAgent(),
Host: handlers.GetSourceIP(r), Host: handlers.GetSourceIP(entryReq),
} }
sendEvent(evt) sendEvent(evt)
@@ -2702,7 +2769,14 @@ func (api objectAPIHandlers) PutObjectExtractHandler(w http.ResponseWriter, r *h
return nil return nil
} }
if err = untar(ctx, hreader, putObjectTar, opts); err != nil { err = untar(ctx, hreader, putObjectTar, opts)
if code := APIErrorCode(tarS3Err.Load()); code != ErrNone {
s3Err = code
if err == nil {
err = errors.New(errorCodes.ToAPIErr(code).Code)
}
}
if err != nil {
apiErr := errorCodes.ToAPIErr(s3Err) apiErr := errorCodes.ToAPIErr(s3Err)
// If not set, convert or use BadRequest // If not set, convert or use BadRequest
if s3Err == ErrNone { if s3Err == ErrNone {
@@ -2751,18 +2825,25 @@ func (api objectAPIHandlers) DeleteObjectHandler(w http.ResponseWriter, r *http.
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL)
return return
} }
rawReplica := hasReplicaStatus(r.Header)
markerExact := hasReplicationMarker(r.Header)
replicationPermitted := false
if rawReplica || markerExact {
replicationPermitted = replicationPermissionAllowed(ctx, r, bucket, object, policy.ReplicateDeleteAction)
}
if rawReplica && !replicationPermitted {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrAccessDenied), r.URL)
return
}
if _, ok := r.Header[xhttp.MinIOSourceReplicationCheck]; ok { if _, ok := r.Header[xhttp.MinIOSourceReplicationCheck]; ok {
// requests to just validate replication settings and permissions are not allowed to delete data // requests to just validate replication settings and permissions are not allowed to delete data
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrReplicationPermissionCheckError), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrReplicationPermissionCheckError), r.URL)
return return
} }
trustedReplication := markerExact && replicationPermitted
replica := r.Header.Get(xhttp.AmzBucketReplicationStatus) == replication.Replica.String() replica := trustedReplication && rawReplica
if replica { if hasReplicationRequestHeaders(r.Header) {
if s3Error := checkRequestAuthType(ctx, r, policy.ReplicateDeleteAction, bucket, object); s3Error != ErrNone { ctx, r = applyReplicationTrust(ctx, r, trustedReplication, replica)
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL)
return
}
} }
if globalDNSConfig != nil { if globalDNSConfig != nil {
+64 -15
View File
@@ -171,6 +171,21 @@ func (api objectAPIHandlers) NewMultipartUploadHandler(w http.ResponseWriter, r
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL)
return return
} }
rawReplica := hasReplicaStatus(r.Header)
markerExact := hasReplicationMarker(r.Header)
replicationPermitted := false
if rawReplica || markerExact {
replicationPermitted = replicationPermissionAllowed(ctx, r, bucket, object, policy.ReplicateObjectAction)
}
if rawReplica && !replicationPermitted {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrAccessDenied), r.URL)
return
}
trustedReplication := markerExact && replicationPermitted
replicaTrusted := trustedReplication && rawReplica
if hasReplicationRequestHeaders(r.Header) {
ctx, r = applyReplicationTrust(ctx, r, trustedReplication, replicaTrusted)
}
// Check if bucket encryption is enabled // Check if bucket encryption is enabled
sseConfig, _ := globalBucketSSEConfigSys.Get(bucket) sseConfig, _ := globalBucketSSEConfigSys.Get(bucket)
@@ -205,7 +220,6 @@ func (api objectAPIHandlers) NewMultipartUploadHandler(w http.ResponseWriter, r
return return
} }
_, sourceReplReq := r.Header[xhttp.MinIOSourceReplicationRequest]
ssecRepHeaders := []string{ ssecRepHeaders := []string{
"X-Minio-Replication-Server-Side-Encryption-Seal-Algorithm", "X-Minio-Replication-Server-Side-Encryption-Seal-Algorithm",
"X-Minio-Replication-Server-Side-Encryption-Sealed-Key", "X-Minio-Replication-Server-Side-Encryption-Sealed-Key",
@@ -218,7 +232,7 @@ func (api objectAPIHandlers) NewMultipartUploadHandler(w http.ResponseWriter, r
break break
} }
} }
if !ssecRep || !sourceReplReq { if !ssecRep || !replicaTrusted {
if err = setEncryptionMetadata(r, bucket, object, encMetadata); err != nil { if err = setEncryptionMetadata(r, bucket, object, encMetadata); err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return return
@@ -242,24 +256,23 @@ func (api objectAPIHandlers) NewMultipartUploadHandler(w http.ResponseWriter, r
return return
} }
} }
if r.Header.Get(xhttp.AmzBucketReplicationStatus) == replication.Replica.String() { if replicaTrusted {
if s3Err := isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.ReplicateObjectAction); s3Err != ErrNone {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Err), r.URL)
return
}
if err = extractReplicationMetadataFromMime(ctx, textproto.MIMEHeader(r.Header), metadata); err != nil { if err = extractReplicationMetadataFromMime(ctx, textproto.MIMEHeader(r.Header), metadata); err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return return
} }
metadata[xhttp.AmzBucketReplicationStatus] = replication.Replica.String()
metadata[ReservedMetadataPrefixLower+ReplicaStatus] = replication.Replica.String() metadata[ReservedMetadataPrefixLower+ReplicaStatus] = replication.Replica.String()
metadata[ReservedMetadataPrefixLower+ReplicaTimestamp] = UTCNow().Format(time.RFC3339Nano) metadata[ReservedMetadataPrefixLower+ReplicaTimestamp] = UTCNow().Format(time.RFC3339Nano)
} else {
delete(metadata, xhttp.AmzBucketReplicationStatus)
} }
retPerms := isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.PutObjectRetentionAction) retPerms := isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.PutObjectRetentionAction)
holdPerms := isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.PutObjectLegalHoldAction) holdPerms := isPutActionAllowed(ctx, getRequestAuthType(r), bucket, object, r, policy.PutObjectLegalHoldAction)
getObjectInfo := objectAPI.GetObjectInfo getObjectInfo := objectAPI.GetObjectInfo
retentionMode, retentionDate, legalHold, s3Err := checkPutObjectLockAllowed(ctx, r, bucket, object, getObjectInfo, retPerms, holdPerms) retentionMode, retentionDate, legalHold, s3Err := checkPutObjectLockAllowed(ctx, r, bucket, object, getObjectInfo, retPerms, holdPerms, replicaTrusted)
if s3Err == ErrNone && retentionMode.Valid() { if s3Err == ErrNone && retentionMode.Valid() {
metadata[strings.ToLower(xhttp.AmzObjectLockMode)] = string(retentionMode) metadata[strings.ToLower(xhttp.AmzObjectLockMode)] = string(retentionMode)
metadata[strings.ToLower(xhttp.AmzObjectLockRetainUntilDate)] = amztime.ISO8601Format(retentionDate.UTC()) metadata[strings.ToLower(xhttp.AmzObjectLockRetainUntilDate)] = amztime.ISO8601Format(retentionDate.UTC())
@@ -408,6 +421,14 @@ func (api objectAPIHandlers) CopyObjectPartHandler(w http.ResponseWriter, r *htt
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL)
return return
} }
if hasReplicaStatus(r.Header) &&
!replicationPermissionAllowed(ctx, r, dstBucket, dstObject, policy.ReplicateObjectAction) {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrAccessDenied), r.URL)
return
}
if hasReplicationRequestHeaders(r.Header) {
ctx, r = applyReplicationTrust(ctx, r, false, false)
}
uploadID := r.Form.Get(xhttp.UploadID) uploadID := r.Form.Get(xhttp.UploadID)
partIDString := r.Form.Get(xhttp.PartNumber) partIDString := r.Form.Get(xhttp.PartNumber)
@@ -849,6 +870,22 @@ func (api objectAPIHandlers) PutObjectPartHandler(w http.ResponseWriter, r *http
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL) writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
return return
} }
rawReplica := hasReplicaStatus(r.Header)
markerExact := hasReplicationMarker(r.Header)
replicationPermitted := false
if rawReplica || markerExact {
replicationPermitted = replicationPermissionAllowed(ctx, r, bucket, object, policy.ReplicateObjectAction)
}
if rawReplica && !replicationPermitted {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrAccessDenied), r.URL)
return
}
trustedReplication := markerExact && replicationPermitted
storedReplica := mi.UserDefined[xhttp.AmzBucketReplicationStatus] == replication.Replica.String()
replicaTrusted := trustedReplication && storedReplica
if hasReplicationRequestHeaders(r.Header) {
ctx, r = applyReplicationTrust(ctx, r, trustedReplication, replicaTrusted)
}
// Read compression metadata preserved in the init multipart for the decision. // Read compression metadata preserved in the init multipart for the decision.
_, isCompressed := mi.UserDefined[ReservedMetadataPrefix+"compression"] _, isCompressed := mi.UserDefined[ReservedMetadataPrefix+"compression"]
@@ -917,11 +954,9 @@ func (api objectAPIHandlers) PutObjectPartHandler(w http.ResponseWriter, r *http
pReader.setChecksumReader(checksumReader) pReader.setChecksumReader(checksumReader)
_, isEncrypted := crypto.IsEncrypted(mi.UserDefined) _, isEncrypted := crypto.IsEncrypted(mi.UserDefined)
_, replicationStatus := mi.UserDefined[xhttp.AmzBucketReplicationStatus]
_, sourceReplReq := r.Header[xhttp.MinIOSourceReplicationRequest]
var objectEncryptionKey crypto.ObjectKey var objectEncryptionKey crypto.ObjectKey
if isEncrypted { if isEncrypted {
if !crypto.SSEC.IsRequested(r.Header) && crypto.SSEC.IsEncrypted(mi.UserDefined) && !replicationStatus { if !crypto.SSEC.IsRequested(r.Header) && crypto.SSEC.IsEncrypted(mi.UserDefined) && !replicaTrusted {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrSSEMultipartEncrypted), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrSSEMultipartEncrypted), r.URL)
return return
} }
@@ -941,7 +976,7 @@ func (api objectAPIHandlers) PutObjectPartHandler(w http.ResponseWriter, r *http
} }
} }
if !sourceReplReq || !crypto.SSEC.IsEncrypted(mi.UserDefined) { if !replicaTrusted || !crypto.SSEC.IsEncrypted(mi.UserDefined) {
// Calculating object encryption key // Calculating object encryption key
key, err = decryptObjectMeta(key, bucket, object, mi.UserDefined) key, err = decryptObjectMeta(key, bucket, object, mi.UserDefined)
if err != nil { if err != nil {
@@ -1000,7 +1035,7 @@ func (api objectAPIHandlers) PutObjectPartHandler(w http.ResponseWriter, r *http
} }
opts.IndexCB = idxCb opts.IndexCB = idxCb
opts.ReplicationRequest = sourceReplReq opts.ReplicationRequest = trustedReplication
putObjectPart := objectAPI.PutObjectPart putObjectPart := objectAPI.PutObjectPart
partInfo, err := putObjectPart(ctx, bucket, object, uploadID, partID, pReader, opts) partInfo, err := putObjectPart(ctx, bucket, object, uploadID, partID, pReader, opts)
@@ -1075,6 +1110,20 @@ func (api objectAPIHandlers) CompleteMultipartUploadHandler(w http.ResponseWrite
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Error), r.URL)
return return
} }
rawReplica := hasReplicaStatus(r.Header)
markerExact := hasReplicationMarker(r.Header)
replicationPermitted := false
if rawReplica || markerExact {
replicationPermitted = replicationPermissionAllowed(ctx, r, bucket, object, policy.ReplicateObjectAction)
}
if rawReplica && !replicationPermitted {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrAccessDenied), r.URL)
return
}
trustedReplication := markerExact && replicationPermitted
if hasReplicationRequestHeaders(r.Header) {
ctx, r = applyReplicationTrust(ctx, r, trustedReplication, trustedReplication && rawReplica)
}
// Get upload id. // Get upload id.
uploadID, _, _, _, s3Error := getObjectResources(r.Form) uploadID, _, _, _, s3Error := getObjectResources(r.Form)
@@ -1117,7 +1166,7 @@ func (api objectAPIHandlers) CompleteMultipartUploadHandler(w http.ResponseWrite
return return
} }
if _, _, _, s3Err := checkPutObjectLockAllowed(ctx, r, bucket, object, objectAPI.GetObjectInfo, ErrNone, ErrNone); s3Err != ErrNone { if _, _, _, s3Err := checkPutObjectLockAllowed(ctx, r, bucket, object, objectAPI.GetObjectInfo, ErrNone, ErrNone, false); s3Err != ErrNone {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Err), r.URL) writeErrorResponse(ctx, w, errorCodes.ToAPIErr(s3Err), r.URL)
return return
} }
@@ -1200,7 +1249,7 @@ func (api objectAPIHandlers) CompleteMultipartUploadHandler(w http.ResponseWrite
if dsc := mustReplicate(ctx, bucket, object, objInfo.getMustReplicateOptions(replication.ObjectReplicationType, opts)); dsc.ReplicateAny() { if dsc := mustReplicate(ctx, bucket, object, objInfo.getMustReplicateOptions(replication.ObjectReplicationType, opts)); dsc.ReplicateAny() {
scheduleReplication(ctx, objInfo, objectAPI, dsc, replication.ObjectReplicationType) scheduleReplication(ctx, objInfo, objectAPI, dsc, replication.ObjectReplicationType)
} }
if _, ok := r.Header[xhttp.MinIOSourceReplicationRequest]; ok { if isTrustedReplication(ctx) {
actualSize, _ := objInfo.GetActualSize() actualSize, _ := objInfo.GetActualSize()
defer globalReplicationStats.Load().UpdateReplicaStat(bucket, actualSize) defer globalReplicationStats.Load().UpdateReplicaStat(bucket, actualSize)
} }
+44
View File
@@ -33,6 +33,7 @@ import (
"time" "time"
"github.com/dustin/go-humanize" "github.com/dustin/go-humanize"
xhttp "github.com/minio/minio/internal/http"
) )
const ( const (
@@ -184,6 +185,49 @@ func TestPostPolicyBucketHandler(t *testing.T) {
ExecObjectLayerTest(t, testPostPolicyBucketHandler) ExecObjectLayerTest(t, testPostPolicyBucketHandler)
} }
func TestPostPolicyCannotForgeReplicationStatus(t *testing.T) {
ExecObjectLayerTest(t, testPostPolicyCannotForgeReplicationStatus)
}
func testPostPolicyCannotForgeReplicationStatus(obj ObjectLayer, instanceType string, t TestErrHandler) {
if err := newTestConfig(globalMinioDefaultRegion, obj); err != nil {
t.Fatalf("Initializing config.json failed")
}
bucketName := getRandomBucketName()
if err := obj.MakeBucket(context.Background(), bucketName, MakeBucketOptions{}); err != nil {
t.Fatalf("%s: make bucket: %v", instanceType, err)
}
apiRouter := initTestAPIEndPoints(obj, []string{"PostPolicy"})
credentials := globalActiveCred
now := UTCNow()
region := globalMinioDefaultRegion
objectPrefix := "post-policy-replication-status"
policyBytes := buildGenericPolicy(now, credentials.AccessKey, region, bucketName, objectPrefix, false)
policyText := strings.TrimSuffix(string(policyBytes), "]}") +
`,["eq","$x-amz-replication-status","REPLICA"]]}`
req, err := newPostRequestV4Generic("", bucketName, objectPrefix, []byte("post policy payload"),
credentials.AccessKey, credentials.SecretKey, region, now, []byte(policyText),
map[string]string{xhttp.AmzBucketReplicationStatus: "REPLICA"}, false, false, false)
if err != nil {
t.Fatalf("%s: create post request: %v", instanceType, err)
}
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
if rec.Code != http.StatusNoContent {
t.Fatalf("%s: POST status %d: %s", instanceType, rec.Code, rec.Body.String())
}
info, err := obj.GetObjectInfo(context.Background(), bucketName, objectPrefix+"/upload.txt", ObjectOptions{})
if err != nil {
t.Fatalf("%s: get object info: %v", instanceType, err)
}
if got := info.UserDefined[xhttp.AmzBucketReplicationStatus]; got != "" {
t.Fatalf("%s: forged replication status persisted as %q", instanceType, got)
}
if !info.ReplicationStatus.Empty() {
t.Fatalf("%s: forged replication status reached ObjectInfo: %q", instanceType, info.ReplicationStatus)
}
}
// testPostPolicyBucketHandler - Tests validate post policy handler uploading objects. // testPostPolicyBucketHandler - Tests validate post policy handler uploading objects.
func testPostPolicyBucketHandler(obj ObjectLayer, instanceType string, t TestErrHandler) { func testPostPolicyBucketHandler(obj ObjectLayer, instanceType string, t TestErrHandler) {
if err := newTestConfig(globalMinioDefaultRegion, obj); err != nil { if err := newTestConfig(globalMinioDefaultRegion, obj); err != nil {
+125
View File
@@ -0,0 +1,125 @@
// Copyright (c) 2015-2026 MinIO, Inc.
// Copyright (c) 2026 PGSTY
//
// 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.
package cmd
import (
"context"
"net/http"
objectreplication "github.com/minio/minio/internal/bucket/replication"
xhttp "github.com/minio/minio/internal/http"
"github.com/minio/minio/internal/logger"
"github.com/minio/pkg/v3/policy"
)
type (
replicationTrustKey struct{}
replicaTrustKey struct{}
)
// hasReplicationMarker reports whether the internal replication marker has
// its one accepted wire value. Header presence alone is never a trust signal.
func hasReplicationMarker(h http.Header) bool {
values, ok := h[http.CanonicalHeaderKey(xhttp.MinIOSourceReplicationRequest)]
return ok && len(values) == 1 && values[0] == "true"
}
func hasReplicationMarkerHeader(h http.Header) bool {
_, ok := h[http.CanonicalHeaderKey(xhttp.MinIOSourceReplicationRequest)]
return ok
}
func hasReplicaStatus(h http.Header) bool {
return h.Get(xhttp.AmzBucketReplicationStatus) == objectreplication.Replica.String()
}
func withReplicationTrust(ctx context.Context, trusted, replicaTrusted bool) context.Context {
ctx = context.WithValue(ctx, replicationTrustKey{}, trusted)
return context.WithValue(ctx, replicaTrustKey{}, trusted && replicaTrusted)
}
func isTrustedReplication(ctx context.Context) bool {
trusted, _ := ctx.Value(replicationTrustKey{}).(bool)
return trusted
}
func isReplicaTrusted(ctx context.Context) bool {
trusted, _ := ctx.Value(replicaTrustKey{}).(bool)
return trusted
}
// replicationPermissionAllowed must be called only after the request's
// existing authentication/signature path has succeeded and populated ReqInfo.
// Replication peers are authenticated principals; an anonymous bucket-policy
// grant must not turn client-controlled internal headers into trusted state.
func replicationPermissionAllowed(ctx context.Context, r *http.Request, bucket, object string, action policy.Action) bool {
reqInfo := logger.GetReqInfo(ctx)
if reqInfo == nil || reqInfo.Cred.AccessKey == "" {
return false
}
reqInfo.BucketName = bucket
reqInfo.ObjectName = object
return authorizeRequest(ctx, r, action) == ErrNone
}
// replicationRequestHeaders are internal request controls. They are removed
// only after signature verification when a request has not earned replication
// trust. Public S3/SSE/checksum headers, proxy loop guards, and replication
// validity/readiness probes are intentionally not listed here.
var replicationRequestHeaders = []string{
xhttp.MinIOSourceReplicationRequest,
xhttp.MinIOSourceETag,
xhttp.MinIOSourceMTime,
xhttp.MinIOSourceDeleteMarker,
xhttp.MinIOSourceDeleteMarkerDelete,
xhttp.MinIOSourceTaggingTimestamp,
xhttp.MinIOSourceObjectRetentionTimestamp,
xhttp.MinIOSourceObjectLegalHoldTimestamp,
"X-Minio-Replication-Server-Side-Encryption-Sealed-Key",
"X-Minio-Replication-Server-Side-Encryption-Seal-Algorithm",
"X-Minio-Replication-Server-Side-Encryption-Iv",
"X-Minio-Replication-Encrypted-Multipart",
xhttp.MinIOReplicationActualObjectSize,
ReplicationSsecChecksumHeader,
xhttp.AmzBucketReplicationStatus,
}
func stripReplicationRequestHeaders(h http.Header) {
for _, name := range replicationRequestHeaders {
h.Del(name)
}
}
func hasReplicationRequestHeaders(h http.Header) bool {
for _, name := range replicationRequestHeaders {
if _, ok := h[http.CanonicalHeaderKey(name)]; ok {
return true
}
}
return false
}
func cloneRequestWithoutReplicationHeaders(ctx context.Context, r *http.Request) *http.Request {
clone := r.Clone(ctx)
stripReplicationRequestHeaders(clone.Header)
return clone
}
// applyReplicationTrust binds the handler context to the effective request.
// The context marker is the authorization source of truth; header removal is
// defense in depth for option builders and future call sites.
func applyReplicationTrust(ctx context.Context, r *http.Request, trusted, replicaTrusted bool) (context.Context, *http.Request) {
ctx = withReplicationTrust(ctx, trusted, replicaTrusted)
if trusted {
return ctx, r.WithContext(ctx)
}
return ctx, cloneRequestWithoutReplicationHeaders(ctx, r)
}
+832
View File
@@ -0,0 +1,832 @@
// Copyright (c) 2015-2026 MinIO, Inc.
// Copyright (c) 2026 PGSTY
//
// 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.
package cmd
import (
"archive/tar"
"bytes"
"context"
"crypto/md5"
"encoding/base64"
"encoding/xml"
"io"
"net/http"
"net/http/httptest"
"net/url"
"strconv"
"strings"
"testing"
"time"
"github.com/minio/madmin-go/v3"
"github.com/minio/minio/internal/auth"
"github.com/minio/minio/internal/crypto"
xhttp "github.com/minio/minio/internal/http"
"github.com/minio/minio/internal/kms"
"github.com/minio/pkg/v3/policy"
)
func TestAPIReplicationTrustProtectsSSECReads(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPIReplicationTrustProtectsSSECReads,
})
}
func testAPIReplicationTrustProtectsSSECReads(_ ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, credentials auth.Credentials, t *testing.T,
) {
previousTLS := globalIsTLS
globalIsTLS = true
defer func() { globalIsTLS = previousTLS }()
key := bytes.Repeat([]byte{0x31}, 32)
keyMD5 := md5.Sum(key)
data := bytes.Repeat([]byte("replication-trust-ssec-"), 256)
sseHeaders := map[string]string{
xhttp.AmzServerSideEncryptionCustomerAlgorithm: xhttp.AmzEncryptionAES,
xhttp.AmzServerSideEncryptionCustomerKey: base64.StdEncoding.EncodeToString(key),
xhttp.AmzServerSideEncryptionCustomerKeyMD5: base64.StdEncoding.EncodeToString(keyMD5[:]),
}
object := "replication-trust/ssec-read"
putCopyChecksumSource(t, apiRouter, credentials, bucketName, object, data, sseHeaders)
readerOnly := newObjectAttributesAuthzUser(t, instanceType, bucketName, `"s3:GetObject"`)
replicator := newObjectAttributesAuthzUser(t, instanceType, bucketName, `"s3:GetObject","s3:ReplicateObject"`)
marker := map[string]string{xhttp.MinIOSourceReplicationRequest: "true"}
wrongCaseMarker := map[string]string{xhttp.MinIOSourceReplicationRequest: "TRUE"}
conditionalMarker := map[string]string{
xhttp.MinIOSourceReplicationRequest: "true",
xhttp.IfNoneMatch: "*",
}
for _, test := range []struct {
name string
method string
creds auth.Credentials
headers map[string]string
wantStatus int
wantPlain bool
wantCipher bool
}{
{name: "get/reader/fake-marker", method: http.MethodGet, creds: readerOnly, headers: marker, wantStatus: http.StatusBadRequest},
{name: "get/replicator/trusted", method: http.MethodGet, creds: replicator, headers: marker, wantStatus: http.StatusOK, wantCipher: true},
{name: "get/root/trusted", method: http.MethodGet, creds: credentials, headers: marker, wantStatus: http.StatusOK, wantCipher: true},
{name: "get/replicator/wrong-case", method: http.MethodGet, creds: replicator, headers: wrongCaseMarker, wantStatus: http.StatusBadRequest},
{name: "get/reader/key", method: http.MethodGet, creds: readerOnly, headers: sseHeaders, wantStatus: http.StatusOK, wantPlain: true},
{name: "head/reader/fake-marker", method: http.MethodHead, creds: readerOnly, headers: marker, wantStatus: http.StatusBadRequest},
{name: "head/reader/conditional-oracle", method: http.MethodHead, creds: readerOnly, headers: conditionalMarker, wantStatus: http.StatusBadRequest},
{name: "head/replicator/trusted", method: http.MethodHead, creds: replicator, headers: marker, wantStatus: http.StatusOK},
{name: "head/replicator/wrong-case", method: http.MethodHead, creds: replicator, headers: wrongCaseMarker, wantStatus: http.StatusBadRequest},
} {
t.Run(test.name, func(t *testing.T) {
req, err := newTestSignedRequestV4(test.method, getGetObjectURL("", bucketName, object), 0, nil,
test.creds.AccessKey, test.creds.SecretKey, test.headers)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
if rec.Code != test.wantStatus {
t.Fatalf("%s: status %d, want %d: %s", instanceType, rec.Code, test.wantStatus, rec.Body.String())
}
if test.wantPlain && !bytes.Equal(rec.Body.Bytes(), data) {
t.Fatal("ordinary SSE-C GET did not return plaintext")
}
if test.wantCipher && (len(rec.Body.Bytes()) == 0 || bytes.Equal(rec.Body.Bytes(), data)) {
t.Fatal("trusted replication GET did not return ciphertext")
}
})
}
}
func TestReplicationTrustControlsInternalOptionsAndEvents(t *testing.T) {
mtime := time.Date(2026, 8, 31, 12, 34, 56, 123, time.UTC)
headers := make(http.Header)
headers.Set(xhttp.MinIOSourceReplicationRequest, "true")
headers.Set(xhttp.MinIOSourceETag, "source-etag")
headers.Set(xhttp.MinIOSourceMTime, mtime.Format(time.RFC3339Nano))
headers.Set(xhttp.MinIOReplicationActualObjectSize, "123")
headers.Set(ReplicationSsecChecksumHeader, "checksum")
ordinary, err := putOptsFromHeaders(t.Context(), headers, nil, false)
if err != nil {
t.Fatal(err)
}
if ordinary.ReplicationRequest || ordinary.PreserveETag != "" || !ordinary.MTime.IsZero() {
t.Fatalf("ordinary options trusted internal headers: %#v", ordinary)
}
trusted, err := putOptsFromHeaders(t.Context(), headers, nil, true)
if err != nil {
t.Fatal(err)
}
if !trusted.ReplicationRequest || trusted.PreserveETag != "source-etag" || !trusted.MTime.Equal(mtime) {
t.Fatalf("trusted options lost source state: %#v", trusted)
}
completeReq := &http.Request{Header: headers.Clone(), Form: make(url.Values)}
ordinaryComplete, err := completeMultipartOpts(t.Context(), completeReq, "bucket", "object")
if err != nil {
t.Fatal(err)
}
if ordinaryComplete.ReplicationRequest || len(ordinaryComplete.UserDefined) != 0 {
t.Fatalf("ordinary completion trusted internal metadata: %#v", ordinaryComplete)
}
trustedCompleteCtx := withReplicationTrust(t.Context(), true, false)
trustedComplete, err := completeMultipartOpts(trustedCompleteCtx, completeReq.WithContext(trustedCompleteCtx), "bucket", "object")
if err != nil {
t.Fatal(err)
}
if !trustedComplete.ReplicationRequest || trustedComplete.UserDefined[ReservedMetadataPrefix+"Actual-Object-Size"] != "123" ||
trustedComplete.UserDefined[ReplicationSsecChecksumHeader] != "checksum" {
t.Fatalf("trusted completion lost internal metadata: %#v", trustedComplete)
}
req := &http.Request{Header: headers.Clone(), Form: make(url.Values)}
req = req.WithContext(context.Background())
if _, ok := extractReqParams(req)[xhttp.MinIOSourceReplicationRequest]; ok {
t.Fatal("untrusted marker suppressed events")
}
trustedCtx := withReplicationTrust(req.Context(), true, false)
req = req.WithContext(trustedCtx)
if _, ok := extractReqParams(req)[xhttp.MinIOSourceReplicationRequest]; !ok {
t.Fatal("trusted replication marker was not propagated to events")
}
}
func TestAPIPutObjectReplicationTrust(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPIPutObjectReplicationTrust,
})
}
func testAPIPutObjectReplicationTrust(obj ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, _ auth.Credentials, t *testing.T,
) {
putOnly := newObjectAttributesAuthzUser(t, instanceType, bucketName, `"s3:PutObject"`)
replicator := newObjectAttributesAuthzUser(t, instanceType, bucketName, `"s3:PutObject","s3:ReplicateObject"`)
payload := []byte("replication trust put payload")
sourceMTime := time.Date(2024, 1, 2, 3, 4, 5, 6, time.UTC)
request := func(t *testing.T, object string, creds auth.Credentials, status string, check bool) *httptest.ResponseRecorder {
t.Helper()
headers := map[string]string{
xhttp.MinIOSourceReplicationRequest: "true",
xhttp.MinIOSourceETag: "source-etag",
xhttp.MinIOSourceMTime: sourceMTime.Format(time.RFC3339Nano),
}
if status != "" {
headers[xhttp.AmzBucketReplicationStatus] = status
}
if check {
headers[xhttp.MinIOSourceReplicationCheck] = "true"
}
req, err := newTestSignedRequestV4(http.MethodPut, getPutObjectURL("", bucketName, object),
int64(len(payload)), bytes.NewReader(payload), creds.AccessKey, creds.SecretKey, headers)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
return rec
}
t.Run("untrusted marker is ordinary", func(t *testing.T) {
object := "replication-trust/put-ordinary"
if rec := request(t, object, putOnly, "PENDING", false); rec.Code != http.StatusOK {
t.Fatalf("status %d: %s", rec.Code, rec.Body.String())
}
info, err := obj.GetObjectInfo(t.Context(), bucketName, object, ObjectOptions{})
if err != nil {
t.Fatal(err)
}
if info.ETag == "source-etag" || info.ModTime.Equal(sourceMTime) {
t.Fatalf("untrusted source state was preserved: ETag=%q MTime=%v", info.ETag, info.ModTime)
}
assertObjectMetadataKeysAbsent(t, info.UserDefined, xhttp.AmzBucketReplicationStatus)
})
t.Run("unauthorized replica is denied", func(t *testing.T) {
object := "replication-trust/put-denied-replica"
if rec := request(t, object, putOnly, "REPLICA", false); rec.Code != http.StatusForbidden {
t.Fatalf("status %d, want 403: %s", rec.Code, rec.Body.String())
}
if _, err := obj.GetObjectInfo(t.Context(), bucketName, object, ObjectOptions{}); err == nil {
t.Fatal("unauthorized replica write created an object")
}
})
t.Run("trusted batch preserves source state", func(t *testing.T) {
object := "replication-trust/put-batch"
if rec := request(t, object, replicator, "", false); rec.Code != http.StatusOK {
t.Fatalf("status %d: %s", rec.Code, rec.Body.String())
}
info, err := obj.GetObjectInfo(t.Context(), bucketName, object, ObjectOptions{})
if err != nil {
t.Fatal(err)
}
if info.ETag != "source-etag" || !info.ModTime.Equal(sourceMTime) {
t.Fatalf("trusted source state lost: ETag=%q MTime=%v", info.ETag, info.ModTime)
}
assertObjectMetadataKeysAbsent(t, info.UserDefined, xhttp.AmzBucketReplicationStatus)
})
t.Run("trusted replica persists replica state", func(t *testing.T) {
object := "replication-trust/put-replica"
if rec := request(t, object, replicator, "REPLICA", false); rec.Code != http.StatusOK {
t.Fatalf("status %d: %s", rec.Code, rec.Body.String())
}
info, err := obj.GetObjectInfo(t.Context(), bucketName, object, ObjectOptions{})
if err != nil {
t.Fatal(err)
}
if info.UserDefined[xhttp.AmzBucketReplicationStatus] != "REPLICA" {
t.Fatalf("replica status not persisted: %#v", info.UserDefined)
}
})
for _, test := range []struct {
name string
creds auth.Credentials
wantStatus int
}{
{name: "validity check requires ReplicateObject", creds: putOnly, wantStatus: http.StatusForbidden},
{name: "validity check succeeds for replicator", creds: replicator, wantStatus: http.StatusBadRequest},
} {
t.Run(test.name, func(t *testing.T) {
object := "replication-trust/put-check-" + strconv.Itoa(test.wantStatus)
if rec := request(t, object, test.creds, "REPLICA", true); rec.Code != test.wantStatus {
t.Fatalf("status %d, want %d: %s", rec.Code, test.wantStatus, rec.Body.String())
}
if _, err := obj.GetObjectInfo(t.Context(), bucketName, object, ObjectOptions{}); err == nil {
t.Fatal("replication validity check created an object")
}
})
}
}
func TestAPISnowballReplicationTrustIsPerEntry(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPISnowballReplicationTrustIsPerEntry,
})
}
func testAPISnowballReplicationTrustIsPerEntry(obj ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, _ auth.Credentials, t *testing.T,
) {
const (
allowedPrefix = "snowball/allowed/"
deniedPrefix = "snowball/denied/"
sourceETag = "0123456789abcdef0123456789abcdef"
)
creds := newSnowballReplicationTrustUser(t, instanceType, bucketName, allowedPrefix)
var body bytes.Buffer
tw := tar.NewWriter(&body)
objects := make([]struct {
name string
trusted bool
}, 0, 32)
for i := 0; i < 16; i++ {
for _, entry := range []struct {
prefix string
trusted bool
}{
{prefix: allowedPrefix, trusted: true},
{prefix: deniedPrefix},
} {
name := entry.prefix + strconv.Itoa(i)
data := []byte("snowball replication trust " + name)
if err := tw.WriteHeader(&tar.Header{Name: name, Mode: 0o600, Size: int64(len(data))}); err != nil {
t.Fatal(err)
}
if _, err := tw.Write(data); err != nil {
t.Fatal(err)
}
objects = append(objects, struct {
name string
trusted bool
}{name: name, trusted: entry.trusted})
}
}
if err := tw.Close(); err != nil {
t.Fatal(err)
}
headers := map[string]string{
xhttp.AmzSnowballExtract: "true",
xhttp.MinIOSourceReplicationRequest: "true",
xhttp.MinIOSourceETag: sourceETag,
}
for _, test := range []struct {
name string
trailer bool
}{
{name: "signed-v4"},
{name: "streaming-unsigned-trailer", trailer: true},
} {
t.Run(test.name, func(t *testing.T) {
var req *http.Request
var err error
if test.trailer {
req, err = newStreamingUnsignedTrailerRequest(http.MethodPut,
getPutObjectURL("", bucketName, "snowball.tar"), body.Bytes(), UTCNow())
if err == nil {
for name, value := range headers {
req.Header.Set(name, value)
}
err = signRequestV4(req, creds.AccessKey, creds.SecretKey)
}
} else {
req, err = newTestSignedRequestV4(http.MethodPut, getPutObjectURL("", bucketName, "snowball.tar"),
int64(body.Len()), bytes.NewReader(body.Bytes()), creds.AccessKey, creds.SecretKey, headers)
}
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("%s: Snowball PUT status %d: %s", instanceType, rec.Code, rec.Body.String())
}
for _, object := range objects {
info, err := obj.GetObjectInfo(t.Context(), bucketName, object.name, ObjectOptions{})
if err != nil {
t.Fatalf("%s: get %s: %v", instanceType, object.name, err)
}
if object.trusted && info.ETag != sourceETag {
t.Errorf("%s: trusted entry %s ETag = %q, want source ETag", instanceType, object.name, info.ETag)
}
if !object.trusted && info.ETag == sourceETag {
t.Errorf("%s: untrusted entry %s preserved source ETag", instanceType, object.name)
}
}
})
}
}
func TestAPISnowballInheritsBucketEncryption(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPISnowballInheritsBucketEncryption,
})
}
func testAPISnowballInheritsBucketEncryption(obj ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, credentials auth.Credentials, t *testing.T,
) {
previousKMS := GlobalKMS
GlobalKMS = kms.NewStub("snowball-default-encryption")
defer func() { GlobalKMS = previousKMS }()
sseXML := []byte(`<ServerSideEncryptionConfiguration xmlns="http://s3.amazonaws.com/doc/2006-03-01/"><Rule><ApplyServerSideEncryptionByDefault><SSEAlgorithm>AES256</SSEAlgorithm></ApplyServerSideEncryptionByDefault></Rule></ServerSideEncryptionConfiguration>`)
if _, err := globalBucketMetadataSys.Update(t.Context(), bucketName, bucketSSEConfig, sseXML); err != nil {
t.Fatalf("%s: configure bucket encryption: %v", instanceType, err)
}
var body bytes.Buffer
tw := tar.NewWriter(&body)
objects := []string{"encrypted/one", "encrypted/two"}
for _, object := range objects {
data := []byte("snowball bucket encryption " + object)
if err := tw.WriteHeader(&tar.Header{Name: object, Mode: 0o600, Size: int64(len(data))}); err != nil {
t.Fatal(err)
}
if _, err := tw.Write(data); err != nil {
t.Fatal(err)
}
}
if err := tw.Close(); err != nil {
t.Fatal(err)
}
req, err := newTestSignedRequestV4(http.MethodPut, getPutObjectURL("", bucketName, "encrypted-snowball.tar"),
int64(body.Len()), bytes.NewReader(body.Bytes()), credentials.AccessKey, credentials.SecretKey,
map[string]string{xhttp.AmzSnowballExtract: "true"})
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("%s: Snowball PUT status %d: %s", instanceType, rec.Code, rec.Body.String())
}
for _, object := range objects {
info, err := obj.GetObjectInfo(t.Context(), bucketName, object, ObjectOptions{})
if err != nil {
t.Fatalf("%s: get %s: %v", instanceType, object, err)
}
if _, encrypted := crypto.IsEncrypted(info.UserDefined); !encrypted {
t.Errorf("%s: extracted entry %s did not inherit bucket encryption", instanceType, object)
}
}
}
func newSnowballReplicationTrustUser(t *testing.T, instanceType, bucketName, allowedPrefix string) auth.Credentials {
t.Helper()
ctx := t.Context()
accessKey, secretKey, err := auth.GenerateCredentials()
if err != nil {
t.Fatalf("%s: generate credentials: %v", instanceType, err)
}
creds := auth.Credentials{AccessKey: accessKey, SecretKey: secretKey}
if _, err = globalIAMSys.CreateUser(ctx, creds.AccessKey, madmin.AddOrUpdateUserReq{
SecretKey: creds.SecretKey,
Status: madmin.AccountEnabled,
}); err != nil {
t.Fatalf("%s: create Snowball user: %v", instanceType, err)
}
policyJSON := `{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": ["s3:PutObject"],
"Resource": ["arn:aws:s3:::` + bucketName + `/*"]
},
{
"Effect": "Allow",
"Action": ["s3:ReplicateObject"],
"Resource": ["arn:aws:s3:::` + bucketName + `/` + allowedPrefix + `*"]
}
]
}`
parsed, err := policy.ParseConfig(strings.NewReader(policyJSON))
if err != nil {
t.Fatalf("%s: parse Snowball policy: %v", instanceType, err)
}
policyName := "snowball-replication-trust-" + mustGetUUID()
if _, err = globalIAMSys.SetPolicy(ctx, policyName, *parsed); err != nil {
t.Fatalf("%s: install Snowball policy: %v", instanceType, err)
}
if _, err = globalIAMSys.PolicyDBSet(ctx, creds.AccessKey, policyName, regUser, false); err != nil {
t.Fatalf("%s: attach Snowball policy: %v", instanceType, err)
}
return creds
}
func TestAPICopyObjectMarkerOnlyDoesNotCopyCiphertext(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPICopyObjectMarkerOnlyDoesNotCopyCiphertext,
})
}
func testAPICopyObjectMarkerOnlyDoesNotCopyCiphertext(obj ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, credentials auth.Credentials, t *testing.T,
) {
previousTLS := globalIsTLS
globalIsTLS = true
defer func() { globalIsTLS = previousTLS }()
key := bytes.Repeat([]byte{0x57}, 32)
keyMD5 := md5.Sum(key)
data := bytes.Repeat([]byte("copy marker-only plaintext "), 256)
sseHeaders := map[string]string{
xhttp.AmzServerSideEncryptionCustomerAlgorithm: xhttp.AmzEncryptionAES,
xhttp.AmzServerSideEncryptionCustomerKey: base64.StdEncoding.EncodeToString(key),
xhttp.AmzServerSideEncryptionCustomerKeyMD5: base64.StdEncoding.EncodeToString(keyMD5[:]),
}
srcObject := "replication-trust/copy-ssec-source"
dstObject := "replication-trust/copy-marker-only"
putCopyChecksumSource(t, apiRouter, credentials, bucketName, srcObject, data, sseHeaders)
replicator := newObjectAttributesAuthzUser(t, instanceType, bucketName,
`"s3:GetObject","s3:PutObject","s3:ReplicateObject"`)
headers := map[string]string{
xhttp.AmzCopySource: url.QueryEscape(SlashSeparator + bucketName + SlashSeparator + srcObject),
xhttp.MinIOSourceReplicationRequest: "true",
xhttp.AmzServerSideEncryptionCopyCustomerAlgorithm: xhttp.AmzEncryptionAES,
xhttp.AmzServerSideEncryptionCopyCustomerKey: base64.StdEncoding.EncodeToString(key),
xhttp.AmzServerSideEncryptionCopyCustomerKeyMD5: base64.StdEncoding.EncodeToString(keyMD5[:]),
}
req, err := newTestSignedRequestV4(http.MethodPut, getCopyObjectURL("", bucketName, dstObject), 0, nil,
replicator.AccessKey, replicator.SecretKey, headers)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("%s: CopyObject status %d: %s", instanceType, rec.Code, rec.Body.String())
}
assertObjectContents(t, obj, bucketName, dstObject, data)
}
func TestAPIDeleteObjectReplicationTrust(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPIDeleteObjectReplicationTrust,
})
}
func testAPIDeleteObjectReplicationTrust(obj ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, _ auth.Credentials, t *testing.T,
) {
deleteOnly := newObjectAttributesAuthzUser(t, instanceType, bucketName, `"s3:DeleteObject"`)
replicator := newObjectAttributesAuthzUser(t, instanceType, bucketName, `"s3:DeleteObject","s3:ReplicateDelete"`)
payload := []byte("delete replication trust")
put := func(t *testing.T, object string) {
t.Helper()
if _, err := obj.PutObject(t.Context(), bucketName, object,
mustGetPutObjReader(t, bytes.NewReader(payload), int64(len(payload)), "", ""), ObjectOptions{}); err != nil {
t.Fatal(err)
}
}
remove := func(t *testing.T, object string, creds auth.Credentials, versionID string, deleteMarker, check bool) *httptest.ResponseRecorder {
t.Helper()
headers := map[string]string{
xhttp.MinIOSourceReplicationRequest: "true",
xhttp.AmzBucketReplicationStatus: "REPLICA",
}
if deleteMarker {
headers[xhttp.MinIOSourceDeleteMarker] = "true"
}
if check {
headers[xhttp.MinIOSourceReplicationCheck] = "true"
}
target := getDeleteObjectURL("", bucketName, object)
if versionID != "" {
target += "?" + url.Values{xhttp.VersionID: {versionID}}.Encode()
}
req, err := newTestSignedRequestV4(http.MethodDelete, target,
0, nil, creds.AccessKey, creds.SecretKey, headers)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
apiRouter.ServeHTTP(rec, req)
return rec
}
t.Run("replica status without ReplicateDelete is denied", func(t *testing.T) {
object := "replication-trust/delete-denied"
put(t, object)
if rec := remove(t, object, deleteOnly, "", false, false); rec.Code != http.StatusForbidden {
t.Fatalf("status %d, want 403: %s", rec.Code, rec.Body.String())
}
if _, err := obj.GetObjectInfo(t.Context(), bucketName, object, ObjectOptions{}); err != nil {
t.Fatalf("denied delete removed object: %v", err)
}
})
t.Run("trusted replica delete remains supported", func(t *testing.T) {
object := "replication-trust/delete-allowed"
put(t, object)
if rec := remove(t, object, replicator, "", false, false); rec.Code != http.StatusNoContent {
t.Fatalf("status %d, want 204: %s", rec.Code, rec.Body.String())
}
})
for _, shape := range []struct {
name string
deleteMarker bool
}{
{name: "delete-marker", deleteMarker: true},
{name: "version-purge"},
} {
t.Run("validity check/"+shape.name, func(t *testing.T) {
for _, test := range []struct {
name string
creds auth.Credentials
wantStatus int
}{
{name: "requires ReplicateDelete", creds: deleteOnly, wantStatus: http.StatusForbidden},
{name: "succeeds for replicator", creds: replicator, wantStatus: http.StatusBadRequest},
} {
t.Run(test.name, func(t *testing.T) {
object := "replication-trust/delete-check-" + shape.name + "-" + strconv.Itoa(test.wantStatus)
put(t, object)
if rec := remove(t, object, test.creds, mustGetUUID(), shape.deleteMarker, true); rec.Code != test.wantStatus {
t.Fatalf("status %d, want %d: %s", rec.Code, test.wantStatus, rec.Body.String())
}
if _, err := obj.GetObjectInfo(t.Context(), bucketName, object, ObjectOptions{}); err != nil {
t.Fatalf("replication validity check removed object: %v", err)
}
})
}
})
}
}
func TestAPISSECMultipartReplicationTrust(t *testing.T) {
defer DetectTestLeak(t)()
ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{
t: t,
objAPITest: testAPISSECMultipartReplicationTrust,
})
}
func testAPISSECMultipartReplicationTrust(obj ObjectLayer, instanceType, bucketName string,
apiRouter http.Handler, credentials auth.Credentials, t *testing.T,
) {
previousTLS := globalIsTLS
globalIsTLS = true
defer func() { globalIsTLS = previousTLS }()
replicator := newObjectAttributesAuthzUser(t, instanceType, bucketName, `"s3:PutObject","s3:GetObject","s3:ReplicateObject"`)
putOnly := newObjectAttributesAuthzUser(t, instanceType, bucketName, `"s3:PutObject"`)
key := bytes.Repeat([]byte{0x42}, 32)
keyMD5 := md5.Sum(key)
data := bytes.Repeat([]byte("trusted multipart replication "), 4096)
sseHeaders := map[string]string{
xhttp.AmzServerSideEncryptionCustomerAlgorithm: xhttp.AmzEncryptionAES,
xhttp.AmzServerSideEncryptionCustomerKey: base64.StdEncoding.EncodeToString(key),
xhttp.AmzServerSideEncryptionCustomerKeyMD5: base64.StdEncoding.EncodeToString(keyMD5[:]),
}
// A marker alone must not let an ordinary writer upload raw bytes into an
// SSE-C multipart upload without presenting the customer key.
fakeObject := "replication-trust/ssec-multipart-fake"
fakeNewReq, err := newTestSignedRequestV4(http.MethodPost, getNewMultipartURL("", bucketName, fakeObject),
0, nil, credentials.AccessKey, credentials.SecretKey, sseHeaders)
if err != nil {
t.Fatal(err)
}
fakeNewRec := httptest.NewRecorder()
apiRouter.ServeHTTP(fakeNewRec, fakeNewReq)
if fakeNewRec.Code != http.StatusOK {
t.Fatalf("fake-path NewMultipart status %d: %s", fakeNewRec.Code, fakeNewRec.Body.String())
}
var fakeInit InitiateMultipartUploadResponse
if err = xmlDecoder(fakeNewRec.Body, &fakeInit, int64(fakeNewRec.Body.Len())); err != nil {
t.Fatal(err)
}
fakePartReq, err := newTestSignedRequestV4(http.MethodPut,
getPutObjectPartURL("", bucketName, fakeObject, fakeInit.UploadID, "1"), int64(len(data)), bytes.NewReader(data),
putOnly.AccessKey, putOnly.SecretKey, map[string]string{xhttp.MinIOSourceReplicationRequest: "true"})
if err != nil {
t.Fatal(err)
}
fakePartRec := httptest.NewRecorder()
apiRouter.ServeHTTP(fakePartRec, fakePartReq)
if fakePartRec.Code != http.StatusBadRequest {
t.Fatalf("fake marker PutPart status %d, want 400: %s", fakePartRec.Code, fakePartRec.Body.String())
}
object := "replication-trust/ssec-multipart"
// Create the source as a real SSE-C multipart object so the encrypted part
// layout and metadata match what the replication worker reads.
newReq, err := newTestSignedRequestV4(http.MethodPost, getNewMultipartURL("", bucketName, object),
0, nil, credentials.AccessKey, credentials.SecretKey, sseHeaders)
if err != nil {
t.Fatal(err)
}
newRec := httptest.NewRecorder()
apiRouter.ServeHTTP(newRec, newReq)
if newRec.Code != http.StatusOK {
t.Fatalf("source NewMultipart status %d: %s", newRec.Code, newRec.Body.String())
}
var sourceInit InitiateMultipartUploadResponse
if err = xmlDecoder(newRec.Body, &sourceInit, int64(newRec.Body.Len())); err != nil {
t.Fatal(err)
}
partReq, err := newTestSignedRequestV4(http.MethodPut,
getPutObjectPartURL("", bucketName, object, sourceInit.UploadID, "1"), int64(len(data)), bytes.NewReader(data),
credentials.AccessKey, credentials.SecretKey, sseHeaders)
if err != nil {
t.Fatal(err)
}
partRec := httptest.NewRecorder()
apiRouter.ServeHTTP(partRec, partReq)
if partRec.Code != http.StatusOK {
t.Fatalf("source PutPart status %d: %s", partRec.Code, partRec.Body.String())
}
sourcePartETag := canonicalizeETag(partRec.Header()[xhttp.ETag][0])
sourceCompleteBody, err := xml.Marshal(CompleteMultipartUpload{Parts: []CompletePart{{PartNumber: 1, ETag: sourcePartETag}}})
if err != nil {
t.Fatal(err)
}
completeReq, err := newTestSignedRequestV4(http.MethodPost,
getCompleteMultipartUploadURL("", bucketName, object, sourceInit.UploadID), int64(len(sourceCompleteBody)),
bytes.NewReader(sourceCompleteBody), credentials.AccessKey, credentials.SecretKey, sseHeaders)
if err != nil {
t.Fatal(err)
}
completeRec := httptest.NewRecorder()
apiRouter.ServeHTTP(completeRec, completeReq)
if completeRec.Code != http.StatusOK {
t.Fatalf("source Complete status %d: %s", completeRec.Code, completeRec.Body.String())
}
gr, err := obj.GetObjectNInfo(t.Context(), bucketName, object, nil, http.Header{}, ObjectOptions{ReplicationRequest: true})
if err != nil {
t.Fatal(err)
}
sourceInfo := gr.ObjInfo
rawPart, err := io.ReadAll(gr)
gr.Close()
if err != nil {
t.Fatal(err)
}
if len(rawPart) == 0 || bytes.Equal(rawPart, data) {
t.Fatal("source replication read did not return encrypted bytes")
}
replicationOpts, isMP, err := putReplicationOpts(t.Context(), "", sourceInfo)
if err != nil {
t.Fatal(err)
}
if !isMP {
t.Fatal("SSE-C multipart source was not recognized as multipart")
}
replicationOpts.Internal.SourceMTime = time.Time{}
replicationHeaders := make(map[string]string)
for name, values := range replicationOpts.Header() {
if len(values) > 0 {
replicationHeaders[name] = values[0]
}
}
// Start the destination upload over the same key. The existing object stays
// readable until Complete, so buffering rawPart above mirrors a remote peer.
replNewReq, err := newTestSignedRequestV4(http.MethodPost, getNewMultipartURL("", bucketName, object),
0, nil, replicator.AccessKey, replicator.SecretKey, replicationHeaders)
if err != nil {
t.Fatal(err)
}
replNewRec := httptest.NewRecorder()
apiRouter.ServeHTTP(replNewRec, replNewReq)
if replNewRec.Code != http.StatusOK {
t.Fatalf("replica NewMultipart status %d: %s", replNewRec.Code, replNewRec.Body.String())
}
var replicaInit InitiateMultipartUploadResponse
if err = xmlDecoder(replNewRec.Body, &replicaInit, int64(replNewRec.Body.Len())); err != nil {
t.Fatal(err)
}
replPartHeaders := map[string]string{xhttp.MinIOSourceReplicationRequest: "true"}
replPartReq, err := newTestSignedRequestV4(http.MethodPut,
getPutObjectPartURL("", bucketName, object, replicaInit.UploadID, "1"), int64(len(rawPart)), bytes.NewReader(rawPart),
replicator.AccessKey, replicator.SecretKey, replPartHeaders)
if err != nil {
t.Fatal(err)
}
replPartRec := httptest.NewRecorder()
apiRouter.ServeHTTP(replPartRec, replPartReq)
if replPartRec.Code != http.StatusOK {
t.Fatalf("replica PutPart status %d: %s", replPartRec.Code, replPartRec.Body.String())
}
replPartETag := canonicalizeETag(replPartRec.Header()[xhttp.ETag][0])
replCompleteBody, err := xml.Marshal(CompleteMultipartUpload{Parts: []CompletePart{{PartNumber: 1, ETag: replPartETag}}})
if err != nil {
t.Fatal(err)
}
actualSize, err := sourceInfo.GetActualSize()
if err != nil {
t.Fatal(err)
}
replCompleteHeaders := map[string]string{
xhttp.MinIOSourceReplicationRequest: "true",
xhttp.MinIOSourceMTime: sourceInfo.ModTime.Format(time.RFC3339Nano),
xhttp.MinIOSourceETag: sourceInfo.ETag,
xhttp.MinIOReplicationActualObjectSize: strconv.FormatInt(actualSize, 10),
}
replCompleteReq, err := newTestSignedRequestV4(http.MethodPost,
getCompleteMultipartUploadURL("", bucketName, object, replicaInit.UploadID), int64(len(replCompleteBody)),
bytes.NewReader(replCompleteBody), replicator.AccessKey, replicator.SecretKey, replCompleteHeaders)
if err != nil {
t.Fatal(err)
}
replCompleteRec := httptest.NewRecorder()
apiRouter.ServeHTTP(replCompleteRec, replCompleteReq)
if replCompleteRec.Code != http.StatusOK {
t.Fatalf("replica Complete status %d: %s", replCompleteRec.Code, replCompleteRec.Body.String())
}
getReq, err := newTestSignedRequestV4(http.MethodGet, getGetObjectURL("", bucketName, object),
0, nil, credentials.AccessKey, credentials.SecretKey, sseHeaders)
if err != nil {
t.Fatal(err)
}
getRec := httptest.NewRecorder()
apiRouter.ServeHTTP(getRec, getReq)
if getRec.Code != http.StatusOK {
t.Fatalf("GET replicated object status %d: %s", getRec.Code, getRec.Body.String())
}
if !bytes.Equal(getRec.Body.Bytes(), data) {
t.Fatal("replicated SSE-C multipart object did not decrypt to source plaintext")
}
}
+10 -9
View File
@@ -93,15 +93,16 @@ func testSiteReplicationStatusAccountsPerSiteAndSurvivesMalformedConfig(obj Obje
CreatedAt: localMeta.Created, CreatedAt: localMeta.Created,
}, },
remoteBucket: { remoteBucket: {
Bucket: remoteBucket, Bucket: remoteBucket,
CreatedAt: remoteBucketMeta.Created, CreatedAt: remoteBucketMeta.Created,
Tags: encode(tagXML), Tags: encode(tagXML),
Versioning: encode(versioningXML), Versioning: encode(versioningXML),
ObjectLockConfig: encode(objectLockXML), ObjectLockConfig: encode(objectLockXML),
SSEConfig: encode(sseXML), SSEConfig: encode(sseXML),
QuotaConfig: encode(quotaJSON), QuotaConfig: encode(quotaJSON),
Policy: remotePolicy, Policy: remotePolicy,
CorsConfig: encode([]byte(testSiteReplicationCORSDoc)), CorsConfig: encode([]byte(testSiteReplicationCORSDoc)),
CorsConfigUpdatedAt: remoteBucketMeta.Created,
}, },
}, },
} }
+6 -7
View File
@@ -3557,13 +3557,12 @@ func (c *SiteReplicationSys) siteReplicationStatus(ctx context.Context, objAPI O
logInvalid("sse", err) logInvalid("sse", err)
} }
} }
if s.CorsConfig != nil { corsState, err := corsReplicationStateFromInfo(s.SRBucketInfo)
if _, err := decodeCORSReplicationPayload(s.CorsConfig); err == nil { if err != nil {
validCorsCfg[i] = true logInvalid("cors", err)
corsCfgCount++ } else if corsState.kind == corsReplicationLive {
} else { validCorsCfg[i] = true
logInvalid("cors", err) corsCfgCount++
}
} }
ss, ok := info.StatsSummary[s.DeploymentID] ss, ok := info.StatsSummary[s.DeploymentID]
if !ok { if !ok {
+2 -6
View File
@@ -26,13 +26,11 @@ import (
"io" "io"
"maps" "maps"
"net/http" "net/http"
"net/textproto"
"strings" "strings"
"time" "time"
"github.com/beevik/ntp" "github.com/beevik/ntp"
"github.com/minio/minio/internal/amztime" "github.com/minio/minio/internal/amztime"
xhttp "github.com/minio/minio/internal/http"
"github.com/minio/minio/internal/logger" "github.com/minio/minio/internal/logger"
"github.com/minio/pkg/v3/env" "github.com/minio/pkg/v3/env"
@@ -435,7 +433,7 @@ func IsObjectLockRequested(h http.Header) bool {
} }
// ParseObjectLockRetentionHeaders parses http headers to extract retention mode and retention date // ParseObjectLockRetentionHeaders parses http headers to extract retention mode and retention date
func ParseObjectLockRetentionHeaders(h http.Header) (rmode RetMode, r RetentionDate, err error) { func ParseObjectLockRetentionHeaders(h http.Header, allowPastRetainDate bool) (rmode RetMode, r RetentionDate, err error) {
retMode := h.Get(AmzObjectLockMode) retMode := h.Get(AmzObjectLockMode)
dateStr := h.Get(AmzObjectLockRetainUntilDate) dateStr := h.Get(AmzObjectLockRetainUntilDate)
if len(retMode) == 0 || len(dateStr) == 0 { if len(retMode) == 0 || len(dateStr) == 0 {
@@ -455,15 +453,13 @@ func ParseObjectLockRetentionHeaders(h http.Header) (rmode RetMode, r RetentionD
if err != nil { if err != nil {
return rmode, r, ErrInvalidRetentionDate return rmode, r, ErrInvalidRetentionDate
} }
_, replReq := h[textproto.CanonicalMIMEHeaderKey(xhttp.MinIOSourceReplicationRequest)]
t, err := UTCNowNTP() t, err := UTCNowNTP()
if err != nil { if err != nil {
lockLogIf(context.Background(), err) lockLogIf(context.Background(), err)
return rmode, r, ErrPastObjectLockRetainDate return rmode, r, ErrPastObjectLockRetainDate
} }
if retDate.Before(t) && !replReq { if retDate.Before(t) && !allowPastRetainDate {
return rmode, r, ErrPastObjectLockRetainDate return rmode, r, ErrPastObjectLockRetainDate
} }
+9 -1
View File
@@ -386,7 +386,7 @@ func TestParseObjectLockRetentionHeaders(t *testing.T) {
} }
for i, tt := range tests { for i, tt := range tests {
_, _, err := ParseObjectLockRetentionHeaders(tt.header) _, _, err := ParseObjectLockRetentionHeaders(tt.header, false)
//nolint:gocritic //nolint:gocritic
if tt.expectedErr == nil { if tt.expectedErr == nil {
if err != nil { if err != nil {
@@ -398,6 +398,14 @@ func TestParseObjectLockRetentionHeaders(t *testing.T) {
t.Fatalf("Case %d error: expected = %v, got = %v", i, tt.expectedErr, err) t.Fatalf("Case %d error: expected = %v, got = %v", i, tt.expectedErr, err)
} }
} }
past := http.Header{
xhttp.AmzObjectLockMode: []string{"governance"},
xhttp.AmzObjectLockRetainUntilDate: []string{"2017-01-02T15:04:05Z"},
}
if _, _, err := ParseObjectLockRetentionHeaders(past, true); err != nil {
t.Fatalf("trusted replica past retention date: %v", err)
}
} }
func TestGetObjectRetentionMeta(t *testing.T) { func TestGetObjectRetentionMeta(t *testing.T) {