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
8 changes: 8 additions & 0 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,14 @@ make soak # 5 minutes of load with a node restarted every 30 seconds
make soak SOAK=4h # the full run the numbers in the docs come from
```

Benchmarks are opt-in too. `make memtier` needs `memtier_benchmark` and `redis-cli` on `PATH`
(`brew install memtier_benchmark redis`):

```sh
make bench # Go benchmarks for the store, RESP, and a three-node cluster
make memtier BENCH_MODE=cluster # 60 seconds of memtier load; memory, durable, or cluster
```

### Pre-commit Hooks (optional)

```sh
Expand Down
24 changes: 22 additions & 2 deletions Makefile
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
.PHONY: all setup clean build test lint lint-fix vet coverage race soak notice licenses verify-release-artifacts
.PHONY: all setup clean build test lint lint-fix vet coverage race soak bench memtier notice licenses verify-release-artifacts

# How long `make soak` runs for. The full run behind the numbers in the docs is SOAK=4h.
SOAK ?= 5m
Expand All @@ -7,6 +7,16 @@ SOAK ?= 5m
# back instantly, which is the crash-loop condition the heap numbers in the docs were taken under.
SOAK_DOWN ?= 10s

# How many times `make bench` repeats each benchmark. Comparing two builds with benchstat wants
# ten or so; one is enough to see a number.
BENCH_COUNT ?= 1

# What `make memtier` starts: memory, durable, or cluster. The performance page has all three.
BENCH_MODE ?= durable

# How many seconds `make memtier` keeps the load on.
BENCH_TIME ?= 60

all: vet lint test build

notice:
Expand Down Expand Up @@ -54,4 +64,14 @@ race:
# test binary and falls back to the current directory. -timeout 0 leaves the deadline to the
# test, so raising SOAK never means recalculating a second number.
soak:
@go test . ./internal/cluster -run TestSoak -soak $(SOAK) -soak-down $(SOAK_DOWN) -timeout 0 -v
@go test ./pkg/kvs ./internal/cluster -run TestSoak -soak $(SOAK) -soak-down $(SOAK_DOWN) -timeout 0 -v

# Go benchmarks for the store, the RESP path, and a three-node cluster. -run '^$' skips the
# tests in the same packages, which would otherwise run first and cost more than the benchmarks.
bench:
@go test ./pkg/kvs ./internal/server ./internal/cluster -run '^$$' -bench . -benchmem -count $(BENCH_COUNT)

# Real load from memtier_benchmark against a kvs built from this tree. Needs memtier_benchmark
# and redis-cli on PATH.
memtier: build
@./scripts/memtier.sh $(BENCH_MODE) $(BENCH_TIME)
5 changes: 5 additions & 0 deletions changes/unreleased/Added-20260929-234842.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
kind: Added
body: '`INFO` now reports `total_connections_received`, `total_commands_processed`, and `used_memory`, and `make bench` and `make memtier` measure the store, the RESP path, and a three-node cluster.'
time: 2026-09-29T23:48:42.733703+09:00
custom:
Issue: "270"
5 changes: 5 additions & 0 deletions changes/unreleased/Fixed-20260929-234843.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
kind: Fixed
body: '`make soak` runs again; it named the repository root, which has held no Go files since the layout change.'
time: 2026-09-29T23:48:43.760181+09:00
custom:
Issue: "270"
22 changes: 22 additions & 0 deletions internal/cluster/bench_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
package cluster

import (
"strconv"
"testing"
)

// BenchmarkClusterPut is one write through consensus on a three-node cluster on this machine:
// a round trip to a majority and an fsync on each of them. It is the number the durable
// single-node figure is compared against.
func BenchmarkClusterPut(b *testing.B) {
leader := waitForLeader(b, startCluster(b, 3))

i := 0
for b.Loop() {
key := "bench-" + strconv.Itoa(i%1000)
if err := leader.store.Put(key, "value"); err != nil {
b.Fatalf("Put(%q) error = %v", key, err)
}
i++
}
}
18 changes: 9 additions & 9 deletions internal/cluster/cluster_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -144,7 +144,7 @@ func newStore() *kvs.Store {
return store
}

