diff --git a/README.md b/README.md index 82b619bb..5bda7444 100644 --- a/README.md +++ b/README.md @@ -791,6 +791,142 @@ The analysis throughput represents the object count and sizes as they are writte Request times shown with `--analyze.v` represents request time for each fan-out call. +## ATOMIC + +The atomic benchmark checks whether an S3 server keeps the consistency +guarantees that Amazon S3 has provided since December 2020: + +- An overwrite replaces the whole object at once. +- A read that starts after a write succeeds sees that write. +- A listing that starts after a PUT or DELETE succeeds shows that change. + +Many applications depend on these guarantees, for example the Hadoop S3A +connector used by Spark. When a server breaks them, the application does not +see an error. It reads old or partial data, or misses a file. + +Other warp benchmarks measure speed and check only object sizes. This benchmark +checks the content, metadata, ETag and listing entry of every response. + +### Run the benchmark + +List each server directly with `--host`, not a load balancer in front of them. +Warp can then send a write to one server and the next read to another: + +``` +λ warp atomic --host=10.0.0.{1...8}:9000 --host-select=roundrobin \ + --objects=16 --obj.size=16MiB --obj.randsize --concurrent=64 --duration=5m +``` + +Use a bucket that holds no other data. Warp empties the bucket before the run, +as it does for every benchmark. After the run, it deletes only the keys it wrote. + +### How warp tells which PUT a response came from + +Every 4 KiB block of an uploaded object starts with a stamp. The stamp names +the PUT that wrote the object and the key it was written to. Warp also stores +the PUT's name in the `X-Amz-Meta-Warp-Atomic` metadata. Warp can therefore +look at any response on its own and tell which PUT produced each part of it. +Uploads always use a single part. + +When the run starts, warp uploads one object and reads it back. It stops with +an error in two cases: + +- The server does not return the `Warp-Atomic` metadata. +- The ETag is not the MD5 of the body. Some servers use other ETags. Rerun + with `--no-etag-md5` for those servers. + +### Operations + +Warp picks each operation at random, in the proportions set by the +`--*-distrib` flags: + +- **PUT** overwrites one of the keys. By default, warp then reads the key back + at once. When `--host` lists more than one server, the read goes to a + different server than the PUT. +- **GET** reads a key and checks the body, metadata and ETag. +- **STAT** reads a key's metadata and checks it. +- **LIST** runs a cycle of five requests: + 1. PUT a new key. + 2. List the key's prefix. The new key must appear. + 3. DELETE the new key. + 4. List the prefix again. The key must be gone. + 5. List the overwritten keys. Each entry must show a current PUT. + + With more than one server in `--host`, steps 2 and 4 go to a different server + than the write before them. + +When the run ends, warp reads every key once more after all writes have stopped. + +### Violations + +Warp reports each failed check as an error that starts with +`atomic :`. The message names the server that answered the read and, +where warp knows it, the server that took the write. If there were violations, +warp prints the number in each category at the end of the run. + +A GET or STAT counts as a violation when it returns: + +| Category | What the server returned | +| --------------- | ----------------------------------------------------------- | +| `torn` | a body made of blocks from two different PUTs | +| `corrupt` | a block that no PUT wrote | +| `length` | a body shorter or longer than the PUT that wrote it | +| `wrong-key` | a body that was written to a different key | +| `missing` | `NoSuchKey` for a key that exists for the whole run | +| `meta-missing` | an object without the `Warp-Atomic` metadata | +| `meta-mismatch` | metadata from one PUT with the body or size of another | +| `etag-mismatch` | an ETag different from the one the server gave for that PUT | +| `etag-md5` | an ETag that is not the MD5 of the body | +| `stale` | a PUT that another PUT replaced before the read started | +| `phantom` | a PUT that this warp process never sent | + +A listing counts as a violation when it shows: + +| Category | What the listing showed | +| -------------- | -------------------------------------------------------------- | +| `list-missing` | no entry for a key whose PUT succeeded | +| `list-deleted` | an entry for a key whose DELETE succeeded | +| `list-size` | a size that belongs to a different PUT than the listed ETag | +| `list-etag` | an ETag different from the one the server gave for the new key | +| `list-stale` | a PUT that another PUT replaced before the listing started | + +### When a read counts as stale + +Warp flags a read as `stale` only when it is certain that the returned data was +already replaced. That is the case when another PUT to the same key did both of +these things: + +- It started after the returned PUT succeeded. +- It succeeded before the read started. + +Warp does not flag these cases: + +- The read overlaps two PUTs. The server may apply concurrent writes in either + order. +- The read returns a PUT that failed. A failed PUT may still have been applied. + +Warp judges staleness only against PUTs sent by the same warp process, because +it compares their times. In distributed mode, each client checks its own PUTs. +The body, metadata and ETag checks apply to every read. + +### Parameters + +- `--objects=N` sets the number of keys. The default is 16. Fewer keys means + more writes to each key. +- `--obj.size=N` sets the object size. The default is 4MiB. Add `--obj.randsize` + so that an overwrite also changes the size. +- `--block.size=N` sets how often the stamp repeats. The default is 4096 bytes. +- `--put-distrib`, `--get-distrib`, `--stat-distrib` and `--list-distrib` set the + mix of operations. The defaults are 40, 45, 15 and 10. +- `--read-after-write` reads each key back straight after a successful PUT. It + is on by default. Turn it off with `--read-after-write=false`. +- `--seed=N` sets the seed for each thread's choice of key, operation and object + size. The default, 0, picks a random seed. Every violation and the final + summary show the seed. Rerunning with the same seed, concurrency and number + of clients repeats each thread's sequence of operations. It does not repeat + the timing between threads and servers, so a violation may not recur. +- `--no-etag-md5` turns off the check that the ETag is the MD5 of the body. + # Analysis When benchmarks have finished all request data will be saved to a file and an analysis will be shown. diff --git a/cli/atomic.go b/cli/atomic.go new file mode 100644 index 00000000..79c54655 --- /dev/null +++ b/cli/atomic.go @@ -0,0 +1,162 @@ +/* + * Warp (C) 2026 MinIO, Inc. + * + * 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 cli + +import ( + "math" + + "github.com/minio/cli" + "github.com/minio/mc/pkg/probe" + "github.com/minio/minio-go/v7" + "github.com/minio/pkg/v3/console" + "github.com/minio/warp/pkg/bench" + "github.com/minio/warp/pkg/generator" +) + +var atomicFlags = []cli.Flag{ + cli.IntFlag{ + Name: "objects", + Value: 16, + Usage: "Number of keys overwritten concurrently. Fewer keys means more contention per key.", + }, + cli.StringFlag{ + Name: "obj.size", + Value: "4MiB", + Usage: "Size of each generated object. Combine with --obj.randsize so an overwrite changes the length. All sizes are base 2 binary.", + }, + cli.IntFlag{ + Name: "block.size", + Value: 4096, + Usage: "Every block of this many bytes carries the identity of the PUT it belongs to.", + }, + cli.Float64Flag{ + Name: "put-distrib", + Usage: "The amount of PUT operations.", + Value: 40, + }, + cli.Float64Flag{ + Name: "get-distrib", + Usage: "The amount of GET operations.", + Value: 45, + }, + cli.Float64Flag{ + Name: "stat-distrib", + Usage: "The amount of STAT operations.", + Value: 15, + }, + cli.Float64Flag{ + Name: "list-distrib", + Usage: "The amount of LIST cycles. Each PUTs a new key, lists it, deletes it, lists again, and lists the overwritten keys.", + Value: 10, + }, + cli.BoolTFlag{ + Name: "read-after-write", + Usage: "Follow every acknowledged PUT with a GET of the same key, from another host when --host lists several.", + }, + cli.Uint64Flag{ + Name: "seed", + Usage: "Seed for each thread's choice of key, operation and object size. 0 picks one at random. Violations report the seed.", + }, + cli.BoolFlag{ + Name: "no-etag-md5", + Usage: "Do not require the ETag to be the MD5 of the body, for servers that use opaque ETags.", + }, +} + +var AtomicCombinedFlags = combineFlags(globalFlags, ioFlags, atomicFlags, genFlags, benchFlags, analyzeFlags) + +var atomicCmd = cli.Command{ + Name: "atomic", + Usage: "verify that concurrent overwrites are atomic and reads are not stale", + Action: mainAtomic, + Before: setGlobalsFromContext, + Flags: AtomicCombinedFlags, + CustomHelpTemplate: `NAME: + {{.HelpName}} - {{.Usage}} + +USAGE: + {{.HelpName}} [FLAGS] + -> see https://github.com/minio/warp#atomic + +Every PUT writes a body in which each block names the PUT it belongs to, +and records that name in object metadata. Every GET and STAT checks that +body, length, metadata and ETag all belong to one PUT, and that the PUT +had not already been overwritten when the read started. With +--list-distrib, LIST cycles check list-after-write and list-after-delete +on new keys, and that listings of the overwritten keys are not stale. +Violations are reported as errors prefixed with "atomic :". + +Staleness is judged only against PUTs issued by the same warp process; +body, metadata and ETag checks apply to every object read. + +FLAGS: + {{range .VisibleFlags}}{{.}} + {{end}}`, +} + +func mainAtomic(ctx *cli.Context) error { + checkAtomicSyntax(ctx) + sse := newSSE(ctx) + sizeFn, err := generator.SizeFn(genOptions(ctx, "obj.size")...) + fatalIf(probe.NewError(err), "Invalid object size") + b := bench.Atomic{ + Common: getCommon(ctx, newGenSource(ctx, "obj.size")), + Keys: ctx.Int("objects"), + BlockSize: ctx.Int("block.size"), + CheckMD5: !ctx.Bool("no-etag-md5"), + ReadAfterWrite: ctx.BoolT("read-after-write"), + Prefix: ctx.String("prefix"), + PutWeight: ctx.Float64("put-distrib"), + GetWeight: ctx.Float64("get-distrib"), + StatWeight: ctx.Float64("stat-distrib"), + ListWeight: ctx.Float64("list-distrib"), + Seed: ctx.Uint64("seed"), + SizeFn: sizeFn, + GetOpts: minio.GetObjectOptions{ServerSideEncryption: sse}, + StatOpts: minio.StatObjectOptions{ServerSideEncryption: sse}, + } + return runBench(ctx, &b) +} + +func checkAtomicSyntax(ctx *cli.Context) { + if ctx.NArg() > 0 { + console.Fatal("Command takes no arguments") + } + if ctx.Int("objects") < 1 { + console.Fatal("At least one object must be tested") + } + if ctx.Int("block.size") < bench.AtomicMinBlockSize { + console.Fatalf("--block.size must be at least %d\n", bench.AtomicMinBlockSize) + } + total := 0.0 + for _, f := range []string{"put-distrib", "get-distrib", "stat-distrib", "list-distrib"} { + w := ctx.Float64(f) + if math.IsNaN(w) || math.IsInf(w, 0) || w < 0 { + console.Fatalf("--%s must be a finite number of at least 0\n", f) + } + total += w + } + if math.IsInf(total, 0) { + console.Fatal("the --*-distrib values add up to more than can be represented") + } + if ctx.Float64("list-distrib") == 0 && (ctx.Float64("put-distrib") == 0 || ctx.Float64("get-distrib")+ctx.Float64("stat-distrib") == 0) { + console.Fatal("set --list-distrib, or both --put-distrib and one of --get-distrib or --stat-distrib") + } + checkAnalyze(ctx) + checkBenchmark(ctx) +} diff --git a/cli/cli.go b/cli/cli.go index 10731367..4909fe39 100644 --- a/cli/cli.go +++ b/cli/cli.go @@ -94,6 +94,7 @@ func init() { listCmd, statCmd, versionedCmd, + atomicCmd, retentionCmd, multipartCmd, multipartPutCmd, diff --git a/cli/generator.go b/cli/generator.go index c36e8be7..c8e4d18e 100644 --- a/cli/generator.go +++ b/cli/generator.go @@ -45,6 +45,13 @@ var genFlags = []cli.Flag{ // newGenSource returns a new generator func newGenSource(ctx *cli.Context, sizeField string) func() generator.Source { + src, err := generator.NewFn(genOptions(ctx, sizeField)...) + fatalIf(probe.NewError(err), "Unable to create data generator") + return src +} + +// genOptions returns the generator options set by the generator flags. +func genOptions(ctx *cli.Context, sizeField string) []generator.Option { prefixSize := 8 if ctx.Bool("noprefix") { prefixSize = 0 @@ -94,10 +101,7 @@ func newGenSource(ctx *cli.Context, sizeField string) func() generator.Source { opts = append([]generator.Option{g.Apply()}, append(opts, generator.WithRandomSize(ctx.Bool("obj.randsize")))...) } - - src, err := generator.NewFn(opts...) - fatalIf(probe.NewError(err), "Unable to create data generator") - return src + return opts } // toSize converts a size indication to bytes. diff --git a/pkg/bench/atomic.go b/pkg/bench/atomic.go new file mode 100644 index 00000000..f77c6059 --- /dev/null +++ b/pkg/bench/atomic.go @@ -0,0 +1,1080 @@ +/* + * Warp (C) 2026 MinIO, Inc. + * + * 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 bench + +import ( + "bytes" + "context" + "crypto/md5" + "encoding/binary" + "encoding/hex" + "errors" + "fmt" + "hash" + "hash/crc32" + "io" + "math/rand" + "net/http" + "sort" + "strconv" + "strings" + "sync" + "time" + + "github.com/minio/minio-go/v7" +) + +const ( + atomicMagic = "WAT1" + atomicHeaderSize = 48 + atomicMetaKey = "Warp-Atomic" + + // AtomicMinBlockSize is the smallest block that still fits a stamp. + AtomicMinBlockSize = atomicHeaderSize +) + +// Atomic violation categories, reported as the prefix of Operation.Err. +const ( + AtomicTorn = "torn" + AtomicWrongKey = "wrong-key" + AtomicCorrupt = "corrupt" + AtomicLength = "length" + AtomicMissing = "missing" + AtomicMetaMissing = "meta-missing" + AtomicMetaMismatch = "meta-mismatch" + AtomicETagMismatch = "etag-mismatch" + AtomicETagMD5 = "etag-md5" + AtomicStale = "stale" + AtomicPhantom = "phantom" + AtomicListMissing = "list-missing" + AtomicListDeleted = "list-deleted" + AtomicListSize = "list-size" + AtomicListETag = "list-etag" + AtomicListStale = "list-stale" +) + +var atomicCRC = crc32.MakeTable(crc32.Castagnoli) + +// AtomicID identifies one PUT: the writer that issued it and its sequence +// number on that writer. +type AtomicID struct { + Writer, Gen uint64 +} + +func (a AtomicID) String() string { + return strconv.FormatUint(a.Writer, 16) + "." + strconv.FormatUint(a.Gen, 16) +} + +// ParseAtomicID parses the form produced by AtomicID.String. +func ParseAtomicID(s string) (AtomicID, bool) { + w, g, ok := strings.Cut(s, ".") + if !ok { + return AtomicID{}, false + } + wv, err1 := strconv.ParseUint(w, 16, 64) + gv, err2 := strconv.ParseUint(g, 16, 64) + if err1 != nil || err2 != nil { + return AtomicID{}, false + } + return AtomicID{Writer: wv, Gen: gv}, true +} + +// atomicStamp is what every block of an object carries, so any block +// identifies the PUT it came from without shared state. +type atomicStamp struct { + ID AtomicID + Size int64 + KeyHash uint64 +} + +func atomicKeyHash(key string) uint64 { + return uint64(crc32.Checksum([]byte(key), atomicCRC))<<32 | uint64(len(key)) +} + +func (s atomicStamp) fillBlock(dst []byte, idx uint64) { + var hdr [atomicHeaderSize]byte + copy(hdr[:4], atomicMagic) + binary.LittleEndian.PutUint64(hdr[4:], s.ID.Writer) + binary.LittleEndian.PutUint64(hdr[12:], s.ID.Gen) + binary.LittleEndian.PutUint64(hdr[20:], uint64(s.Size)) + binary.LittleEndian.PutUint64(hdr[28:], idx) + binary.LittleEndian.PutUint64(hdr[36:], s.KeyHash) + binary.LittleEndian.PutUint32(hdr[44:], crc32.Checksum(hdr[:44], atomicCRC)) + n := copy(dst, hdr[:]) + if n == len(dst) { + return + } + x := s.ID.Writer*0x9e3779b97f4a7c15 ^ s.ID.Gen*0xbf58476d1ce4e5b9 ^ idx*0x94d049bb133111eb + rest := dst[n:] + for len(rest) >= 8 { + x = splitmix64(x) + binary.LittleEndian.PutUint64(rest, x) + rest = rest[8:] + } + if len(rest) > 0 { + var tail [8]byte + binary.LittleEndian.PutUint64(tail[:], splitmix64(x)) + copy(rest, tail[:]) + } +} + +func splitmix64(x uint64) uint64 { + x += 0x9e3779b97f4a7c15 + x = (x ^ (x >> 30)) * 0xbf58476d1ce4e5b9 + x = (x ^ (x >> 27)) * 0x94d049bb133111eb + return x ^ (x >> 31) +} + +func parseAtomicHeader(b []byte) (st atomicStamp, idx uint64, ok bool) { + if len(b) < atomicHeaderSize || string(b[:4]) != atomicMagic { + return st, 0, false + } + if crc32.Checksum(b[:44], atomicCRC) != binary.LittleEndian.Uint32(b[44:]) { + return st, 0, false + } + st.ID.Writer = binary.LittleEndian.Uint64(b[4:]) + st.ID.Gen = binary.LittleEndian.Uint64(b[12:]) + st.Size = int64(binary.LittleEndian.Uint64(b[20:])) + st.KeyHash = binary.LittleEndian.Uint64(b[36:]) + return st, binary.LittleEndian.Uint64(b[28:]), true +} + +// atomicReader streams the stamped content of one PUT. +type atomicReader struct { + st atomicStamp + blockSize int64 + off int64 + buf []byte + bufIdx int64 + last time.Time +} + +func newAtomicReader(st atomicStamp, blockSize int) *atomicReader { + return &atomicReader{st: st, blockSize: int64(blockSize), buf: make([]byte, blockSize), bufIdx: -1} +} + +func (r *atomicReader) Read(p []byte) (int, error) { + if r.off >= r.st.Size { + return 0, io.EOF + } + bi := r.off / r.blockSize + blk := r.buf[:min(r.blockSize, r.st.Size-bi*r.blockSize)] + if bi != r.bufIdx { + r.st.fillBlock(blk, uint64(bi)) + r.bufIdx = bi + } + n := copy(p, blk[r.off-bi*r.blockSize:]) + r.off += int64(n) + r.last = time.Now() + return n, nil +} + +func (r *atomicReader) Seek(offset int64, whence int) (int64, error) { + switch whence { + case io.SeekStart: + case io.SeekCurrent: + offset += r.off + case io.SeekEnd: + offset += r.st.Size + default: + return 0, errors.New("invalid whence") + } + if offset < 0 { + return 0, errors.New("negative position") + } + r.off = offset + return offset, nil +} + +func (r *atomicReader) LastByte() *time.Time { + return &r.last +} + +// atomicVerifier checks a streamed body against the stamp in its first +// block and records the first fault found. +type atomicVerifier struct { + blockSize int + buf []byte + expect []byte + n int + idx uint64 + total int64 + first *atomicStamp + md5 hash.Hash + fault string + detail string +} + +func newAtomicVerifier(blockSize int) *atomicVerifier { + return &atomicVerifier{ + blockSize: blockSize, + buf: make([]byte, blockSize), + expect: make([]byte, blockSize), + md5: md5.New(), + } +} + +func (v *atomicVerifier) Write(p []byte) (int, error) { + v.md5.Write(p) + written := len(p) + for len(p) > 0 { + c := copy(v.buf[v.n:], p) + v.n += c + p = p[c:] + if v.n == v.blockSize { + v.checkBlock(v.buf) + v.n = 0 + } + } + return written, nil +} + +func (v *atomicVerifier) setFault(kind, detail string) { + if v.fault == "" { + v.fault, v.detail = kind, detail + } +} + +func (v *atomicVerifier) checkBlock(b []byte) { + idx := v.idx + v.idx++ + v.total += int64(len(b)) + if idx == 0 { + st, bi, ok := parseAtomicHeader(b) + if !ok || bi != 0 { + v.setFault(AtomicCorrupt, "no valid stamp in first block") + return + } + v.first = &st + } + if v.first == nil || v.fault != "" { + return + } + start := int64(idx) * int64(v.blockSize) + if start >= v.first.Size { + return + } + want := v.expect[:min(int64(len(b)), v.first.Size-start)] + v.first.fillBlock(want, idx) + if bytes.Equal(b[:len(want)], want) { + return + } + if st, bi, ok := parseAtomicHeader(b); ok && st.ID != v.first.ID { + v.setFault(AtomicTorn, fmt.Sprintf("block %d is block %d of %s (size %d), object starts as %s (size %d)", idx, bi, st.ID, st.Size, v.first.ID, v.first.Size)) + return + } + v.setFault(AtomicCorrupt, fmt.Sprintf("block %d of %s does not match its content", idx, v.first.ID)) +} + +// finish checks the trailing partial block and the total length. It returns +// the stamp of the object, or nil when no stamp could be read. +func (v *atomicVerifier) finish() *atomicStamp { + if v.n > 0 { + v.checkBlock(v.buf[:v.n]) + v.n = 0 + } + if v.first == nil { + if v.total == 0 { + v.setFault(AtomicLength, "empty body") + } + return nil + } + if v.total != v.first.Size { + v.setFault(AtomicLength, fmt.Sprintf("body of %s is %d bytes, stamp says %d", v.first.ID, v.total, v.first.Size)) + } + return v.first +} + +func (v *atomicVerifier) md5Hex() string { + return hex.EncodeToString(v.md5.Sum(nil)) +} + +type atomicWrite struct { + start, end time.Time + acked bool + etag string + size int64 + endpoint string +} + +type atomicKeyHistory struct { + writes map[AtomicID]*atomicWrite + // pruned holds, per writer, the highest generation dropped from writes. + // A writer issues one PUT at a time, so every lower generation of that + // writer on the key was superseded too. + pruned map[uint64]uint64 + // prunedETags maps the ETags of recently pruned writes to their IDs, so + // a listing that still shows one can be named stale. + prunedETags map[string]AtomicID + prunedETagOrder []string +} + +const atomicPrunedETags = 4096 + +// atomicHistory records the real-time interval of every PUT this process +// issued, to decide whether a read returned data that was already overwritten. +// Timestamps from other warp clients are never compared against it. +type atomicHistory struct { + mu sync.Mutex + owner uint64 + keys map[string]*atomicKeyHistory + reads map[uint64]time.Time + nextRead uint64 + sincePrun map[string]int +} + +func newAtomicHistory(owner uint64) *atomicHistory { + return &atomicHistory{ + owner: owner, + keys: make(map[string]*atomicKeyHistory), + reads: make(map[uint64]time.Time), + sincePrun: make(map[string]int), + } +} + +func (h *atomicHistory) key(k string) *atomicKeyHistory { + kh := h.keys[k] + if kh == nil { + kh = &atomicKeyHistory{ + writes: make(map[AtomicID]*atomicWrite), + pruned: make(map[uint64]uint64), + prunedETags: make(map[string]AtomicID), + } + h.keys[k] = kh + } + return kh +} + +func (h *atomicHistory) begin(k string, id AtomicID, start time.Time, size int64, endpoint string) { + h.mu.Lock() + h.key(k).writes[id] = &atomicWrite{start: start, size: size, endpoint: endpoint} + h.mu.Unlock() +} + +func (h *atomicHistory) endpoint(k string, id AtomicID) string { + h.mu.Lock() + defer h.mu.Unlock() + if w := h.key(k).writes[id]; w != nil { + return w.endpoint + } + return "" +} + +// finish resolves a PUT. A PUT that failed stays unresolved forever, since +// it may still have been applied. +func (h *atomicHistory) finish(k string, id AtomicID, end time.Time, acked bool, etag string) { + h.mu.Lock() + defer h.mu.Unlock() + w := h.key(k).writes[id] + if w == nil || !acked { + return + } + w.end, w.acked, w.etag = end, true, etag + h.sincePrun[k]++ + if h.sincePrun[k] >= 32 { + h.sincePrun[k] = 0 + h.pruneLocked(k, end) + } +} + +// readBegin registers a read and returns its ticket and start time. The start +// is taken under the lock so that no prune runs between the two. +func (h *atomicHistory) readBegin() (uint64, time.Time) { + h.mu.Lock() + defer h.mu.Unlock() + start := time.Now() + return h.readBeginLocked(start), start +} + +func (h *atomicHistory) readBeginAt(start time.Time) uint64 { + h.mu.Lock() + defer h.mu.Unlock() + return h.readBeginLocked(start) +} + +func (h *atomicHistory) readBeginLocked(start time.Time) uint64 { + h.nextRead++ + h.reads[h.nextRead] = start + return h.nextRead +} + +func (h *atomicHistory) readEnd(t uint64) { + h.mu.Lock() + delete(h.reads, t) + h.mu.Unlock() +} + +// supersededBefore returns the acknowledged write on kh with the latest +// start among those that finished before t, excluding skip. +func supersededBefore(kh *atomicKeyHistory, t time.Time, skip AtomicID) (latestID AtomicID, latest atomicWrite) { + for id, a := range kh.writes { + if id == skip || !a.acked || !a.end.Before(t) { + continue + } + if a.start.After(latest.start) { + latestID, latest = id, *a + } + } + return latestID, latest +} + +// pruneLocked drops writes that every current and future read must see as +// overwritten: acknowledged writes that ended before another acknowledged +// write started, where that write finished before the oldest active read. +func (h *atomicHistory) pruneLocked(k string, now time.Time) { + cutoff := now + for _, t := range h.reads { + if t.Before(cutoff) { + cutoff = t + } + } + kh := h.keys[k] + _, latest := supersededBefore(kh, cutoff, AtomicID{}) + s := latest.start + for id, w := range kh.writes { + if !w.acked || !w.end.Before(s) || id.Writer>>32 != h.owner { + continue + } + delete(kh.writes, id) + if id.Gen > kh.pruned[id.Writer] { + kh.pruned[id.Writer] = id.Gen + } + kh.prunedETags[w.etag] = id + kh.prunedETagOrder = append(kh.prunedETagOrder, w.etag) + if len(kh.prunedETagOrder) > atomicPrunedETags { + delete(kh.prunedETags, kh.prunedETagOrder[0]) + kh.prunedETagOrder = kh.prunedETagOrder[1:] + } + } +} + +func (h *atomicHistory) forget(k string) { + h.mu.Lock() + delete(h.keys, k) + delete(h.sincePrun, k) + h.mu.Unlock() +} + +// checkListed classifies a listing entry of key k by the ETag and size it +// shows. Entries whose ETag this process never saw acknowledged are not +// judged: they may belong to a PUT still in flight or to another client. +func (h *atomicHistory) checkListed(k, etag string, size int64, listStart time.Time) (kind, detail string) { + h.mu.Lock() + defer h.mu.Unlock() + kh := h.key(k) + for id, w := range kh.writes { + if !w.acked || w.etag != etag { + continue + } + if w.size != size { + return AtomicListSize, fmt.Sprintf("listed ETag %s of %s (%d bytes) with size %d", etag, id, w.size, size) + } + if aid, a := supersededBefore(kh, listStart, id); w.end.Before(a.start) { + return AtomicListStale, fmt.Sprintf("listed %s, but %s (written via %s) was acknowledged %s before the listing started", + id, aid, a.endpoint, listStart.Sub(a.end).Round(time.Microsecond)) + } + return "", "" + } + if id, ok := kh.prunedETags[etag]; ok { + return AtomicListStale, fmt.Sprintf("listed %s, which a later acknowledged PUT replaced", id) + } + return "", "" +} + +// etag returns the ETag the server acknowledged for id, if known. +func (h *atomicHistory) etag(k string, id AtomicID) (etag string, size int64, ok bool) { + h.mu.Lock() + defer h.mu.Unlock() + w := h.key(k).writes[id] + if w == nil || !w.acked { + return "", 0, false + } + return w.etag, w.size, true +} + +// check classifies a read of id on key k by a request that started at +// readStart. It returns an empty string when the read is admissible. +func (h *atomicHistory) check(k string, id AtomicID, readStart time.Time) (kind, detail string) { + if id.Writer>>32 != h.owner { + return "", "" + } + h.mu.Lock() + defer h.mu.Unlock() + kh := h.key(k) + w := kh.writes[id] + if w == nil { + if id.Gen <= kh.pruned[id.Writer] { + if aid, a := supersededBefore(kh, readStart, id); !a.start.IsZero() { + return AtomicStale, fmt.Sprintf("returned %s, but %s (written via %s) was acknowledged %s before the read started", + id, aid, a.endpoint, readStart.Sub(a.end).Round(time.Microsecond)) + } + return AtomicStale, fmt.Sprintf("returned %s, which a later acknowledged PUT replaced", id) + } + return AtomicPhantom, fmt.Sprintf("returned %s, which this client never wrote to %s", id, k) + } + if !w.acked { + return "", "" + } + if aid, a := supersededBefore(kh, readStart, id); w.end.Before(a.start) { + return AtomicStale, fmt.Sprintf("returned %s, but %s (written via %s) was acknowledged %s before the read started", + id, aid, a.endpoint, readStart.Sub(a.end).Round(time.Microsecond)) + } + return "", "" +} + +// Atomic overwrites a small set of keys concurrently and verifies that every +// GET and STAT returns exactly one whole acknowledged PUT: body, length, +// metadata and ETag all belonging to the same write, and not a write that +// was already overwritten before the read began. +type Atomic struct { + Common + + Keys int + BlockSize int + CheckMD5 bool + // ReadAfterWrite follows every acknowledged PUT with a GET of the same + // key from the same thread, which lands on the next host in --host. + ReadAfterWrite bool + Prefix string + + // Seed drives every thread's choice of key, operation and object size. + // Zero picks a random seed. The seed reproduces each thread's sequence of + // operations, not the interleaving of threads and servers. + Seed uint64 + // SizeFn picks an object size, following the configured size flags. + SizeFn func(rng *rand.Rand) int64 + + PutWeight, GetWeight, StatWeight, ListWeight float64 + + GetOpts minio.GetObjectOptions + StatOpts minio.StatObjectOptions + + hist *atomicHistory + counts map[string]int64 + countsMu sync.Mutex +} + +func (g *Atomic) keyName(i int) string { + return g.Prefix + "atomic-" + strconv.Itoa(i) +} + +func (g *Atomic) writerID(thread uint64) uint64 { + return uint64(g.ClientIdx)<<32 | thread +} + +func (g *Atomic) putOpts(id AtomicID) minio.PutObjectOptions { + opts := g.PutOpts + opts.DisableMultipart = true + opts.ContentType = "application/octet-stream" + opts.UserMetadata = map[string]string{atomicMetaKey: id.String()} + return opts +} + +func (g *Atomic) stamp(rng *rand.Rand, key string, id AtomicID) atomicStamp { + return atomicStamp{ID: id, Size: max(g.SizeFn(rng), int64(atomicHeaderSize)), KeyHash: atomicKeyHash(key)} +} + +// threadRNG returns the random source for one thread of this client. +func (g *Atomic) threadRNG(thread uint64) *rand.Rand { + return rand.New(rand.NewSource(int64(splitmix64(g.Seed ^ splitmix64(g.writerID(thread)))))) +} + +// Prepare creates the bucket and writes every key once so reads never race +// the first PUT. +func (g *Atomic) Prepare(ctx context.Context) error { + g.counts = make(map[string]int64) + if g.Keys < 1 { + return errors.New("at least one key is required") + } + if g.BlockSize < AtomicMinBlockSize { + return fmt.Errorf("block size must be at least %d bytes", AtomicMinBlockSize) + } + if g.PutWeight < 0 || g.GetWeight < 0 || g.StatWeight < 0 || g.ListWeight < 0 { + return errors.New("distributions cannot be negative") + } + if g.ListWeight == 0 && (g.PutWeight == 0 || g.GetWeight+g.StatWeight == 0) { + return errors.New("need a LIST distribution, or both PUT and GET or STAT distributions") + } + if g.SizeFn == nil { + return errors.New("no object size function") + } + for g.Seed == 0 { + g.Seed = rand.Uint64() + } + g.hist = newAtomicHistory(uint64(g.ClientIdx)) + if err := g.createEmptyBucket(ctx); err != nil { + return err + } + g.UpdateStatus(fmt.Sprint("Writing ", g.Keys, " keys, seed ", g.Seed)) + writer := g.writerID(1<<32 - 1) + rng := g.threadRNG(1<<32 - 1) + for i := range g.Keys { + id := AtomicID{Writer: writer, Gen: uint64(i + 1)} + key := g.keyName(i) + st := g.stamp(rng, key, id) + cl, done := g.Client() + start := time.Now() + g.hist.begin(key, id, start, st.Size, cl.EndpointURL().String()) + res, err := cl.PutObject(ctx, g.Bucket, key, newAtomicReader(st, g.BlockSize), st.Size, g.putOpts(id)) + if err == nil && i == 0 { + err = g.probe(ctx, cl, key, st, res.ETag) + } + done() + if err != nil { + return err + } + end := time.Now() + g.hist.finish(key, id, end, true, res.ETag) + g.prepareProgress(float64(i+1) / float64(g.Keys)) + } + return nil +} + +// probe fails the run early when the server drops the metadata or uses +// ETags that are not MD5, since every read would then stop at that check +// and never reach the staleness check. +func (g *Atomic) probe(ctx context.Context, cl *minio.Client, key string, st atomicStamp, etag string) error { + info, err := cl.StatObject(ctx, g.Bucket, key, g.StatOpts) + if err != nil { + return fmt.Errorf("stat after upload: %w", err) + } + if _, ok := ParseAtomicID(info.UserMetadata[atomicMetaKey]); !ok { + return fmt.Errorf("server did not return the %s metadata written with the object", atomicMetaKey) + } + if !g.CheckMD5 || g.PutOpts.ServerSideEncryption != nil { + return nil + } + h := md5.New() + if _, err := io.Copy(h, newAtomicReader(st, g.BlockSize)); err != nil { + return err + } + if sum := hex.EncodeToString(h.Sum(nil)); sum != etag { + return fmt.Errorf("server ETag %s is not the MD5 of the body (%s), rerun with --no-etag-md5", etag, sum) + } + return nil +} + +func (g *Atomic) count(kind string) { + g.countsMu.Lock() + g.counts[kind]++ + g.countsMu.Unlock() +} + +func (g *Atomic) violation(op *Operation, kind, detail string) { + g.count(kind) + op.Err = "atomic " + kind + ": " + op.File + ": " + detail + " [read via " + op.Endpoint + ", seed " + strconv.FormatUint(g.Seed, 10) + "]" + g.Error(op.Err) +} + +// Start will execute the main benchmark. +// Operations should begin executing when the start channel is closed. +func (g *Atomic) Start(ctx context.Context, wait chan struct{}) error { + var wg sync.WaitGroup + wg.Add(g.Concurrency) + c := g.Collector + if g.AutoTermDur > 0 { + ctx = c.AutoTerm(ctx, "", g.AutoTermScale, autoTermCheck, autoTermSamples, g.AutoTermDur) + } + nonTerm := context.Background() + total := g.PutWeight + g.GetWeight + g.StatWeight + g.ListWeight + for i := range g.Concurrency { + go func(i int) { + defer wg.Done() + rcv := c.Receiver() + rng := g.threadRNG(uint64(i)) + writer := g.writerID(uint64(i)) + var gen uint64 + done := ctx.Done() + <-wait + for { + select { + case <-done: + return + default: + } + if g.rpsLimit(ctx) != nil { + return + } + key := g.keyName(rng.Intn(g.Keys)) + pick := rng.Float64() * total + switch { + case pick < g.PutWeight: + gen++ + op := g.put(nonTerm, uint32(i), g.stamp(rng, key, AtomicID{Writer: writer, Gen: gen}), key) + rcv <- op + if g.ReadAfterWrite && op.Err == "" { + rcv <- g.get(nonTerm, uint32(i), key, op.Endpoint) + } + case pick < g.PutWeight+g.GetWeight: + rcv <- g.get(nonTerm, uint32(i), key, "") + case pick < g.PutWeight+g.GetWeight+g.StatWeight: + rcv <- g.stat(nonTerm, uint32(i), key) + default: + gen++ + g.listCycle(nonTerm, uint32(i), rng, AtomicID{Writer: writer, Gen: gen}, rcv) + } + } + }(i) + } + wg.Wait() + g.finalPass() + return nil +} + +func (g *Atomic) put(ctx context.Context, thread uint32, st atomicStamp, key string) Operation { + id := st.ID + r := newAtomicReader(st, g.BlockSize) + client, clDone := g.Client() + defer clDone() + op := Operation{ + OpType: http.MethodPut, + Thread: thread, + Size: st.Size, + File: key, + ObjPerOp: 1, + Endpoint: client.EndpointURL().String(), + } + op.Start = time.Now() + g.hist.begin(key, id, op.Start, st.Size, op.Endpoint) + res, err := client.PutObject(ctx, g.Bucket, key, r, st.Size, g.putOpts(id)) + op.End = time.Now() + op.LastByte = r.LastByte() + if err != nil { + op.Err = err.Error() + g.Error("upload error: ", err) + } + g.hist.finish(key, id, op.End, err == nil, res.ETag) + return op +} + +// clientAvoiding returns a client for a host other than avoid when there +// is more than one host. Rejected clients stay checked out until it returns, +// so a selector that favors idle hosts moves off the avoided one. +func (g *Atomic) clientAvoiding(avoid string) (*minio.Client, func()) { + var rejected []func() + defer func() { + for _, done := range rejected { + done() + } + }() + for range 16 { + client, clDone := g.Client() + if avoid == "" || client.EndpointURL().String() != avoid { + return client, clDone + } + rejected = append(rejected, clDone) + } + return g.Client() +} + +func (g *Atomic) get(ctx context.Context, thread uint32, key, avoid string) Operation { + client, clDone := g.clientAvoiding(avoid) + defer clDone() + op := Operation{ + OpType: http.MethodGet, + Thread: thread, + File: key, + ObjPerOp: 1, + Endpoint: client.EndpointURL().String(), + } + var t uint64 + t, op.Start = g.hist.readBegin() + defer g.hist.readEnd(t) + g.verifyGet(ctx, client, &op) + return op +} + +func (g *Atomic) verifyGet(ctx context.Context, client *minio.Client, op *Operation) { + fbr := firstByteRecorder{} + v := newAtomicVerifier(g.BlockSize) + obj, err := client.GetObject(ctx, g.Bucket, op.File, g.GetOpts) + var info minio.ObjectInfo + if err == nil { + fbr.r = obj + _, err = io.Copy(v, &fbr) + if err == nil { + info, err = obj.Stat() + } + obj.Close() + } + op.FirstByte = fbr.t + op.End = time.Now() + op.Size = v.total + int64(v.n) + if err != nil { + if minio.ToErrorResponse(err).Code == "NoSuchKey" { + g.violation(op, AtomicMissing, "key written before the run returned NoSuchKey") + return + } + if errors.Is(err, io.ErrUnexpectedEOF) && op.Size > 0 { + v.finish() + v.setFault(AtomicLength, "body ended early") + g.violation(op, v.fault, v.detail+": "+err.Error()) + return + } + op.Err = err.Error() + g.Error("download error: ", err) + return + } + st := v.finish() + if v.fault == "" && st.KeyHash != atomicKeyHash(op.File) { + v.setFault(AtomicWrongKey, "body "+st.ID.String()+" was written to a different key") + } + if v.fault != "" { + g.violation(op, v.fault, v.detail+" (ETag "+info.ETag+", meta "+info.UserMetadata[atomicMetaKey]+")") + return + } + metaID, ok := ParseAtomicID(info.UserMetadata[atomicMetaKey]) + if !ok { + g.violation(op, AtomicMetaMissing, "no "+atomicMetaKey+" metadata on body "+st.ID.String()) + return + } + if metaID != st.ID { + g.violation(op, AtomicMetaMismatch, "metadata says "+metaID.String()+", body is "+st.ID.String()+g.writtenVia(op.File, metaID, st.ID)) + return + } + if etag, _, ok := g.hist.etag(op.File, st.ID); ok && etag != info.ETag { + g.violation(op, AtomicETagMismatch, fmt.Sprintf("body is %s acked with ETag %s, GET returned ETag %s", st.ID, etag, info.ETag)) + return + } + if g.CheckMD5 && g.PutOpts.ServerSideEncryption == nil && v.md5Hex() != info.ETag { + g.violation(op, AtomicETagMD5, fmt.Sprintf("ETag %s, body %s has MD5 %s", info.ETag, st.ID, v.md5Hex())) + return + } + if kind, detail := g.hist.check(op.File, st.ID, op.Start); kind != "" { + g.violation(op, kind, detail+g.writtenVia(op.File, st.ID)) + } +} + +func (g *Atomic) writtenVia(key string, ids ...AtomicID) string { + var parts []string + for _, id := range ids { + if ep := g.hist.endpoint(key, id); ep != "" { + parts = append(parts, id.String()+" written via "+ep) + } + } + if len(parts) == 0 { + return "" + } + return " (" + strings.Join(parts, ", ") + ")" +} + +func (g *Atomic) stat(ctx context.Context, thread uint32, key string) Operation { + client, clDone := g.Client() + defer clDone() + op := Operation{ + OpType: "STAT", + Thread: thread, + File: key, + ObjPerOp: 1, + Endpoint: client.EndpointURL().String(), + } + var t uint64 + t, op.Start = g.hist.readBegin() + defer g.hist.readEnd(t) + info, err := client.StatObject(ctx, g.Bucket, key, g.StatOpts) + op.End = time.Now() + if err != nil { + if minio.ToErrorResponse(err).Code == "NoSuchKey" { + g.violation(&op, AtomicMissing, "key written before the run returned NoSuchKey") + return op + } + op.Err = err.Error() + g.Error("stat error: ", err) + return op + } + id, ok := ParseAtomicID(info.UserMetadata[atomicMetaKey]) + if !ok { + g.violation(&op, AtomicMetaMissing, "no "+atomicMetaKey+" metadata, ETag "+info.ETag) + return op + } + if etag, size, ok := g.hist.etag(key, id); ok { + if etag != info.ETag { + g.violation(&op, AtomicETagMismatch, fmt.Sprintf("metadata says %s acked with ETag %s, STAT returned ETag %s", id, etag, info.ETag)) + return op + } + if size != info.Size { + g.violation(&op, AtomicMetaMismatch, fmt.Sprintf("metadata says %s of %d bytes, STAT returned %d", id, size, info.Size)) + return op + } + } + if kind, detail := g.hist.check(key, id, op.Start); kind != "" { + g.violation(&op, kind, detail) + } + return op +} + +func (g *Atomic) listDir() string { + return g.Prefix + "atomic.list/" +} + +// listCycle checks list-after-write and list-after-delete on a new key, then +// that a listing of the overwritten keys shows no replaced PUT. Each listing +// goes to a host other than the one that took the preceding write. +func (g *Atomic) listCycle(ctx context.Context, thread uint32, rng *rand.Rand, id AtomicID, rcv chan<- Operation) { + dir := g.listDir() + strconv.FormatUint(id.Writer, 16) + "/" + key := dir + strconv.FormatUint(id.Gen, 16) + defer g.hist.forget(key) + + put := g.put(ctx, thread, g.stamp(rng, key, id), key) + rcv <- put + if put.Err != "" { + return + } + etag, size, _ := g.hist.etag(key, id) + op, objs, release := g.list(ctx, thread, dir, put.Endpoint) + if op.Err == "" { + switch obj, ok := objs[key]; { + case !ok: + g.violation(&op, AtomicListMissing, fmt.Sprintf("%s acknowledged via %s is not listed", id, put.Endpoint)) + case obj.Size != size: + g.violation(&op, AtomicListSize, fmt.Sprintf("%s acknowledged via %s with %d bytes, listed with %d", id, put.Endpoint, size, obj.Size)) + case obj.ETag != etag: + g.violation(&op, AtomicListETag, fmt.Sprintf("%s acknowledged via %s with ETag %s, listed with %s", id, put.Endpoint, etag, obj.ETag)) + } + } + release() + rcv <- op + + del := g.remove(ctx, thread, key) + rcv <- del + if del.Err != "" { + return + } + op, objs, release = g.list(ctx, thread, dir, del.Endpoint) + if _, ok := objs[key]; ok && op.Err == "" { + g.violation(&op, AtomicListDeleted, fmt.Sprintf("%s deleted via %s is still listed", id, del.Endpoint)) + } + release() + rcv <- op + + op, objs, release = g.list(ctx, thread, g.Prefix+"atomic-", "") + if op.Err == "" { + for i := range g.Keys { + k := g.keyName(i) + obj, ok := objs[k] + if !ok { + op.File = k + g.violation(&op, AtomicListMissing, "key written before the run is not listed") + break + } + if kind, detail := g.hist.checkListed(k, obj.ETag, obj.Size, op.Start); kind != "" { + op.File = k + g.violation(&op, kind, detail) + break + } + } + } + release() + rcv <- op +} + +// list returns the listing and a function that ends the read. Call it after +// every check against the listing, so no prune runs before they finish. +func (g *Atomic) list(ctx context.Context, thread uint32, prefix, avoid string) (Operation, map[string]minio.ObjectInfo, func()) { + client, clDone := g.clientAvoiding(avoid) + defer clDone() + op := Operation{ + OpType: "LIST", + Thread: thread, + File: prefix, + Endpoint: client.EndpointURL().String(), + } + var t uint64 + t, op.Start = g.hist.readBegin() + objs := make(map[string]minio.ObjectInfo) + for obj := range client.ListObjects(ctx, g.Bucket, minio.ListObjectsOptions{Prefix: prefix}) { + if obj.Err != nil { + op.Err = obj.Err.Error() + g.Error("list error: ", obj.Err) + break + } + if op.FirstByte == nil { + now := time.Now() + op.FirstByte = &now + } + objs[obj.Key] = obj + } + op.End = time.Now() + op.ObjPerOp = len(objs) + return op, objs, func() { g.hist.readEnd(t) } +} + +func (g *Atomic) remove(ctx context.Context, thread uint32, key string) Operation { + client, clDone := g.Client() + defer clDone() + op := Operation{ + OpType: http.MethodDelete, + Thread: thread, + File: key, + ObjPerOp: 1, + Endpoint: client.EndpointURL().String(), + } + op.Start = time.Now() + err := client.RemoveObject(ctx, g.Bucket, key, minio.RemoveObjectOptions{}) + op.End = time.Now() + if err != nil { + op.Err = err.Error() + g.Error("delete error: ", err) + } + return op +} + +// finalPass reads every key once all writers have stopped, so the last +// acknowledged state is checked without any write in flight. +func (g *Atomic) finalPass() { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute) + defer cancel() + g.UpdateStatus("Verifying final state of all keys") + for i := range g.Keys { + client, clDone := g.Client() + op := Operation{File: g.keyName(i), Endpoint: client.EndpointURL().String(), Start: time.Now()} + g.verifyGet(ctx, client, &op) + clDone() + } + g.summary() +} + +func (g *Atomic) summary() { + g.countsMu.Lock() + defer g.countsMu.Unlock() + if len(g.counts) == 0 { + return + } + kinds := make([]string, 0, len(g.counts)) + for k := range g.counts { + kinds = append(kinds, k) + } + sort.Strings(kinds) + var sb strings.Builder + for _, k := range kinds { + fmt.Fprintf(&sb, " %s=%d", k, g.counts[k]) + } + g.Error("atomic violations:" + sb.String() + " seed=" + strconv.FormatUint(g.Seed, 10)) +} + +// Cleanup deletes the keys this benchmark wrote and nothing else. +func (g *Atomic) Cleanup(ctx context.Context) { + cl, done := g.Client() + defer done() + for i := range g.Keys { + if err := cl.RemoveObject(ctx, g.Bucket, g.keyName(i), minio.RemoveObjectOptions{}); err != nil { + g.Error("cleanup: ", err) + } + } + g.deleteAllInBucket(ctx, strings.TrimSuffix(g.listDir(), "/")) +} diff --git a/pkg/bench/atomic_test.go b/pkg/bench/atomic_test.go new file mode 100644 index 00000000..57056feb --- /dev/null +++ b/pkg/bench/atomic_test.go @@ -0,0 +1,358 @@ +/* + * Warp (C) 2026 MinIO, Inc. + * + * 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 bench + +import ( + "crypto/md5" + "encoding/hex" + "io" + "math/rand" + "slices" + "strconv" + "testing" + "time" + + "github.com/minio/minio-go/v7" +) + +const testBlock = 4096 + +func atomicBody(t *testing.T, st atomicStamp) []byte { + t.Helper() + b, err := io.ReadAll(newAtomicReader(st, testBlock)) + if err != nil { + t.Fatal(err) + } + if int64(len(b)) != st.Size { + t.Fatalf("reader produced %d bytes, want %d", len(b), st.Size) + } + return b +} + +func verify(body []byte, chunk int) *atomicVerifier { + v := newAtomicVerifier(testBlock) + for len(body) > 0 { + n := min(chunk, len(body)) + v.Write(body[:n]) + body = body[n:] + } + v.finish() + return v +} + +func TestAtomicVerifier(t *testing.T) { + key := "atomic-0" + a := atomicStamp{ID: AtomicID{Writer: 1, Gen: 7}, Size: 10*testBlock + 123, KeyHash: atomicKeyHash(key)} + b := atomicStamp{ID: AtomicID{Writer: 2, Gen: 9}, Size: 10*testBlock + 123, KeyHash: atomicKeyHash(key)} + bodyA, bodyB := atomicBody(t, a), atomicBody(t, b) + + spliced := append(append([]byte{}, bodyA[:5*testBlock]...), bodyB[5*testBlock:]...) + unaligned := append(append([]byte{}, bodyA[:5*testBlock+100]...), bodyB[5*testBlock+100:]...) + flipped := append([]byte{}, bodyA...) + flipped[3*testBlock+999] ^= 1 + longer := atomicBody(t, atomicStamp{ID: a.ID, Size: a.Size + testBlock, KeyHash: a.KeyHash}) + + tests := []struct { + name string + body []byte + want string + }{ + {"whole", bodyA, ""}, + {"spliced at block", spliced, AtomicTorn}, + {"spliced inside block", unaligned, AtomicCorrupt}, + {"bit flip", flipped, AtomicCorrupt}, + {"truncated", bodyA[:len(bodyA)-1], AtomicLength}, + {"truncated at block", bodyA[:4*testBlock], AtomicLength}, + {"extended", append(append([]byte{}, bodyA...), 0), AtomicLength}, + {"stamp says shorter", longer[:a.Size], ""}, + {"no stamp", make([]byte, 2*testBlock), AtomicCorrupt}, + {"empty", nil, AtomicLength}, + } + for _, tc := range tests { + for _, chunk := range []int{1, 1000, testBlock, 1 << 20} { + if tc.name == "stamp says shorter" { + continue + } + v := verify(tc.body, chunk) + if v.fault != tc.want { + t.Errorf("%s (chunk %d): fault %q (%s), want %q", tc.name, chunk, v.fault, v.detail, tc.want) + } + } + } + + v := verify(longer[:a.Size], testBlock) + if v.fault != AtomicLength { + t.Errorf("prefix of a longer write: fault %q, want %q", v.fault, AtomicLength) + } + + v = verify(bodyA, 777) + sum := md5.Sum(bodyA) + if v.md5Hex() != hex.EncodeToString(sum[:]) { + t.Errorf("md5 %s, want %x", v.md5Hex(), sum) + } + if st := v.first; st == nil || st.ID != a.ID || st.KeyHash != atomicKeyHash(key) { + t.Errorf("stamp %+v, want %+v", st, a) + } + if v.first.KeyHash == atomicKeyHash("atomic-1") { + t.Error("distinct keys hash equal") + } +} + +func TestAtomicReaderSeek(t *testing.T) { + st := atomicStamp{ID: AtomicID{Writer: 3, Gen: 1}, Size: 3*testBlock + 5} + want := atomicBody(t, st) + r := newAtomicReader(st, testBlock) + io.CopyN(io.Discard, r, 2*testBlock+17) + if _, err := r.Seek(testBlock-3, io.SeekStart); err != nil { + t.Fatal(err) + } + got, _ := io.ReadAll(r) + if string(got) != string(want[testBlock-3:]) { + t.Fatal("content after seek differs") + } +} + +func TestAtomicStamp(t *testing.T) { + st := atomicStamp{ID: AtomicID{Writer: 5<<32 | 3, Gen: 42}, Size: 1234, KeyHash: 99} + blk := make([]byte, testBlock) + st.fillBlock(blk, 6) + got, idx, ok := parseAtomicHeader(blk) + if !ok || got != st || idx != 6 { + t.Fatalf("parsed %+v idx %d ok %v", got, idx, ok) + } + blk[10] ^= 0x80 + if _, _, ok := parseAtomicHeader(blk); ok { + t.Fatal("corrupted header accepted") + } + id, ok := ParseAtomicID(st.ID.String()) + if !ok || id != st.ID { + t.Fatalf("round trip %v -> %v", st.ID, id) + } + for _, s := range []string{"", "1", "x.1", "1.y", "1.2.3"} { + if _, ok := ParseAtomicID(s); ok { + t.Errorf("ParseAtomicID(%q) accepted", s) + } + } +} + +func TestAtomicHistory(t *testing.T) { + base := time.Unix(1000, 0) + at := func(ms int) time.Time { return base.Add(time.Duration(ms) * time.Millisecond) } + const owner = 5 + w := func(thread, gen uint64) AtomicID { return AtomicID{Writer: owner<<32 | thread, Gen: gen} } + const k = "key" + + h := newAtomicHistory(owner) + h.begin(k, w(1, 1), at(0), 10, "") + h.finish(k, w(1, 1), at(10), true, "e1") + h.begin(k, w(2, 1), at(20), 10, "") + h.finish(k, w(2, 1), at(30), true, "e2") + h.begin(k, w(3, 1), at(5), 10, "") + h.finish(k, w(3, 1), at(50), true, "e3") + h.begin(k, w(4, 1), at(60), 10, "") + h.begin(k, w(5, 1), at(1), 10, "") + h.finish(k, w(5, 1), at(2), false, "") + + tests := []struct { + name string + id AtomicID + readStart time.Time + want string + }{ + {"latest before any overwrite", w(1, 1), at(15), ""}, + {"overwritten before read", w(1, 1), at(31), AtomicStale}, + {"read overlaps overwrite", w(1, 1), at(25), ""}, + {"concurrent writes, either may win", w(3, 1), at(31), ""}, + {"overlapping write wins", w(2, 1), at(51), ""}, + {"in flight", w(4, 1), at(0), ""}, + {"failed write may land", w(5, 1), at(1000), ""}, + {"never written", w(1, 99), at(40), AtomicPhantom}, + {"other client", AtomicID{Writer: 6<<32 | 1, Gen: 1}, at(40), ""}, + } + for _, tc := range tests { + if got, detail := h.check(k, tc.id, tc.readStart); got != tc.want { + t.Errorf("%s: %q (%s), want %q", tc.name, got, detail, tc.want) + } + } + if etag, size, ok := h.etag(k, w(2, 1)); !ok || etag != "e2" || size != 10 { + t.Errorf("etag %q %d %v", etag, size, ok) + } + if _, _, ok := h.etag(k, w(4, 1)); ok { + t.Error("etag of in-flight write reported") + } +} + +func TestAtomicHistoryPrune(t *testing.T) { + base := time.Unix(1000, 0) + at := func(ms int) time.Time { return base.Add(time.Duration(ms) * time.Millisecond) } + const owner = 1 + id := func(gen uint64) AtomicID { return AtomicID{Writer: owner<<32 | 7, Gen: gen} } + const k = "key" + + h := newAtomicHistory(owner) + slow := h.readBeginAt(at(15)) + for g := uint64(1); g <= 100; g++ { + ms := int(g) * 10 + h.begin(k, id(g), at(ms), 1, "") + h.finish(k, id(g), at(ms+5), true, "") + } + if got, _ := h.check(k, id(1), at(15)); got != "" { + t.Fatalf("read active since before the overwrite flagged %q", got) + } + if got, _ := h.check(k, id(2), at(15)); got != "" { + t.Fatalf("write overlapping an active read flagged %q", got) + } + h.readEnd(slow) + + h.begin(k, id(101), at(1010), 1, "") + h.finish(k, id(101), at(1015), true, "") + for g := uint64(102); g <= 140; g++ { + ms := int(g) * 10 + h.begin(k, id(g), at(ms), 1, "") + h.finish(k, id(g), at(ms+5), true, "") + } + if n := len(h.keys[k].writes); n > 40 { + t.Fatalf("%d writes retained after pruning", n) + } + for _, g := range []uint64{1, 50, 100} { + if got, _ := h.check(k, id(g), at(2000)); got != AtomicStale { + t.Errorf("pruned gen %d: %q, want %q", g, got, AtomicStale) + } + } + if got, _ := h.check(k, id(140), at(2000)); got != "" { + t.Errorf("latest write flagged %q", got) + } + if got, _ := h.check(k, id(500), at(2000)); got != AtomicPhantom { + t.Errorf("unissued gen: %q, want %q", got, AtomicPhantom) + } +} + +func TestAtomicCheckListed(t *testing.T) { + base := time.Unix(1000, 0) + at := func(ms int) time.Time { return base.Add(time.Duration(ms) * time.Millisecond) } + const owner = 2 + id := func(gen uint64) AtomicID { return AtomicID{Writer: owner<<32 | 1, Gen: gen} } + const k = "key" + + h := newAtomicHistory(owner) + h.begin(k, id(1), at(0), 100, "http://a") + h.finish(k, id(1), at(10), true, "e1") + h.begin(k, id(2), at(20), 200, "http://b") + h.finish(k, id(2), at(30), true, "e2") + h.begin(k, id(3), at(40), 300, "http://c") + + tests := []struct { + name string + etag string + size int64 + start time.Time + want string + }{ + {"current", "e2", 200, at(35), ""}, + {"replaced before listing", "e1", 100, at(35), AtomicListStale}, + {"listing overlaps overwrite", "e1", 100, at(25), ""}, + {"size of another write", "e2", 100, at(35), AtomicListSize}, + {"in flight, ETag unknown", "e3", 300, at(45), ""}, + } + for _, tc := range tests { + if got, detail := h.checkListed(k, tc.etag, tc.size, tc.start); got != tc.want { + t.Errorf("%s: %q (%s), want %q", tc.name, got, detail, tc.want) + } + } + + for g := uint64(4); g <= 80; g++ { + ms := int(g) * 100 + h.begin(k, id(g), at(ms), 1, "") + h.finish(k, id(g), at(ms+5), true, "e"+strconv.FormatUint(g, 10)) + } + if _, ok := h.keys[k].writes[id(1)]; ok { + t.Fatal("write 1 was not pruned") + } + if got, _ := h.checkListed(k, "e1", 100, at(100000)); got != AtomicListStale { + t.Errorf("pruned ETag: %q, want %q", got, AtomicListStale) + } + h.forget(k) + if _, ok := h.keys[k]; ok { + t.Error("forget kept the key") + } +} + +func TestAtomicSeed(t *testing.T) { + sizes := func(g *Atomic, thread uint64) []int64 { + rng := g.threadRNG(thread) + out := make([]int64, 8) + for i := range out { + out[i] = g.stamp(rng, "k", AtomicID{}).Size + } + return out + } + sizeFn := func(rng *rand.Rand) int64 { return 1 + rng.Int63n(1<<20) } + a := &Atomic{Seed: 42, SizeFn: sizeFn} + b := &Atomic{Seed: 42, SizeFn: sizeFn} + if !slices.Equal(sizes(a, 3), sizes(b, 3)) { + t.Fatal("same seed and thread gave different sequences") + } + for name, other := range map[string][]int64{ + "thread": sizes(a, 4), + "seed": sizes(&Atomic{Seed: 43, SizeFn: sizeFn}, 3), + "client": sizes(&Atomic{Seed: 42, SizeFn: sizeFn, Common: Common{ClientIdx: 1}}, 3), + } { + if slices.Equal(sizes(a, 3), other) { + t.Errorf("different %s gave the same sequence", name) + } + } + if got := (&Atomic{SizeFn: func(*rand.Rand) int64 { return 1 }}).stamp(rand.New(rand.NewSource(1)), "k", AtomicID{}).Size; got != atomicHeaderSize { + t.Errorf("size below the stamp was not raised: %d", got) + } +} + +func TestAtomicClientAvoiding(t *testing.T) { + hosts := []string{"a:9000", "b:9000"} + clients := make([]*minio.Client, 0, len(hosts)) + for _, host := range hosts { + cl, err := minio.New(host, &minio.Options{}) + if err != nil { + t.Fatal(err) + } + clients = append(clients, cl) + } + running := []int{0, 5} + leastRunning := func() (*minio.Client, func()) { + idx := 0 + if running[1] < running[0] { + idx = 1 + } + running[idx]++ + return clients[idx], func() { running[idx]-- } + } + g := &Atomic{Common: Common{Client: leastRunning}} + avoid := clients[0].EndpointURL().String() + cl, done := g.clientAvoiding(avoid) + if cl.EndpointURL().String() == avoid { + t.Fatal("returned the avoided host") + } + done() + if running[0] != 0 || running[1] != 5 { + t.Errorf("clients not released: running %v", running) + } + if cl, done := g.clientAvoiding(""); cl != clients[0] { + t.Error("with nothing to avoid, did not take the idle host") + } else { + done() + } +} diff --git a/pkg/generator/generator.go b/pkg/generator/generator.go index 61228574..784df7f8 100644 --- a/pkg/generator/generator.go +++ b/pkg/generator/generator.go @@ -123,6 +123,18 @@ func New(opts ...Option) (Source, error) { return options.src(options) } +// SizeFn returns the object size picker the options configure, so a caller +// can draw sizes from its own random source. +func SizeFn(opts ...Option) (func(rng *rand.Rand) int64, error) { + options := defaultOptions() + for _, ofn := range opts { + if err := ofn(&options); err != nil { + return nil, err + } + } + return options.getSize, nil +} + // NewFn return data source. func NewFn(opts ...Option) (func() Source, error) { options := defaultOptions() @@ -175,7 +187,7 @@ func randASCIIBytes(dst []byte, rng *rand.Rand) { func GetExpRandSize(rng *rand.Rand, minSize, maxSize int64) int64 { if maxSize-minSize < 10 { if maxSize-minSize <= 0 { - return 0 + return maxSize } return 1 + minSize + rng.Int63n(maxSize-minSize) } diff --git a/pkg/generator/generator_test.go b/pkg/generator/generator_test.go index bdd35f62..f418ecdf 100644 --- a/pkg/generator/generator_test.go +++ b/pkg/generator/generator_test.go @@ -19,6 +19,7 @@ package generator import ( "io" + "math/rand" "testing" ) @@ -166,3 +167,35 @@ func BenchmarkWithRandomData(b *testing.B) { }) } } + +func TestSizeFn(t *testing.T) { + fixed, err := SizeFn(WithRandomData().Apply(), WithSize(1000)) + if err != nil { + t.Fatal(err) + } + rng := rand.New(rand.NewSource(1)) + if got := fixed(rng); got != 1000 { + t.Fatalf("fixed size %d", got) + } + random, err := SizeFn(WithRandomData().Apply(), WithMinMaxSize(500, 5000), WithRandomSize(true)) + if err != nil { + t.Fatal(err) + } + equal, err := SizeFn(WithRandomData().Apply(), WithMinMaxSize(4096, 4096), WithRandomSize(true)) + if err != nil { + t.Fatal(err) + } + if got := equal(rng); got != 4096 { + t.Errorf("equal bounds gave size %d, want 4096", got) + } + a, b := rand.New(rand.NewSource(7)), rand.New(rand.NewSource(7)) + for range 100 { + x, y := random(a), random(b) + if x != y { + t.Fatal("same random source gave different sizes") + } + if x < 500 || x > 5000 { + t.Fatalf("size %d outside 500-5000", x) + } + } +}