fix: retain quorate null versions after newer minorities

Recount only when the existing selection is below quorum and each original
nonempty stream has one ordinary null version with identical EC headers.
Reuse the existing header grouping and leave history pruning unchanged.

Add deterministic signed LIST, complete-metadata permutation, excluded-history,
object-read, scanner/heal, Walk and migration-consumer regressions.

Signed-off-by: Feng Ruohang <rh@vonng.com>
(cherry picked from commit 9804e4deb4d5a18cca640be72bff9b47a415f7de)
Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
Feng Ruohang
2026-09-16 21:07:17 +08:00
parent fced863036
commit 8d06424b12
3 changed files with 789 additions and 35 deletions
+433
View File
@@ -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 <http://www.gnu.org/licenses/>.
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)
}
+274
View File
@@ -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 <http://www.gnu.org/licenses/>.
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)
}
}
}
})
}
}
+82 -35
View File
@@ -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