package cmd import ( "encoding/base64" "encoding/json" "fmt" "net/http" "net/http/httptest" "testing" "time" "github.com/minio/madmin-go/v3" "github.com/minio/minio/internal/auth" ) func TestIssue77CurrentHealInputs(t *testing.T) { ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: review77HealInputs}) } func review77HealInputs(obj ObjectLayer, backend, bucket string, _ http.Handler, cred auth.Credentials, t *testing.T) { ctx := t.Context() localID, remoteID := globalDeploymentID(), "issue77-remote" created := UTCNow().Add(-time.Hour) at := created.Add(10 * time.Minute) newerAt := at.Add(time.Minute) enc := func(b []byte) *string { s := base64.StdEncoding.EncodeToString(b); return &s } base := newBucketMetadata(bucket) base.SetCreatedAt(created) base.defaultTimestamps() c := &SiteReplicationSys{enabled: true} status := func(local, remote madmin.SRBucketInfo, mismatch madmin.SRBucketStatsSummary) srStatusInfo { local.Bucket, remote.Bucket = bucket, bucket return srStatusInfo{Sites: map[string]madmin.PeerInfo{localID: {Name: "local"}, remoteID: {Name: "remote"}}, BucketStats: map[string]map[string]srBucketStatsSummary{bucket: {localID: {SRBucketStatsSummary: mismatch, meta: srBucketMetaInfo{SRBucketInfo: local, DeploymentID: localID}}, remoteID: {meta: srBucketMetaInfo{SRBucketInfo: remote, DeploymentID: remoteID}}}}} } t.Run(backend+"/quota-tombstone-cache", func(t *testing.T) { quota, _ := json.Marshal(madmin.BucketQuota{Quota: 1024, Type: madmin.HardQuota}) meta := base meta.QuotaConfigJSON, meta.QuotaConfigUpdatedAt = quota, at if err := globalBucketMetadataSys.save(ctx, meta); err != nil { t.Fatal(err) } s := status(madmin.SRBucketInfo{CreatedAt: created, QuotaConfig: enc(quota), QuotaConfigUpdatedAt: at}, madmin.SRBucketInfo{CreatedAt: created, QuotaConfigUpdatedAt: newerAt}, madmin.SRBucketStatsSummary{QuotaCfgMismatch: true}) if err := c.healBucketQuotaConfig(ctx, obj, bucket, s); err != nil { t.Fatal(err) } cached, _, err := globalBucketMetadataSys.GetQuotaConfig(ctx, bucket) if err != nil { t.Fatal(err) } disk, err := loadBucketMetadata(ctx, obj, bucket) if err != nil { t.Fatal(err) } if len(disk.QuotaConfigJSON) != 0 { t.Fatal("quota tombstone not written to disk") } if cached != nil && cached.Quota != 0 { t.Errorf("QUOTA_CACHE: disk has no quota, cache still enforces %d", cached.Quota) } }) t.Run(backend+"/same-payload-time-barrier", func(t *testing.T) { xml := []byte(`keysame`) meta := base meta.TaggingConfigXML, meta.TaggingConfigUpdatedAt = xml, at if err := globalBucketMetadataSys.save(ctx, meta); err != nil { t.Fatal(err) } // Even forcing mismatch=true cannot make the existing heal advance a // timestamp when the live payload bytes already match. s := status(madmin.SRBucketInfo{CreatedAt: created, Tags: enc(xml), TagConfigUpdatedAt: at}, madmin.SRBucketInfo{CreatedAt: created, Tags: enc(xml), TagConfigUpdatedAt: newerAt}, madmin.SRBucketStatsSummary{TagMismatch: true}) if err := c.healTagMetadata(ctx, obj, bucket, s); err != nil { t.Fatal(err) } after, err := loadBucketMetadata(ctx, obj, bucket) if err != nil { t.Fatal(err) } if !after.TaggingConfigUpdatedAt.Equal(newerAt) { t.Errorf("BARRIER_NOT_HEALED: same payload remains at %s, latest is %s", after.TaggingConfigUpdatedAt, newerAt) } }) t.Run(backend+"/creation-default-selection", func(t *testing.T) { policy := []byte(fmt.Sprintf(`{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*","Action":"s3:GetObject","Resource":"arn:aws:s3:::%s/*"}]}`, bucket)) meta := base meta.PolicyConfigJSON, meta.PolicyConfigUpdatedAt = policy, at baselineAt := newerAt.Add(time.Hour) s := status(madmin.SRBucketInfo{CreatedAt: created, Policy: policy, PolicyUpdatedAt: at}, madmin.SRBucketInfo{CreatedAt: baselineAt, PolicyUpdatedAt: baselineAt}, madmin.SRBucketStatsSummary{PolicyMismatch: true}) wrong := 0 for i := 0; i < 32; i++ { if err := globalBucketMetadataSys.save(ctx, meta); err != nil { t.Fatal(err) } if err := c.healBucketPolicies(ctx, obj, bucket, s); err != nil { t.Fatal(err) } after, err := loadBucketMetadata(ctx, obj, bucket) if err != nil { t.Fatal(err) } if len(after.PolicyConfigJSON) == 0 { wrong++ } } if wrong > 0 { t.Errorf("DEFAULT_SELECTED: creation-default erased valid policy in %d/32 heal rounds", wrong) } }) t.Run(backend+"/tag-remote-source-time", func(t *testing.T) { var got []madmin.SRBucketMeta remote := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { var item madmin.SRBucketMeta if err := json.NewDecoder(r.Body).Decode(&item); err != nil { t.Error(err) w.WriteHeader(http.StatusBadRequest) return } got = append(got, item) w.WriteHeader(http.StatusOK) })) defer remote.Close() svc, err := auth.CreateCredentials("issue77-heal-service", "issue77-heal-service-secret") if err != nil { t.Fatal(err) } svc.ParentUser = cred.AccessKey if _, err = globalIAMSys.store.AddServiceAccount(ctx, svc); err != nil { t.Fatal(err) } defer globalIAMSys.DeleteServiceAccount(ctx, svc.AccessKey, false) peers := map[string]madmin.PeerInfo{localID: {Name: "local", DeploymentID: localID}, remoteID: {Name: "remote", DeploymentID: remoteID, Endpoint: remote.URL}} remoteClient := &SiteReplicationSys{enabled: true, state: srState{ServiceAccountAccessKey: svc.AccessKey, Peers: peers}} xml := []byte(`keysame`) s := status(madmin.SRBucketInfo{CreatedAt: created, Tags: enc(xml), TagConfigUpdatedAt: newerAt}, madmin.SRBucketInfo{CreatedAt: created}, madmin.SRBucketStatsSummary{}) v := s.BucketStats[bucket][remoteID] v.TagMismatch = true s.BucketStats[bucket][remoteID] = v if err := remoteClient.healTagMetadata(ctx, obj, bucket, s); err != nil { t.Fatal(err) } if len(got) != 1 { t.Fatalf("remote received %d events", len(got)) } if !got[0].UpdatedAt.Equal(newerAt) { t.Errorf("TAG_WIRE_TIME: sent %s, want %s", got[0].UpdatedAt, newerAt) } }) } func TestIssue77CurrentPeerCheckBeforeLock(t *testing.T) { ExecObjectLayerAPITest(ExecObjectLayerAPITestArgs{t: t, objAPITest: func(obj ObjectLayer, backend, bucket string, _ http.Handler, _ auth.Credentials, t *testing.T) { t.Run(backend, func(t *testing.T) { ctx := t.Context() meta, err := loadBucketMetadata(ctx, obj, bucket) if err != nil { t.Fatal(err) } oldAt, newAt := meta.Created.Add(time.Hour), meta.Created.Add(2*time.Hour) oldXML := base64.StdEncoding.EncodeToString([]byte(`keyolder`)) newXML := []byte(`keynewest`) lockCtx, unlock, err := lockBucketMetadata(ctx, obj, bucket) if err != nil { t.Fatal(err) } locked := true defer func() { if locked { unlock() } }() ready := make(chan struct{}, 1) hook := func(name string) { if name == bucket { select { case ready <- struct{}{}: default: } } } lockBucketMetadataAcquireHook.Store(&hook) defer lockBucketMetadataAcquireHook.Store(nil) done := make(chan error, 1) go func() { done <- globalSiteReplicationSys.PeerBucketTaggingHandler(ctx, bucket, &oldXML, oldAt) }() select { case <-ready: case <-time.After(5 * time.Second): t.Fatal("older event never reached metadata lock") } // A newer writer commits while holding the existing metadata lock. // The old peer event has already checked the pre-commit cache. meta.TaggingConfigXML, meta.TaggingConfigUpdatedAt = newXML, newAt if err := globalBucketMetadataSys.saveMetadata(lockCtx, obj, meta); err != nil { t.Fatal(err) } unlock() locked = false select { case err := <-done: if err != nil { t.Fatal(err) } case <-time.After(5 * time.Second): t.Fatal("older event did not finish") } after, err := loadBucketMetadata(ctx, obj, bucket) if err != nil { t.Fatal(err) } if string(after.TaggingConfigXML) != string(newXML) { t.Errorf("CHECK_OUTSIDE_LOCK: queued older event replaced newer committed tags with %q", after.TaggingConfigXML) } }) }}) }