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)
			}
		})
	}})
}
