mirror of
https://github.com/pgsty/minio.git
synced 2026-09-26 04:45:59 +03:00
114dc10529
Signed-off-by: Feng Ruohang <rh@vonng.com>
214 lines
8.3 KiB
Plaintext
214 lines
8.3 KiB
Plaintext
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(`<Tagging><TagSet><Tag><Key>key</Key><Value>same</Value></Tag></TagSet></Tagging>`)
|
|
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(`<Tagging><TagSet><Tag><Key>key</Key><Value>same</Value></Tag></TagSet></Tagging>`)
|
|
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(`<Tagging><TagSet><Tag><Key>key</Key><Value>older</Value></Tag></TagSet></Tagging>`))
|
|
newXML := []byte(`<Tagging><TagSet><Tag><Key>key</Key><Value>newest</Value></Tag></TagSet></Tagging>`)
|
|
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)
|
|
}
|
|
})
|
|
}})
|
|
}
|