func startCluster(t *testing.T, size int) []*testNode {
func startCluster(t testing.TB, size int) []*testNode {
t.Helper()

nodes := make([]*testNode, 0, size)
Expand All @@ -171,7 +171,7 @@ func startCluster(t *testing.T, size int) []*testNode {
return nodes
}

func (n *testNode) start(t *testing.T, bootstrap bool) {
func (n *testNode) start(t testing.TB, bootstrap bool) {
t.Helper()

n.store = newStore()
Expand All @@ -195,7 +195,7 @@ func (n *testNode) start(t *testing.T, bootstrap bool) {
})
}

func (n *testNode) stop(t *testing.T) {
func (n *testNode) stop(t testing.TB) {
t.Helper()

if err := n.node.Close(); err != nil {
Expand All @@ -206,13 +206,13 @@ func (n *testNode) stop(t *testing.T) {

// restart brings the node back on the same address with the same log, which is what a process
// that was killed and started again looks like to the rest of the cluster.
func (n *testNode) restart(t *testing.T) {
func (n *testNode) restart(t testing.TB) {
t.Helper()

n.start(t, false)
}

func (n *testNode) mustEventuallyHold(t *testing.T, key, want string) {
func (n *testNode) mustEventuallyHold(t testing.TB, key, want string) {
t.Helper()

eventually(t, n.id+" to hold "+key, func() bool {
Expand All @@ -222,7 +222,7 @@ func (n *testNode) mustEventuallyHold(t *testing.T, key, want string) {
})
}

func waitForLeader(t *testing.T, nodes []*testNode) *testNode {
func waitForLeader(t testing.TB, nodes []*testNode) *testNode {
t.Helper()

var leader *testNode
Expand All @@ -244,7 +244,7 @@ func waitForLeader(t *testing.T, nodes []*testNode) *testNode {
// waitForWrite offers the write to every node until one takes it, and reports which did. It is
// what a client retrying a rejected write does, so timing it measures the outage a client sees
// rather than the moment the cluster privately agreed on a leader.
func waitForWrite(t *testing.T, nodes []*testNode, key, value string) *testNode {
func waitForWrite(t testing.TB, nodes []*testNode, key, value string) *testNode {
t.Helper()

var taken *testNode
Expand Down Expand Up @@ -276,7 +276,7 @@ func without(nodes []*testNode, excluded *testNode) []*testNode {

// reserveAddr picks a free port and lets go of it, so Raft can bind it and a restarted node can
// bind it again.
func reserveAddr(t *testing.T) string {
func reserveAddr(t testing.TB) string {
t.Helper()

var lc net.ListenConfig
Expand All @@ -295,7 +295,7 @@ func reserveAddr(t *testing.T) string {

// eventually waits for consensus to settle. Elections take as long as they take, so polling is
// the only honest way to wait for one.
func eventually(t *testing.T, what string, check func() bool) {
func eventually(t testing.TB, what string, check func() bool) {
t.Helper()

deadline := time.Now().Add(30 * time.Second)
Expand Down
12 changes: 11 additions & 1 deletion internal/server/resp.go
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,12 @@ type RESPServer struct {
cursors respCursors
scripts respScripts
lastID atomic.Int64
// commands counts every command that reached its handler, for INFO. Queued commands count
// when EXEC runs them and calls made inside a script count too, which is how Redis counts.
//
// ponytail: one shared atomic per command; per-connection counters summed at INFO if it ever
// shows in a profile.
commands atomic.Int64

mu sync.Mutex
conns map[*respConn]struct{}
Expand Down Expand Up @@ -208,7 +214,6 @@ func (s *RESPServer) newConn(netConn net.Conn) *respConn {
server: s,
netConn: netConn,
writer: resp.NewWriter(netConn),
id: s.lastID.Add(1),
authed: s.password == "",
pushes: make(chan respPush, respPushDepth),
done: make(chan struct{}),
Expand All @@ -229,6 +234,9 @@ func (s *RESPServer) track(conn *respConn) bool {
return false
}

// Only once admitted, so INFO's total_connections_received leaves out the refused ones, as
// Redis does.
conn.id = s.lastID.Add(1)
s.conns[conn] = struct{}{}

return true
Expand Down Expand Up @@ -579,6 +587,8 @@ func (c *respConn) dispatch(args [][]byte) error {
return c.writer.WriteSimple("QUEUED")
}

c.server.commands.Add(1)

return cmd.run(c, args)
}

Expand Down
80 changes: 80 additions & 0 deletions internal/server/resp_bench_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,80 @@
package server

import (
"strconv"
"testing"

"github.com/redis/go-redis/v9"
)

// These go through a real client library over loopback TCP, so they price the whole RESP path:
// parsing, dispatch, the store, and the reply. The network is the machine's own and costs little.

// benchKeys is a fixed keyspace overwritten in place, so a run measures the command path rather
// than a map that keeps growing. The store's benchmarks use the same shape.
const benchKeys = 1000

const benchValue = "value"

// benchKeyNames is built once so that strconv is not what the benchmarks end up measuring.
var benchKeyNames = func() []string {
names := make([]string, benchKeys)
for i := range names {
names[i] = "bench-" + strconv.Itoa(i)
}

return names
}()

func benchKey(i int) string { return benchKeyNames[i%benchKeys] }

func BenchmarkRESPSet(b *testing.B) {
client := newGoRedisClient(b, &redis.Options{})
ctx := b.Context()

i := 0
for b.Loop() {
if err := client.Set(ctx, benchKey(i), benchValue, 0).Err(); err != nil {
b.Fatalf("Set(%q) error = %v", benchKey(i), err)
}
i++
}
}

func BenchmarkRESPGet(b *testing.B) {
client := newGoRedisClient(b, &redis.Options{})
ctx := b.Context()

for i := range benchKeys {
if err := client.Set(ctx, benchKey(i), benchValue, 0).Err(); err != nil {
b.Fatalf("Set(%q) error = %v", benchKey(i), err)
}
}

i := 0
for b.Loop() {
if err := client.Get(ctx, benchKey(i)).Err(); err != nil {
b.Fatalf("Get(%q) error = %v", benchKey(i), err)
}
i++
}
}

// BenchmarkRESPSetParallel shares one client, and so its connection pool, across goroutines,
// which is how a service using go-redis sends writes.
func BenchmarkRESPSetParallel(b *testing.B) {
client := newGoRedisClient(b, &redis.Options{})
ctx := b.Context()

// RunParallel, unlike b.Loop, would otherwise time the server starting.
b.ResetTimer()
b.RunParallel(func(pb *testing.PB) {
for i := 0; pb.Next(); i++ {
if err := client.Set(ctx, benchKey(i), benchValue, 0).Err(); err != nil {
b.Errorf("Set(%q) error = %v", benchKey(i), err)

return
}
}
})
}
2 changes: 1 addition & 1 deletion internal/server/resp_client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ import (
// newGoRedisClient starts a RESP server and connects a real client library to it. Hand
// written bytes cover the wire format elsewhere; this exists to check the parts a client
// drives on its own, such as protocol negotiation, pipelining, and subscribe bookkeeping.
func newGoRedisClient(t *testing.T, opts *redis.Options) *redis.Client {
func newGoRedisClient(t testing.TB, opts *redis.Options) *redis.Client {
t.Helper()

var lc net.ListenConfig
Expand Down
19 changes: 19 additions & 0 deletions internal/server/resp_commands.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"fmt"
"maps"
"runtime"
"runtime/metrics"
"strconv"
"strings"
"sync"
Expand Down Expand Up @@ -484,6 +485,14 @@ func (c *respConn) cmdInfo(_ [][]byte) error {
"# Clients",
"connected_clients:" + strconv.Itoa(c.server.connCount()),
"",
"# Stats",
// Connection ids are handed out one per accepted connection, so the last one is the count.
"total_connections_received:" + strconv.FormatInt(c.server.lastID.Load(), 10),
"total_commands_processed:" + strconv.FormatInt(c.server.commands.Load(), 10),
"",
"# Memory",
"used_memory:" + strconv.FormatUint(respUsedMemory(), 10),
"",
}
lines = append(lines, c.server.replicationInfo()...)
lines = append(lines,
Expand All @@ -503,6 +512,16 @@ func (c *respConn) cmdInfo(_ [][]byte) error {
return c.writer.WriteBulkString(strings.Join(lines, respCRLF) + respCRLF)
}

// respUsedMemory is the heap held by objects, live ones and dead ones not yet swept, the closest
// Go has to what Redis reports as used_memory. runtime/metrics reads it without stopping the
// world, which ReadMemStats would.
func respUsedMemory() uint64 {
sample := []metrics.Sample{{Name: "/memory/classes/heap/objects:bytes"}}
metrics.Read(sample)

return sample[0].Value.Uint64()
}

func (c *respConn) cmdConfig(args [][]byte) error {
switch respUpper(args[1]) {
case "GET":
Expand Down
51 changes: 49 additions & 2 deletions internal/server/resp_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -280,8 +280,55 @@ func TestRESPClientAndInfo(t *testing.T) {
t.Fatalf("INFO = %q, want it to contain %q", info, want)
}
}
if !strings.Contains(info, "connected_clients:1") {
t.Fatalf("INFO = %q, want connected_clients:1", info)
// Five CLIENT calls and the INFO itself, which is counted before it reports.
for _, want := range []string{
"connected_clients:1\r\n", "total_connections_received:1\r\n",
"total_commands_processed:6\r\n", "used_memory:",
} {
if !strings.Contains(info, want) {
t.Fatalf("INFO = %q, want it to contain %q", info, want)
}
}
}

// A queued command is counted when EXEC runs it, not when it is queued, so a transaction is not
// counted twice.
func TestRESPInfoCountsQueuedCommandsOnExec(t *testing.T) {
client := newRESPClient(t, kvs.NewStore())

client.do("+OK"+respCRLF, "MULTI")
client.do("+QUEUED"+respCRLF, "SET", "k", "v")
client.do("+QUEUED"+respCRLF, "GET", "k")
client.do("*2"+respCRLF+"+OK"+respCRLF+"$1"+respCRLF+"v"+respCRLF, "EXEC")

client.send("INFO")
if info := client.readBulk(); !strings.Contains(info, "total_commands_processed:5\r\n") {
t.Fatalf("INFO = %q, want total_commands_processed:5", info)
}
}

// A redis.call is a command of its own, so a script that makes one counts twice.
func TestRESPInfoCountsScriptCalls(t *testing.T) {
client := newRESPClient(t, kvs.NewStore())

client.do("+OK"+respCRLF, respCmdEval, "return redis.call('SET','k','v')", "0")

client.send("INFO")
if info := client.readBulk(); !strings.Contains(info, "total_commands_processed:3\r\n") {
t.Fatalf("INFO = %q, want total_commands_processed:3", info)
}
}

// A command refused before its handler runs is not a processed command.
func TestRESPInfoSkipsRejectedCommands(t *testing.T) {
client := newRESPClient(t, kvs.NewStore())

client.do("-ERR unknown command 'NOPE'"+respCRLF, "NOPE")
client.do("-ERR wrong number of arguments for 'echo' command"+respCRLF, "ECHO")

client.send("INFO")
if info := client.readBulk(); !strings.Contains(info, "total_commands_processed:1\r\n") {
t.Fatalf("INFO = %q, want total_commands_processed:1", info)
}
}

Expand Down
Loading
Loading