mirror of
https://github.com/pgsty/minio.git
synced 2026-09-26 04:45:59 +03:00
test: synchronize conditional PUT disk fixtures
Replace unsynchronized getDisks swaps with backing disk-list updates under erasureDisksMu, matching the existing GetDisks reader lock. Apply the same helper to capacity and read-fault adapters while preserving nested restore ordering. Add a regression that overlaps fixture changes with the real IAM Walk reader, and run the conditional PUT suite under the race detector in CI. The regression reproduces the old fixture race; ten fixed race iterations pass without warnings. Production conditional PUT behavior is unchanged. Refs #199 Signed-off-by: Feng Ruohang <rh@vonng.com>
This commit is contained in:
@@ -90,6 +90,9 @@ jobs:
|
|||||||
- name: Run S3 Select tests under race detector
|
- name: Run S3 Select tests under race detector
|
||||||
run: go test -race ./internal/s3select/... -count=1
|
run: go test -race ./internal/s3select/... -count=1
|
||||||
|
|
||||||
|
- name: Run conditional PUT tests under race detector
|
||||||
|
run: go test -race ./cmd -run '^Test(PoolsConditionalPut|SinglePoolConditionalPutHTTP)' -count=1 -timeout=5m
|
||||||
|
|
||||||
crosscompile:
|
crosscompile:
|
||||||
name: Cross Compile
|
name: Cross Compile
|
||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
|
|||||||
@@ -52,18 +52,32 @@ func (d conditionalPutCapacityDisk) DiskInfo(ctx context.Context, opts DiskInfoO
|
|||||||
return info, err
|
return info, err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Keep getDisks immutable while background IAM and storage readers use it.
|
||||||
|
// GetDisks takes this same mutex when it copies the backing disk list.
|
||||||
|
func conditionalPutSwapDisks(pool *erasureSets, object string, wrap func(StorageAPI) StorageAPI) func() {
|
||||||
|
setIndex := pool.getHashedSet(object).setIndex
|
||||||
|
pool.erasureDisksMu.Lock()
|
||||||
|
previous := pool.erasureDisks[setIndex]
|
||||||
|
disks := append([]StorageAPI(nil), previous...)
|
||||||
|
for i, disk := range disks {
|
||||||
|
disks[i] = wrap(disk)
|
||||||
|
}
|
||||||
|
pool.erasureDisks[setIndex] = disks
|
||||||
|
pool.erasureDisksMu.Unlock()
|
||||||
|
return func() {
|
||||||
|
pool.erasureDisksMu.Lock()
|
||||||
|
pool.erasureDisks[setIndex] = previous
|
||||||
|
pool.erasureDisksMu.Unlock()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func conditionalPutPool(t *testing.T, z *erasureServerPools, object string, target int) func() {
|
func conditionalPutPool(t *testing.T, z *erasureServerPools, object string, target int) func() {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
var restore []func()
|
var restore []func()
|
||||||
for i, pool := range z.serverPools {
|
for i, pool := range z.serverPools {
|
||||||
set := pool.getHashedSet(object)
|
restore = append(restore, conditionalPutSwapDisks(pool, object, func(disk StorageAPI) StorageAPI {
|
||||||
previous := set.getDisks
|
return conditionalPutCapacityDisk{StorageAPI: disk, full: i != target}
|
||||||
disks := append([]StorageAPI(nil), previous()...)
|
}))
|
||||||
for j, disk := range disks {
|
|
||||||
disks[j] = conditionalPutCapacityDisk{StorageAPI: disk, full: i != target}
|
|
||||||
}
|
|
||||||
set.getDisks = func() []StorageAPI { return disks }
|
|
||||||
restore = append(restore, func() { set.getDisks = previous })
|
|
||||||
}
|
}
|
||||||
return func() {
|
return func() {
|
||||||
for _, fn := range restore {
|
for _, fn := range restore {
|
||||||
@@ -230,14 +244,10 @@ func TestPoolsConditionalPutUnreadable(t *testing.T) {
|
|||||||
putConsistencyObject(t, z, bucket, object, 1, "current", ObjectOptions{MTime: UTCNow().Add(-time.Minute)})
|
putConsistencyObject(t, z, bucket, object, 1, "current", ObjectOptions{MTime: UTCNow().Add(-time.Minute)})
|
||||||
}
|
}
|
||||||
defer conditionalPutPool(t, z, object, 0)()
|
defer conditionalPutPool(t, z, object, 0)()
|
||||||
set := z.serverPools[faultPool].getHashedSet(object)
|
restoreFault := conditionalPutSwapDisks(z.serverPools[faultPool], object, func(disk StorageAPI) StorageAPI {
|
||||||
original := set.getDisks
|
return consistencyReadFaultDisk{StorageAPI: disk, bucket: bucket, object: object}
|
||||||
disks := append([]StorageAPI(nil), original()...)
|
})
|
||||||
for i := range disks {
|
defer restoreFault()
|
||||||
disks[i] = consistencyReadFaultDisk{StorageAPI: disks[i], bucket: bucket, object: object}
|
|
||||||
}
|
|
||||||
set.getDisks = func() []StorageAPI { return disks }
|
|
||||||
defer func() { set.getDisks = original }()
|
|
||||||
headers := map[string]string{xhttp.IfNoneMatch: "*"}
|
headers := map[string]string{xhttp.IfNoneMatch: "*"}
|
||||||
if match {
|
if match {
|
||||||
headers = map[string]string{xhttp.IfMatch: "*"}
|
headers = map[string]string{xhttp.IfMatch: "*"}
|
||||||
@@ -252,7 +262,7 @@ func TestPoolsConditionalPutUnreadable(t *testing.T) {
|
|||||||
if !isErrReadQuorum(err) || called != 0 {
|
if !isErrReadQuorum(err) || called != 0 {
|
||||||
t.Errorf("lookup error=%v callback calls=%d", err, called)
|
t.Errorf("lookup error=%v callback calls=%d", err, called)
|
||||||
}
|
}
|
||||||
set.getDisks = original
|
restoreFault()
|
||||||
get := multipartConditionRequest(t, router, http.MethodGet, url, "", nil)
|
get := multipartConditionRequest(t, router, http.MethodGet, url, "", nil)
|
||||||
if present {
|
if present {
|
||||||
if get.Code != http.StatusOK || get.Body.String() != "current" || multipartConditionResponseETag(get) != fmt.Sprintf("%x", md5.Sum([]byte("current"))) {
|
if get.Code != http.StatusOK || get.Body.String() != "current" || multipartConditionResponseETag(get) != fmt.Sprintf("%x", md5.Sum([]byte("current"))) {
|
||||||
@@ -595,14 +605,10 @@ func TestPoolsConditionalPutReplicaAvailability(t *testing.T) {
|
|||||||
object := "replica-availability"
|
object := "replica-availability"
|
||||||
oi := putConsistencyObject(t, z, bucket, object, 0, "old", ObjectOptions{Versioned: addressed, MTime: UTCNow().Add(-time.Minute)})
|
oi := putConsistencyObject(t, z, bucket, object, 0, "old", ObjectOptions{Versioned: addressed, MTime: UTCNow().Add(-time.Minute)})
|
||||||
defer conditionalPutPool(t, z, object, 0)()
|
defer conditionalPutPool(t, z, object, 0)()
|
||||||
set := z.serverPools[1].getHashedSet(object)
|
restoreFault := conditionalPutSwapDisks(z.serverPools[1], object, func(disk StorageAPI) StorageAPI {
|
||||||
original := set.getDisks
|
return consistencyReadFaultDisk{StorageAPI: disk, bucket: bucket, object: object}
|
||||||
disks := append([]StorageAPI(nil), original()...)
|
})
|
||||||
for i := range disks {
|
defer restoreFault()
|
||||||
disks[i] = consistencyReadFaultDisk{StorageAPI: disks[i], bucket: bucket, object: object}
|
|
||||||
}
|
|
||||||
set.getDisks = func() []StorageAPI { return disks }
|
|
||||||
defer func() { set.getDisks = original }()
|
|
||||||
headers := map[string]string{
|
headers := map[string]string{
|
||||||
xhttp.MinIOSourceReplicationRequest: "true",
|
xhttp.MinIOSourceReplicationRequest: "true",
|
||||||
xhttp.AmzBucketReplicationStatus: "REPLICA",
|
xhttp.AmzBucketReplicationStatus: "REPLICA",
|
||||||
@@ -618,7 +624,7 @@ func TestPoolsConditionalPutReplicaAvailability(t *testing.T) {
|
|||||||
if put.Code != wantStatus {
|
if put.Code != wantStatus {
|
||||||
t.Fatalf("replica PUT: %d want %d: %s", put.Code, wantStatus, put.Body.String())
|
t.Fatalf("replica PUT: %d want %d: %s", put.Code, wantStatus, put.Body.String())
|
||||||
}
|
}
|
||||||
set.getDisks = original
|
restoreFault()
|
||||||
get := multipartConditionRequest(t, router, http.MethodGet, url, "", nil)
|
get := multipartConditionRequest(t, router, http.MethodGet, url, "", nil)
|
||||||
if get.Code != http.StatusOK || get.Body.String() != wantBody {
|
if get.Code != http.StatusOK || get.Body.String() != wantBody {
|
||||||
t.Fatalf("replica GET: %d %q", get.Code, get.Body.String())
|
t.Fatalf("replica GET: %d %q", get.Code, get.Body.String())
|
||||||
@@ -684,3 +690,32 @@ func TestPoolsConditionalPutDeleteMarkerTie(t *testing.T) {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Overlap fixture changes with the real IAM Walk reader instead of relying on
|
||||||
|
// its periodic refresh timer to expose an unsynchronized disk-adapter swap.
|
||||||
|
func TestPoolsConditionalPutFixtureConcurrentIAM(t *testing.T) {
|
||||||
|
z, _ := consistencyPools(t)
|
||||||
|
conditionalPutBucket(t, z, "unversioned")
|
||||||
|
ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
started, finished := make(chan struct{}), make(chan error, 1)
|
||||||
|
iam := globalIAMSys
|
||||||
|
go func() {
|
||||||
|
close(started)
|
||||||
|
for range 20 {
|
||||||
|
if err := iam.Load(ctx, false); err != nil {
|
||||||
|
finished <- err
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
finished <- nil
|
||||||
|
}()
|
||||||
|
<-started
|
||||||
|
for i := range 5000 {
|
||||||
|
restore := conditionalPutPool(t, z, "fixture-concurrent-iam", i%2)
|
||||||
|
restore()
|
||||||
|
}
|
||||||
|
if err := <-finished; err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user