diff --git a/cmd/metacache-null-quorum_test.go b/cmd/metacache-null-quorum_test.go new file mode 100644 index 000000000..a0aa09308 --- /dev/null +++ b/cmd/metacache-null-quorum_test.go @@ -0,0 +1,433 @@ +// Copyright (c) 2026 Feng Ruohang +// +// This file is part of Silo Object Storage stack +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package cmd + +import ( + "bytes" + "context" + "encoding/xml" + "fmt" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "slices" + "sync" + "testing" + "time" + + madmin "github.com/minio/madmin-go/v3" +) + +// Hold nine completed disk walks while a PUT overwrites one name. The other +// seven walks start after that PUT completes. Both generations come from real +// disk walkers and valid PUT metadata; only the reader schedule is controlled. +type nullQuorumWalkSchedule struct { + bucket string + oldReady chan struct{} + resume chan struct{} + mu sync.Mutex + snapshots [16][]byte +} + +type nullQuorumWalkDisk struct { + StorageAPI + schedule *nullQuorumWalkSchedule + index int +} + +func (d *nullQuorumWalkDisk) WalkDir(ctx context.Context, opts WalkDirOptions, out io.Writer) error { + if opts.Bucket != d.schedule.bucket { + return d.StorageAPI.WalkDir(ctx, opts, out) + } + var stream bytes.Buffer + if d.index < 9 { + if err := d.StorageAPI.WalkDir(ctx, opts, &stream); err != nil { + return err + } + select { + case d.schedule.oldReady <- struct{}{}: + case <-ctx.Done(): + return ctx.Err() + } + } + select { + case <-d.schedule.resume: + case <-ctx.Done(): + return ctx.Err() + } + if d.index >= 9 { + if err := d.StorageAPI.WalkDir(ctx, opts, &stream); err != nil { + return err + } + } + d.schedule.mu.Lock() + d.schedule.snapshots[d.index] = bytes.Clone(stream.Bytes()) + d.schedule.mu.Unlock() + _, err := out.Write(stream.Bytes()) + return err +} + +func nullQuorumBackend(t *testing.T) (*erasureServerPools, string, http.Handler) { + t.Helper() + ctx, cancel := context.WithCancel(t.Context()) + obj, dirs, err := prepareErasure16(ctx) + if err != nil { + cancel() + t.Fatal(err) + } + z := obj.(*erasureServerPools) + previous := newObjectLayerFn() + setObjectLayer(z) + t.Cleanup(func() { + cancel() + z.Shutdown(context.Background()) + removeRoots(dirs) + setObjectLayer(previous) + }) + bucket, router, err := initAPIHandlerTest(ctx, z, nil, MakeBucketOptions{}) + if err != nil { + t.Fatal(err) + } + return z, bucket, router +} + +func nullQuorumRequest(t *testing.T, router http.Handler, method, target string) *httptest.ResponseRecorder { + t.Helper() + req, err := newTestSignedRequestV4(method, target, 0, nil, globalActiveCred.AccessKey, globalActiveCred.SecretKey, nil) + if err != nil { + t.Fatal(err) + } + rec := httptest.NewRecorder() + router.ServeHTTP(rec, req.WithContext(t.Context())) + return rec +} + +func TestListObjectsSingleNullQuorumHTTP(t *testing.T) { + z, bucket, router := nullQuorumBackend(t) + oldBody := bytes.Repeat([]byte{'a'}, 8192) + newBody := bytes.Repeat([]byte{'b'}, 8192) + oldTime := time.Now().UTC().Add(-time.Hour) + newTime := oldTime.Add(time.Minute) + var want []string + for i := range 16 { + name := fmt.Sprintf("fixed/%02d", i) + want = append(want, name) + _, err := z.PutObject(t.Context(), bucket, name, mustGetPutObjReader(t, bytes.NewReader(oldBody), int64(len(oldBody)), "", ""), ObjectOptions{MTime: oldTime}) + if err != nil { + t.Fatal(err) + } + } + const overwritten = "fixed/10" + schedule := &nullQuorumWalkSchedule{bucket: bucket, oldReady: make(chan struct{}, 16), resume: make(chan struct{})} + set := z.serverPools[0].sets[0] + disks := set.getDisks() + wrapped := make([]StorageAPI, len(disks)) + for i, disk := range disks { + wrapped[i] = &nullQuorumWalkDisk{StorageAPI: disk, schedule: schedule, index: i} + } + set.getDisks = func() []StorageAPI { return wrapped } + if globalAPIConfig.getListQuorum() != "strict" { + t.Fatal("test requires strict listing across all sixteen disks") + } + writeDone := make(chan error, 1) + putReader := mustGetPutObjReader(t, bytes.NewReader(newBody), int64(len(newBody)), "", "") + go func() { + defer close(schedule.resume) + ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) + defer cancel() + for range 9 { + select { + case <-schedule.oldReady: + case <-ctx.Done(): + writeDone <- ctx.Err() + return + } + } + _, err := z.PutObject(ctx, bucket, overwritten, putReader, ObjectOptions{MTime: newTime}) + writeDone <- err + }() + rec := nullQuorumRequest(t, router, http.MethodGet, getListObjectsV2URL("", bucket, "fixed/", "1000", "", "", "")) + if err := <-writeDone; err != nil { + t.Fatal(err) + } + var list ListObjectsV2Response + if rec.Code != http.StatusOK || xml.Unmarshal(rec.Body.Bytes(), &list) != nil { + t.Fatalf("LIST: %d %s", rec.Code, rec.Body.String()) + } + var got []string + for _, object := range list.Contents { + got = append(got, object.Key) + } + t.Logf("LIST HTTP=%d KeyCount=%d IsTruncated=%v keys=%v", rec.Code, list.KeyCount, list.IsTruncated, got) + if list.KeyCount != 16 || list.IsTruncated || !slices.Equal(got, want) { + t.Errorf("LIST omitted a name during overwrite: got %v, want %v", got, want) + } + + // Verify the actual reader inputs, including shape and EC, instead of + // assuming the scheduling barrier produced the intended resolver case. + schedule.mu.Lock() + snapshots := schedule.snapshots + schedule.mu.Unlock() + for i, snapshot := range snapshots { + reader := newMetacacheReader(bytes.NewReader(snapshot)) + var found bool + for { + entry, err := reader.next() + if err == io.EOF { + break + } + if err != nil { + t.Fatal(err) + } + if entry.name != overwritten { + continue + } + xl, err := entry.xlmeta() + if err != nil || len(xl.versions) != 1 { + t.Fatalf("disk %d: invalid input shape: %v", i, err) + } + header := xl.versions[0].header + wantTime := oldTime + if i >= 9 { + wantTime = newTime + } + if header.ModTime != wantTime.UnixNano() || header.VersionID != [16]byte{} || header.Type != ObjectType || header.FreeVersion() { + t.Fatalf("disk %d: unexpected header %v", i, header) + } + t.Logf("disk=%d versions=1 header=%v", i, header) + found = true + } + reader.Close() + if !found { + t.Fatalf("disk %d did not emit the overwritten name", i) + } + } + get := nullQuorumRequest(t, router, http.MethodGet, getGetObjectURL("", bucket, overwritten)) + head := nullQuorumRequest(t, router, http.MethodHead, getGetObjectURL("", bucket, overwritten)) + t.Logf("completed PUT readback: GET=%d bytes=%d HEAD=%d length=%s", get.Code, get.Body.Len(), head.Code, head.Header().Get("Content-Length")) + if get.Code != http.StatusOK || !bytes.Equal(get.Body.Bytes(), newBody) || head.Code != http.StatusOK || head.Header().Get("Content-Length") != "8192" { + t.Fatalf("completed PUT was not readable: GET=%d HEAD=%d", get.Code, head.Code) + } +} + +func TestSingleNullQuorumReadAndConsumers(t *testing.T) { + z, bucket, router := nullQuorumBackend(t) + set := z.serverPools[0].sets[0] + disks := set.getDisks() + const object = "object" + oldBody, newBody := bytes.Repeat([]byte{'a'}, 8192), bytes.Repeat([]byte{'b'}, 8192) + oldTime := time.Now().UTC().Add(-time.Hour) + put := func(body []byte, modTime time.Time) { + t.Helper() + _, err := z.PutObject(t.Context(), bucket, object, mustGetPutObjReader(t, bytes.NewReader(body), int64(len(body)), "", ""), ObjectOptions{MTime: modTime}) + if err != nil { + t.Fatal(err) + } + } + put(oldBody, oldTime) + oldMeta := make([][]byte, len(disks)) + for i, disk := range disks { + var err error + oldMeta[i], err = disk.ReadAll(t.Context(), bucket, object+"/"+xlStorageFormatFile) + if err != nil { + t.Fatal(err) + } + } + put(newBody, oldTime.Add(time.Minute)) + newMinorityMeta := mustReadNullQuorumMeta(t, disks[14], bucket, object) + // Model an incomplete overwrite using real inline shards: fifteen old + // copies retain a read quorum; the final disk contains the newer minority. + for i := range 15 { + if err := disks[i].WriteAll(t.Context(), bucket, object+"/"+xlStorageFormatFile, oldMeta[i]); err != nil { + t.Fatal(err) + } + } + checkRead := func(body []byte, modTime time.Time) { + t.Helper() + get := nullQuorumRequest(t, router, http.MethodGet, getGetObjectURL("", bucket, object)) + head := nullQuorumRequest(t, router, http.MethodHead, getGetObjectURL("", bucket, object)) + if get.Code != 200 || !bytes.Equal(get.Body.Bytes(), body) || head.Code != 200 || head.Header().Get("Content-Length") != "8192" { + t.Fatalf("GET/HEAD lost readable quorum: GET=%d HEAD=%d body=%q", get.Code, head.Code, get.Body.String()) + } + oi, err := z.GetObjectInfo(t.Context(), bucket, object, ObjectOptions{}) + if err != nil || !oi.ModTime.Equal(modTime) { + t.Fatalf("wrong readable generation: %v, %v", oi.ModTime, err) + } + } + checkRead(oldBody, oldTime) + // Reading must not reinterpret the minority as absent or delete it. + minority, err := disks[15].ReadVersion(t.Context(), "", bucket, object, "", ReadOptions{}) + if err != nil || !minority.ModTime.Equal(oldTime.Add(time.Minute)) { + t.Fatalf("read mutated the minority: %+v, %v", minority, err) + } + for _, kind := range []string{"rebalance", "decommission"} { + t.Run(kind, func(t *testing.T) { + var names []string + consume := func(entry metaCacheEntry) { + versions, err := entry.fileInfoVersions(bucket) + if err != nil || len(versions.Versions) != 1 || !versions.Versions[0].ModTime.Equal(oldTime) { + t.Errorf("invalid migration metadata: %+v, %v", versions, err) + return + } + // Migration reopens the selected name/version on all drives; + // it must not treat the selected listing metadata as disk agreement. + reader, err := set.GetObjectNInfo(t.Context(), bucket, entry.name, nil, nil, ObjectOptions{VersionID: nullVersionID, NoLock: true, NoDecryption: true}) + if err != nil { + t.Error(err) + return + } + body, err := io.ReadAll(reader) + reader.Close() + if err != nil || !bytes.Equal(body, oldBody) { + t.Errorf("migration could not read selected object data: %v", err) + } + names = append(names, entry.name) + } + var err error + if kind == "rebalance" { + err = set.listObjectsToRebalance(t.Context(), bucket, consume) + } else { + err = set.listObjectsToDecommission(t.Context(), decomBucketInfo{Name: bucket}, consume) + } + if err != nil || !slices.Equal(names, []string{object}) { + t.Fatalf("migration listing dropped the quorum: %v, %v", names, err) + } + }) + } + for _, latestOnly := range []bool{false, true} { + results := make(chan itemOrErr[ObjectInfo], 16) + if err := z.Walk(t.Context(), bucket, "", results, WalkOptions{AskDisks: "strict", LatestOnly: latestOnly}); err != nil { + t.Fatal(err) + } + var names []string + for result := range results { + if result.Err != nil || !result.Item.ModTime.Equal(oldTime) { + t.Fatalf("Walk returned invalid metadata: %+v", result) + } + names = append(names, result.Item.Name) + } + if !slices.Equal(names, []string{object}) { + t.Fatalf("Walk dropped the quorum: %v", names) + } + } + + // Exercise the scanner's actual abandoned-child path. Make the scan drive + // miss the object, while keeping an older quorum and a newer minority on + // the remaining disks. Its queued heal must read all disks and repair both. + if err := disks[14].WriteAll(t.Context(), bucket, object+"/"+xlStorageFormatFile, newMinorityMeta); err != nil { + t.Fatal(err) + } + if err := os.RemoveAll(filepath.Join(disks[15].Endpoint().Path, bucket, object)); err != nil { + t.Fatal(err) + } + scanNullQuorumAbandoned(t, z, bucket, object, oldTime) + for i, disk := range disks { + fi, err := disk.ReadVersion(t.Context(), "", bucket, object, "", ReadOptions{Healing: true}) + if err != nil || !fi.ModTime.Equal(oldTime) { + t.Fatalf("scanner heal did not reconcile disk %d: %v, %v", i, fi.ModTime, err) + } + } + checkRead(oldBody, oldTime) + put(newBody, oldTime.Add(2*time.Minute)) + checkRead(newBody, oldTime.Add(2*time.Minute)) +} + +func mustReadNullQuorumMeta(t *testing.T, disk StorageAPI, bucket, object string) []byte { + t.Helper() + data, err := disk.ReadAll(t.Context(), bucket, object+"/"+xlStorageFormatFile) + if err != nil { + t.Fatal(err) + } + return data +} + +func scanNullQuorumAbandoned(t *testing.T, z *erasureServerPools, bucket, object string, oldTime time.Time) { + t.Helper() + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + disks := z.serverPools[0].sets[0].getDisks() + seq := newBgHealSequence() + defer seq.cancelCtx() + previousState, previousRoutine := globalBackgroundHealState, globalBackgroundHealRoutine + globalBackgroundHealState = &allHealState{healSeqMap: map[string]*healSequence{"test": seq}} + routine := &healRoutine{tasks: make(chan healTask)} + globalBackgroundHealRoutine = routine + defer func() { globalBackgroundHealState, globalBackgroundHealRoutine = previousState, previousRoutine }() + workerDone := make(chan struct{}) + var healed []string + go func() { + defer close(workerDone) + for { + select { + case <-ctx.Done(): + return + case task := <-routine.tasks: + var result madmin.HealResultItem + var err error + if task.object == "" { + result, err = z.HealBucket(ctx, task.bucket, task.opts) + } else { + healed = append(healed, task.object+"/"+task.versionID) + result, err = z.HealObject(ctx, task.bucket, task.object, task.versionID, task.opts) + } + task.respCh <- healResult{result: result, err: err} + } + } + }() + cache := dataUsageCache{Info: dataUsageCacheInfo{Name: bucket}} + cache.replace(bucket, "", dataUsageEntry{}) + cache.replace(bucket+"/"+object, bucket, dataUsageEntry{Size: 8192, Objects: 1, Versions: 1}) + scanner := folderScanner{ + root: disks[15].Endpoint().Path, oldCache: cache, + newCache: dataUsageCache{Info: cache.Info}, updateCache: dataUsageCache{Info: cache.Info}, + disks: disks, disksQuorum: 8, healObjectSelect: 1, + weSleep: func() bool { return false }, shouldHeal: func() bool { return true }, updateCurrentPath: func(string) {}, + getSize: func(item scannerItem) (sizeSummary, error) { + if filepath.Base(item.Path) != xlStorageFormatFile { + return sizeSummary{}, errSkipFile + } + data, err := os.ReadFile(item.Path) + if err != nil { + return sizeSummary{}, err + } + var xl xlMetaV2 + if err := xl.Load(data); err != nil { + return sizeSummary{}, err + } + fi, err := xl.ToFileInfo(bucket, object, "", false, false) + if err != nil || !fi.ModTime.Equal(oldTime) { + return sizeSummary{}, fmt.Errorf("scanner re-read wrong generation: %v, %v", fi.ModTime, err) + } + return sizeSummary{totalSize: fi.Size, versions: 1}, nil + }, + } + var usage dataUsageEntry + err := scanner.scanFolder(ctx, cachedFolder{name: bucket, objectHealProbDiv: 1}, &usage) + cancel() + <-workerDone + if err != nil || len(healed) != 1 || healed[0] != object+"/"+nullVersionID && healed[0] != object+"/" { + t.Fatalf("scanner did not queue a name-based heal: %v, %v", healed, err) + } + flat := scanner.newCache.sizeRecursive(bucket) + if flat == nil || flat.Objects != 1 || flat.Size != 8192 { + t.Fatalf("scanner lost the healed object from usage: %+v", flat) + } + t.Logf("scanner queued %v; healed old quorum, repaired minority/missing copies, and retained one 8192-byte object", healed) +} diff --git a/cmd/xl-storage-format-v2-quorum_test.go b/cmd/xl-storage-format-v2-quorum_test.go new file mode 100644 index 000000000..23a96db7b --- /dev/null +++ b/cmd/xl-storage-format-v2-quorum_test.go @@ -0,0 +1,274 @@ +// Copyright (c) 2026 Feng Ruohang +// +// This file is part of Silo Object Storage stack +// +// This program is free software: you can redistribute it and/or modify +// it under the terms of the GNU Affero General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// This program is distributed in the hope that it will be useful +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Affero General Public License for more details. +// +// You should have received a copy of the GNU Affero General Public License +// along with this program. If not, see . + +package cmd + +import ( + "fmt" + "math/rand" + "reflect" + "slices" + "testing" + "time" +) + +// Use complete metadata so the selected version can also be decoded by LIST +// and the object read path, rather than testing only synthetic shallow headers. +func quorumNullVersion(t testing.TB, generation int64) xlMetaV2ShallowVersion { + t.Helper() + fi := newFileInfo("fixed/object", 12, 4) + fi.Erasure.Index = 1 + fi.DataDir = "11111111-1111-1111-1111-111111111111" + fi.ModTime = time.Unix(generation, 0).UTC() + fi.Size = 8192 + fi.Parts = []ObjectPartInfo{{Number: 1, Size: fi.Size, ActualSize: fi.Size}} + fi.Metadata = map[string]string{"etag": fmt.Sprintf("%032d", generation)} + var xl xlMetaV2 + if err := xl.AddVersion(fi); err != nil { + t.Fatal(err) + } + return xl.versions[0] +} + +// Start with forward, reverse and interleaved orders, then fixed-seed shuffles. +func quorumVersionOrders(input [][]xlMetaV2ShallowVersion) [][][]xlMetaV2ShallowVersion { + orders := [][][]xlMetaV2ShallowVersion{slices.Clone(input), slices.Clone(input)} + slices.Reverse(orders[1]) + interleaved := make([][]xlMetaV2ShallowVersion, 0, len(input)) + for i, j := 0, len(input)-1; i <= j; i, j = i+1, j-1 { + interleaved = append(interleaved, input[i]) + if i != j { + interleaved = append(interleaved, input[j]) + } + } + orders = append(orders, interleaved) + for seed := range int64(32) { + order := slices.Clone(input) + rand.New(rand.NewSource(seed)).Shuffle(len(order), func(i, j int) { + order[i], order[j] = order[j], order[i] + }) + orders = append(orders, order) + } + return orders +} + +func TestMergeXLV2SingleNullQuorum(t *testing.T) { + v := []xlMetaV2ShallowVersion{quorumNullVersion(t, 100), quorumNullVersion(t, 200), quorumNullVersion(t, 300)} + for _, tt := range []struct { + name string + counts [3]int + empty int + quorum int + want int + }{ + {"consistent", [3]int{16}, 0, 8, 0}, + {"old9-new7", [3]int{9, 7}, 0, 8, 0}, + {"old15-new1", [3]int{15, 1}, 0, 8, 0}, + {"old3-new1", [3]int{3, 1}, 0, 3, 0}, + {"two-quorums", [3]int{8, 8}, 0, 8, 1}, + {"empty-streams", [3]int{9, 1}, 6, 8, 0}, + {"two-subquorums", [3]int{7, 7}, 2, 8, -1}, + {"three-subquorums", [3]int{6, 5, 5}, 0, 8, -1}, + } { + t.Run(tt.name, func(t *testing.T) { + var input [][]xlMetaV2ShallowVersion + for g, n := range tt.counts { + for range n { + input = append(input, []xlMetaV2ShallowVersion{v[g]}) + } + } + input = append(input, make([][]xlMetaV2ShallowVersion, tt.empty)...) + want := []xlMetaV2ShallowVersion{} + if tt.want >= 0 { + want = append(want, v[tt.want]) + } + for order, versions := range quorumVersionOrders(input) { + for _, strict := range []bool{false, true} { + for _, requested := range []int{0, 1} { + got := mergeXLV2Versions(tt.quorum, strict, requested, versions...) + if !reflect.DeepEqual(got, want) { + t.Fatalf("order=%d strict=%v requested=%d: got %#v, want %#v", order, strict, requested, got, want) + } + } + } + } + }) + } +} + +func TestMergeXLV2SingleNullHeaderGroups(t *testing.T) { + old := quorumNullVersion(t, 100) + newer := quorumNullVersion(t, 200) + differentSig := old + differentSig.header.Signature[0]++ + differentFlags := old + differentFlags.header.Flags ^= xlFlagInlineData + for _, tt := range []struct { + name string + input [][]xlMetaV2ShallowVersion + strict bool + want []xlMetaV2ShallowVersion + }{ + {"signature-nonstrict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentSig}, {newer}}, false, []xlMetaV2ShallowVersion{differentSig}}, + {"signature-strict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentSig}, {newer}}, true, []xlMetaV2ShallowVersion{}}, + {"flags-nonstrict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentFlags}, {newer}}, false, []xlMetaV2ShallowVersion{}}, + {"flags-strict", [][]xlMetaV2ShallowVersion{{old}, {old}, {differentFlags}, {newer}}, true, []xlMetaV2ShallowVersion{}}, + } { + t.Run(tt.name, func(t *testing.T) { + got := mergeXLV2Versions(3, tt.strict, 0, tt.input...) + if !reflect.DeepEqual(got, tt.want) { + t.Fatalf("got %#v, want %#v", got, tt.want) + } + }) + } +} + +func TestResolveSingleNullQuorum(t *testing.T) { + old, newer := quorumNullVersion(t, 100), quorumNullVersion(t, 200) + input := make([][]xlMetaV2ShallowVersion, 16) + for i := range input { + input[i] = []xlMetaV2ShallowVersion{old} + if i >= 9 { + input[i] = []xlMetaV2ShallowVersion{newer} + } + } + for order, versions := range quorumVersionOrders(input) { + entries := make(metaCacheEntries, len(versions)) + for i := range versions { + xl := &xlMetaV2{versions: versions[i]} + metadata, err := xl.AppendTo(nil) + if err != nil { + t.Fatal(err) + } + entries[i] = metaCacheEntry{name: "fixed/object", metadata: metadata} + } + selected, ok := entries.resolve(&metadataResolutionParams{objQuorum: 8, dirQuorum: 8, requestedVersions: 1}) + if !ok { + t.Fatalf("order %d: a complete null object version has quorum but was omitted", order) + } + xl, err := selected.xlmeta() + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(xl.versions, []xlMetaV2ShallowVersion{old}) { + t.Fatalf("order %d: incorrect complete selected version: %#v", order, xl.versions) + } + fi, err := xl.ToFileInfo("bucket", "fixed/object", "", false, false) + if err != nil || !fi.IsValid() || fi.Size != 8192 || !fi.ModTime.Equal(time.Unix(100, 0)) { + t.Fatalf("order %d: selected metadata cannot be read: %+v, %v", order, fi, err) + } + } +} + +// These are compatibility assertions for inputs outside the narrow repair. +// In particular, the mixed-version outputs below are not a new correctness +// contract for version enumeration; that behavior requires a separate change. +func TestMergeXLV2NullHistoriesUnchanged(t *testing.T) { + old, newer := quorumNullVersion(t, 100), quorumNullVersion(t, 300) + middle := quorumNullVersion(t, 200) + middle.header.VersionID = [16]byte{1} + highest := quorumNullVersion(t, 400) + highest.header.VersionID = [16]byte{2} + for _, tt := range []struct { + name string + input [][]xlMetaV2ShallowVersion + want []xlMetaV2ShallowVersion + }{ + {"mixed-forward", [][]xlMetaV2ShallowVersion{{old}, {old}, {old}, {newer, middle}, {middle}, {middle}}, []xlMetaV2ShallowVersion{middle}}, + {"mixed-reverse", [][]xlMetaV2ShallowVersion{{middle}, {middle}, {newer, middle}, {old}, {old}, {old}}, []xlMetaV2ShallowVersion{old}}, + {"history-pruned-first", [][]xlMetaV2ShallowVersion{{highest, old}, {highest, old}, {highest, old}, {highest, newer}}, []xlMetaV2ShallowVersion{highest}}, + } { + t.Run(tt.name, func(t *testing.T) { + got := mergeXLV2Versions(3, false, 0, tt.input...) + if !reflect.DeepEqual(got, tt.want) { + t.Fatalf("excluded history changed: got %#v, want %#v", got, tt.want) + } + }) + } + for _, kind := range []string{"delete", "free", "legacy", "mixed-ec", "legacy-modern-ec", "nonzero-strict"} { + t.Run(kind, func(t *testing.T) { + a, b := old, newer + switch kind { + case "delete": + b.header.Type = DeleteType + case "free": + b.header.Flags |= xlFlagFreeVersion + case "legacy": + b.header.Type = LegacyType + case "mixed-ec": + b.header.EcN, b.header.EcM = 8, 8 + case "legacy-modern-ec": + b.header.EcN, b.header.EcM = 0, 0 + case "nonzero-strict": + a.header.VersionID, b.header.VersionID = [16]byte{1}, [16]byte{1} + } + for _, requested := range []int{0, 1} { + got := mergeXLV2Versions(3, true, requested, []xlMetaV2ShallowVersion{a}, []xlMetaV2ShallowVersion{a}, []xlMetaV2ShallowVersion{a}, []xlMetaV2ShallowVersion{b}) + if len(got) != 0 { + t.Fatalf("excluded input changed: requested=%d got %#v", requested, got) + } + } + }) + } + // All-legacy EC headers are eligible, unlike mixing legacy and modern EC. + old.header.EcN, old.header.EcM = 0, 0 + newer.header.EcN, newer.header.EcM = 0, 0 + got := mergeXLV2Versions(3, false, 0, []xlMetaV2ShallowVersion{old}, []xlMetaV2ShallowVersion{old}, []xlMetaV2ShallowVersion{old}, []xlMetaV2ShallowVersion{newer}) + if !reflect.DeepEqual(got, []xlMetaV2ShallowVersion{old}) { + t.Fatalf("all-legacy EC lost the old quorum: %#v", got) + } +} + +func BenchmarkResolveSingleNullQuorumPage(b *testing.B) { + old, newer := quorumNullVersion(b, 100), quorumNullVersion(b, 200) + for _, kind := range []string{"consistent", "old-first", "new-first"} { + b.Run(kind, func(b *testing.B) { + var templates [16]metaCacheEntry + for i := range templates { + v := old + if kind == "old-first" && i >= 9 || kind == "new-first" && i < 7 { + v = newer + } + xl := &xlMetaV2{versions: []xlMetaV2ShallowVersion{v}} + data, err := xl.AppendTo(nil) + if err != nil { + b.Fatal(err) + } + templates[i] = metaCacheEntry{metadata: data} + } + var names [1000]string + for i := range names { + names[i] = fmt.Sprintf("fixed/%04d", i) + } + resolver := metadataResolutionParams{objQuorum: 8, dirQuorum: 8, requestedVersions: 1} + b.ReportAllocs() + b.ResetTimer() + for b.Loop() { + for _, name := range names { + entries := templates + for i := range entries { + entries[i].name = name + } + selected, ok := metaCacheEntries(entries[:]).resolve(&resolver) + if ok && selected.reusable { + metaDataPoolPut(selected.metadata) + } + } + } + }) + } +} diff --git a/cmd/xl-storage-format-v2.go b/cmd/xl-storage-format-v2.go index f50267e25..fd6ac1642 100644 --- a/cmd/xl-storage-format-v2.go +++ b/cmd/xl-storage-format-v2.go @@ -1941,6 +1941,10 @@ func mergeXLV2Versions(quorum int, strict bool, requestedVersions int, versions // No need for non-strict checks if quorum is 1. strict = true } + // Keep the original stream shapes: pruning must not make a versioned object + // eligible for the single null-version recount below. + originalVersions := versions + var checkedSingleNull, singleNull bool // Shallow copy input versions = append(make([][]xlMetaV2ShallowVersion, 0, len(versions)), versions...) @@ -2015,44 +2019,23 @@ func mergeXLV2Versions(quorum int, strict bool, requestedVersions int, versions // Version IDs match, but otherwise unable to resolve. // We are either strict, or don't have enough information to match. // Switch to a pure counting algo. - x := make(map[xlMetaV2VersionHeader]int, len(tops)) - for _, a := range tops { - if a.header.VersionID != ver.header.VersionID { - continue - } - if !strict { - // we must match EC, when we are not strict. - if !a.header.matchesEC(ver.header) { - continue - } - - a.header.Signature = [4]byte{} - } - x[a.header]++ - } - latestCount = 0 - for k, v := range x { - if v < latestCount { - continue - } - if v == latestCount && latest.header.sortsBefore(k) { - // Tiebreak, use sort. - continue - } - for _, a := range tops { - hdr := a.header - if !strict { - hdr.Signature = [4]byte{} - } - if hdr == k { - latest = a - } - } - latestCount = v - } + latest, latestCount = countXLV2Versions(tops, ver.header, latest, strict) break } } + if latestCount < quorum { + if !checkedSingleNull { + singleNull = singleNullVersionStreams(originalVersions) + checkedSingleNull = true + } + if singleNull { + // A newer minority at the end can hide an older quorum from + // the selection loop. Recount before discarding the null ID. + if candidate, count := countXLV2Versions(tops, latest.header, latest, strict); count >= quorum { + latest, latestCount = candidate, count + } + } + } if latestCount >= quorum { merged = append(merged, latest) @@ -2113,6 +2096,70 @@ func mergeXLV2Versions(quorum int, strict bool, requestedVersions int, versions return merged } +// singleNullVersionStreams excludes histories and other version types from the +// additional recount. Check the original inputs, before any stream is pruned. +func singleNullVersionStreams(versions [][]xlMetaV2ShallowVersion) bool { + var ec xlMetaV2VersionHeader + var haveEC bool + for _, stream := range versions { + if len(stream) == 0 { + continue + } + if len(stream) != 1 { + return false + } + h := stream[0].header + if h.VersionID != [16]byte{} || h.Type != ObjectType || h.FreeVersion() { + return false + } + if haveEC && (h.EcN != ec.EcN || h.EcM != ec.EcM) { + return false + } + ec, haveEC = h, true + } + return haveEC +} + +// countXLV2Versions selects the most frequent compatible header for reference's +// VersionID, retaining the existing sort tiebreak and last matching entry. +func countXLV2Versions(tops []xlMetaV2ShallowVersion, reference xlMetaV2VersionHeader, latest xlMetaV2ShallowVersion, strict bool) (xlMetaV2ShallowVersion, int) { + x := make(map[xlMetaV2VersionHeader]int, len(tops)) + for _, a := range tops { + if a.header.VersionID != reference.VersionID { + continue + } + if !strict { + // we must match EC, when we are not strict. + if !a.header.matchesEC(reference) { + continue + } + a.header.Signature = [4]byte{} + } + x[a.header]++ + } + var latestCount int + for k, v := range x { + if v < latestCount { + continue + } + if v == latestCount && latest.header.sortsBefore(k) { + // Tiebreak, use sort. + continue + } + for _, a := range tops { + hdr := a.header + if !strict { + hdr.Signature = [4]byte{} + } + if hdr == k { + latest = a + } + } + latestCount = v + } + return latest, latestCount +} + type xlMetaBuf []byte // ToFileInfo converts xlMetaV2 into a common FileInfo datastructure