Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
136 changes: 136 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 <category>:`. 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.
Expand Down
162 changes: 162 additions & 0 deletions cli/atomic.go
Original file line number Diff line number Diff line change
@@ -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 <http://www.gnu.org/licenses/>.
*/

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 <category>:".

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")...)
Comment thread
harshavardhana marked this conversation as resolved.
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)
}
1 change: 1 addition & 0 deletions cli/cli.go
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,7 @@ func init() {
listCmd,
statCmd,
versionedCmd,
atomicCmd,
retentionCmd,
multipartCmd,
multipartPutCmd,
Expand Down
12 changes: 8 additions & 4 deletions cli/generator.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down
Loading
Loading