Files
zcli-playground/valkey/chaos.go
T
2026-06-11 10:52:17 +02:00

823 lines
20 KiB
Go

package main
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"log"
"math/rand"
"net/http"
"os"
"os/signal"
"sort"
"strconv"
"strings"
"sync"
"sync/atomic"
"syscall"
"time"
"github.com/redis/go-redis/v9"
)
// ---------------------------------------------------------------------------
// Chaos / soak engine — shared core for local TUI mode and Zerops worker mode
// ---------------------------------------------------------------------------
const (
chaosCounterKey = "vk:chaos:counter"
chaosLedgerKeyFn = "vk:chaos:ledger:"
chaosBloatKeyFn = "vk:chaos:bloat:"
ledgerCapacity = 5000
tickInterval = time.Second
maxRecentLat = 200
maxGlobalLat = 2000
maxSparkSamples = 60
maxSeriesPoints = 120
maxEvents = 500
httpBatchSize = 200
httpBatchInterval = 100 * time.Millisecond
httpBatchTimeout = 10 * time.Second
httpQueueSize = 10_000
)
// ---- ledger -------------------------------------------------------------
type ledgerEntry struct {
seq int64
value string
ackedAt time.Time
writer string
// Which Valkey instance acked the write (role + short run_id).
// "role:master/abc12345" identifies the master at write time; lets us
// tell post-mortem whether a MISSING was acked by a node that later
// stopped being authoritative.
writeBackend string
}
type ledger struct {
mu sync.RWMutex
entries []ledgerEntry
seq atomic.Int64
}
func newLedger() *ledger { return &ledger{} }
func (l *ledger) next() int64 { return l.seq.Add(1) }
func (l *ledger) append(e ledgerEntry) {
l.mu.Lock()
defer l.mu.Unlock()
if len(l.entries) >= ledgerCapacity {
drop := ledgerCapacity / 10
copy(l.entries, l.entries[drop:])
l.entries = l.entries[:len(l.entries)-drop]
}
l.entries = append(l.entries, e)
}
func (l *ledger) sample() (ledgerEntry, bool) {
l.mu.RLock()
defer l.mu.RUnlock()
if len(l.entries) == 0 {
return ledgerEntry{}, false
}
return l.entries[rand.Intn(len(l.entries))], true
}
func (l *ledger) last() (ledgerEntry, bool) {
l.mu.RLock()
defer l.mu.RUnlock()
n := len(l.entries)
if n == 0 {
return ledgerEntry{}, false
}
return l.entries[n-1], true
}
func (l *ledger) snapshot() []ledgerEntry {
l.mu.RLock()
defer l.mu.RUnlock()
out := make([]ledgerEntry, len(l.entries))
copy(out, l.entries)
return out
}
func ledgerKey(seq int64) string { return chaosLedgerKeyFn + strconv.FormatInt(seq, 10) }
func bloatKey(seq int64) string { return chaosBloatKeyFn + strconv.FormatInt(seq, 10) }
// ---- shared chaos context (workers + sink) ------------------------------
type chaosCtx struct {
ledger *ledger
counter atomic.Int64
sink EventSink
startAt time.Time
}
type outageRec struct {
start, end time.Time
opsLost int64
}
// ---- event sink abstraction --------------------------------------------
//
// Workers ship their events to a remote collector over HTTP via httpSink,
// which batches and POSTs them; the collector serves the dashboard. Events
// flow as opResultMsg / subRecvMsg / ledgerCheckMsg / eventMsg shapes.
type EventSink interface {
OpResult(opResultMsg)
SubRecv(subRecvMsg)
LedgerCheck(ledgerCheckMsg)
Event(eventMsg)
Close() error
}
// ---- HTTP sink (Zerops worker) -----------------------------------------
type httpSink struct {
url string
containerID string
client *http.Client
queue chan wireEvent
done chan struct{}
wg sync.WaitGroup
}
func newHTTPSink(url, containerID string) *httpSink {
s := &httpSink{
url: strings.TrimRight(url, "/") + "/events",
containerID: containerID,
client: &http.Client{Timeout: httpBatchTimeout},
queue: make(chan wireEvent, httpQueueSize),
done: make(chan struct{}),
}
s.wg.Add(1)
go s.runBatcher()
return s
}
func (s *httpSink) runBatcher() {
defer s.wg.Done()
defer close(s.done)
batch := make([]wireEvent, 0, httpBatchSize)
t := time.NewTicker(httpBatchInterval)
defer t.Stop()
flush := func() {
if len(batch) == 0 {
return
}
s.post(batch)
batch = batch[:0]
}
for {
select {
case ev, ok := <-s.queue:
if !ok {
flush()
return
}
batch = append(batch, ev)
if len(batch) >= httpBatchSize {
flush()
}
case <-t.C:
flush()
}
}
}
func (s *httpSink) post(events []wireEvent) {
body, err := json.Marshal(wireBatch{ContainerID: s.containerID, Events: events})
if err != nil {
log.Printf("httpSink: marshal: %v", err)
return
}
// minimal retry: best-effort, drops on persistent failure
for attempt := 0; attempt < 3; attempt++ {
req, _ := http.NewRequest("POST", s.url, bytes.NewReader(body))
req.Header.Set("Content-Type", "application/json")
resp, err := s.client.Do(req)
if err == nil {
io.Copy(io.Discard, resp.Body)
resp.Body.Close()
if resp.StatusCode/100 == 2 {
return
}
}
time.Sleep(time.Duration(100*(attempt+1)) * time.Millisecond)
}
}
func (s *httpSink) enqueue(e wireEvent) {
if e.At.IsZero() {
e.At = time.Now()
}
select {
case s.queue <- e:
default:
// queue full — drop event rather than block the worker
}
}
func (s *httpSink) OpResult(m opResultMsg) {
we := wireEvent{
Type: "op",
Worker: m.worker,
Op: m.op,
At: m.at,
OK: m.err == nil,
DurUS: m.dur.Microseconds(),
}
if m.err != nil {
we.Err = m.err.Error()
}
s.enqueue(we)
}
func (s *httpSink) SubRecv(m subRecvMsg) {
s.enqueue(wireEvent{Type: "sub", Worker: m.worker, Seq: m.seq, At: m.at})
}
func (s *httpSink) LedgerCheck(m ledgerCheckMsg) {
s.enqueue(wireEvent{Type: "ledger", Worker: m.worker, Seq: m.seq, Kind: m.kind, Detail: m.detail, At: time.Now()})
}
func (s *httpSink) Event(m eventMsg) {
s.enqueue(wireEvent{Type: "info", At: m.at, Level: m.level, Text: m.text})
}
func (s *httpSink) Close() error {
close(s.queue)
s.wg.Wait()
return nil
}
// ---- wire format (worker → collector) ----------------------------------
type wireBatch struct {
ContainerID string `json:"container"`
Events []wireEvent `json:"events"`
}
type wireEvent struct {
Type string `json:"type"` // op | sub | ledger | info
Worker string `json:"worker,omitempty"`
At time.Time `json:"at"`
// op
Op string `json:"op,omitempty"`
OK bool `json:"ok,omitempty"`
DurUS int64 `json:"dur_us,omitempty"`
Err string `json:"err,omitempty"`
// sub / ledger
Seq int64 `json:"seq,omitempty"`
Kind string `json:"kind,omitempty"`
Detail string `json:"detail,omitempty"`
// info
Level string `json:"level,omitempty"`
Text string `json:"text,omitempty"`
}
// ---- workers ------------------------------------------------------------
type worker struct {
cfg WorkerConfig
rdb *redis.Client
cc *chaosCtx
bloatBuf []byte // pre-allocated random payload for bloater role
}
func validateWorker(w WorkerConfig) error {
if w.Name == "" {
return fmt.Errorf("missing name")
}
if w.Host == "" || w.Port == 0 {
return fmt.Errorf("missing host/port")
}
switch w.Role {
case "reader", "writer", "readwriter", "publisher", "subscriber", "bloater":
default:
return fmt.Errorf("unknown role %q", w.Role)
}
if (w.Role == "publisher" || w.Role == "subscriber") && w.Channel == "" {
return fmt.Errorf("role %s requires channel", w.Role)
}
if w.Role != "subscriber" && w.Interval <= 0 {
return fmt.Errorf("missing or invalid interval")
}
if w.Role == "bloater" && w.ValueBytes <= 0 {
return fmt.Errorf("role bloater requires valueBytes > 0")
}
return nil
}
func (w *worker) run(ctx context.Context) {
switch w.cfg.Role {
case "writer":
w.runTicked(ctx, w.opWrite)
case "readwriter":
w.runTicked(ctx, w.opReadWrite)
case "reader":
w.runTicked(ctx, w.opRead)
case "publisher":
w.runTicked(ctx, w.opPublish)
case "subscriber":
w.runSubscriber(ctx)
case "bloater":
w.runTicked(ctx, w.opBloat)
}
}
func (w *worker) runTicked(ctx context.Context, op func(context.Context) (string, error)) {
t := time.NewTicker(w.cfg.Interval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
}
opCtx, cancel := context.WithTimeout(ctx, maxDur(w.cfg.Interval, 5*time.Second))
start := time.Now()
name, err := op(opCtx)
cancel()
// If the worker's parent context has been cancelled (SIGTERM /
// Ctrl+C / TUI exit), this op's error is shutdown noise — don't
// flag it as a failure or trigger an outage. We're done.
if ctx.Err() != nil {
return
}
w.cc.sink.OpResult(opResultMsg{
worker: w.cfg.Name,
op: name,
err: err,
dur: time.Since(start),
at: time.Now(),
})
}
}
// All read/write paths route through Client.Conn() so a single TCP socket
// carries the entire op + a probe INFO. HAProxy pins one frontend↔backend
// per TCP connection in tcp mode, so the probe reports the exact Valkey
// instance that handled the user command - critical for diagnosing whether
// a MISSING was answered by the new master, a lagging replica, or a node
// mid-shutdown.
func (w *worker) opWrite(ctx context.Context) (string, error) {
conn := w.rdb.Conn()
defer conn.Close()
if _, err := conn.Incr(ctx, chaosCounterKey).Result(); err != nil {
return "INCR", err
}
w.cc.counter.Add(1)
seq := w.cc.ledger.next()
val := fmt.Sprintf("%d:%d", seq, rand.Int63())
if err := conn.Set(ctx, ledgerKey(seq), val, 0).Err(); err != nil {
return "SET", err
}
backend := backendIdent(ctx, conn)
w.cc.ledger.append(ledgerEntry{
seq: seq, value: val, ackedAt: time.Now(),
writer: w.cfg.Name, writeBackend: backend,
})
return "WRITE", nil
}
func (w *worker) opReadWrite(ctx context.Context) (string, error) {
conn := w.rdb.Conn()
defer conn.Close()
if _, err := conn.Incr(ctx, chaosCounterKey).Result(); err != nil {
return "INCR", err
}
w.cc.counter.Add(1)
seq := w.cc.ledger.next()
val := fmt.Sprintf("%d:%d", seq, rand.Int63())
if err := conn.Set(ctx, ledgerKey(seq), val, 0).Err(); err != nil {
return "SET", err
}
writeBackend := backendIdent(ctx, conn)
w.cc.ledger.append(ledgerEntry{
seq: seq, value: val, ackedAt: time.Now(),
writer: w.cfg.Name, writeBackend: writeBackend,
})
// Read-after-write on the SAME conn → same backend (the master). Any
// mismatch here means the master we just acked from doesn't even agree
// with itself, which would be a real bug rather than replica lag.
got, err := conn.Get(ctx, ledgerKey(seq)).Result()
if err != nil {
return "GET", err
}
if got != val {
return "GET", fmt.Errorf("read-after-write mismatch seq=%d backend=%s", seq, writeBackend)
}
return "RW", nil
}
func (w *worker) opRead(ctx context.Context) (string, error) {
e, ok := w.cc.ledger.sample()
if !ok {
if _, err := w.rdb.Ping(ctx).Result(); err != nil {
return "PING", err
}
return "PING", nil
}
conn := w.rdb.Conn()
defer conn.Close()
got, err := conn.Get(ctx, ledgerKey(e.seq)).Result()
if err == redis.Nil {
readBackend := backendIdent(ctx, conn)
w.cc.sink.LedgerCheck(ledgerCheckMsg{
worker: w.cfg.Name, seq: e.seq, kind: "missing",
detail: fmt.Sprintf("read-from=%s acked %s ago by %s on %s",
readBackend,
time.Since(e.ackedAt).Round(time.Millisecond),
e.writer, e.writeBackend),
})
return "GET", nil
}
if err != nil {
return "GET", err
}
if got != e.value {
readBackend := backendIdent(ctx, conn)
w.cc.sink.LedgerCheck(ledgerCheckMsg{
worker: w.cfg.Name, seq: e.seq, kind: "mismatch",
detail: fmt.Sprintf("read-from=%s expected=%s got=%s (acked by %s on %s)",
readBackend, trim(e.value, 30), trim(got, 30),
e.writer, e.writeBackend),
})
} else {
w.cc.sink.LedgerCheck(ledgerCheckMsg{worker: w.cfg.Name, seq: e.seq, kind: "ok"})
}
return "GET", nil
}
// backendIdent returns "role:run_id_short" for the Valkey instance behind
// the given conn. INFO is best-effort: on probe failure (e.g. conn already
// broken after a NIL response from a node mid-shutdown), returns "?" so
// the missing-detail at least notes that the conn died on the way out.
func backendIdent(ctx context.Context, conn *redis.Conn) string {
info, err := conn.Info(ctx, "server", "replication").Result()
if err != nil {
return "?"
}
var role, runID string
for _, line := range strings.Split(info, "\n") {
line = strings.TrimRight(line, "\r")
switch {
case strings.HasPrefix(line, "run_id:"):
runID = strings.TrimPrefix(line, "run_id:")
case strings.HasPrefix(line, "role:"):
role = strings.TrimPrefix(line, "role:")
}
}
if len(runID) > 8 {
runID = runID[:8]
}
if role == "" && runID == "" {
return "?"
}
return role + "/" + runID
}
func (w *worker) opPublish(ctx context.Context) (string, error) {
seq := w.cc.ledger.next()
payload := fmt.Sprintf("pub-%d", seq)
if err := w.rdb.Publish(ctx, w.cfg.Channel, payload).Err(); err != nil {
return "PUB", err
}
return "PUB", nil
}
// opBloat writes a fixed-size random payload to a fresh key under a
// dedicated "bloat" namespace, no TTL. Each tick adds another key, so
// Valkey memory usage grows monotonically — useful for watching how
// failover, replication, and RAM scaling behave under increasing pressure.
// Bloat keys are intentionally NOT added to the durability ledger; they're
// memory load, not consistency probes.
func (w *worker) opBloat(ctx context.Context) (string, error) {
seq := w.cc.ledger.next()
if err := w.rdb.Set(ctx, bloatKey(seq), w.bloatBuf, 0).Err(); err != nil {
return "BLOAT", err
}
return "BLOAT", nil
}
func (w *worker) runSubscriber(ctx context.Context) {
sub := w.rdb.Subscribe(ctx, w.cfg.Channel)
defer sub.Close()
if _, err := sub.Receive(ctx); err != nil {
// shutdown-time cancellation isn't a real subscribe failure
if ctx.Err() != nil {
return
}
w.cc.sink.OpResult(opResultMsg{
worker: w.cfg.Name, op: "SUBSCRIBE", err: err, at: time.Now(),
})
return
}
ch := sub.Channel()
for {
select {
case <-ctx.Done():
return
case msg, ok := <-ch:
if !ok {
return
}
seq, perr := strconv.ParseInt(strings.TrimPrefix(msg.Payload, "pub-"), 10, 64)
if perr != nil {
continue
}
w.cc.sink.SubRecv(subRecvMsg{worker: w.cfg.Name, seq: seq, at: time.Now()})
}
}
}
// ---- event message types (worker → sink) -------------------------------
type opResultMsg struct {
worker string
op string
err error
dur time.Duration
at time.Time
}
type subRecvMsg struct {
worker string
seq int64
at time.Time
}
type ledgerCheckMsg struct {
worker string
seq int64
kind string
detail string
}
type eventMsg struct {
at time.Time
level string
text string
}
// ---- shared stats types (mirrored server-side by the collector) ---------
type uiStats struct {
role string
addr string
ops, ok, fail int64
consecFail int
inOutage bool
outageStart time.Time
outageOpsLost int64
outages []outageRec
recentLat []time.Duration
prevOps int64
opsHist []float64
verified, missing, mismatch int64
msgs int64
gaps int64
longestGap time.Duration
lastSeq int64
lastSeqAt time.Time
}
type eventLine struct {
at time.Time
level string
text string
}
// ---- runWorker (Zerops worker container) -------------------------------
func runWorker(ws []WorkerConfig, collectorURL string) int {
for _, w := range ws {
if err := validateWorker(w); err != nil {
fmt.Fprintf(os.Stderr, "worker %q: %v\n", w.Name, err)
return 2
}
}
containerID := containerIdent()
log.Printf("worker container_id=%s collector=%s workers=%d", containerID, collectorURL, len(ws))
ctx, cancel := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer cancel()
cc := &chaosCtx{
ledger: newLedger(),
startAt: time.Now(),
}
cc.sink = newHTTPSink(collectorURL, containerID)
workers := buildWorkers(ws, cc)
defer func() {
for _, w := range workers {
_ = w.rdb.Close()
}
}()
// announce ourselves so the collector registers the worker rows immediately
for _, w := range workers {
cc.sink.Event(eventMsg{
at: time.Now(), level: "info",
text: fmt.Sprintf("hello %s/%s role=%s addr=%s:%d db=%d",
containerID, w.cfg.Name, w.cfg.Role, w.cfg.Host, w.cfg.Port, w.cfg.DB),
})
}
var wg sync.WaitGroup
for _, w := range workers {
wg.Add(1)
go func(w *worker) {
defer wg.Done()
w.run(ctx)
}(w)
}
<-ctx.Done()
log.Printf("worker shutting down, draining sink…")
waitWorkers(&wg, 5*time.Second)
// Best-effort durability sweep against this container's own DB before exit.
fd := finalSweep(workers, cc)
cc.sink.Event(eventMsg{
at: time.Now(), level: "info",
text: fmt.Sprintf("[%s] FINAL acked=%d verified=%d missing=%d mismatch=%d local_acked=%d server_counter=%d",
containerID, fd.acked, fd.verified, fd.missing, fd.mismatch, fd.localAcked, fd.serverCounter),
})
_ = cc.sink.Close()
if fd.missing > 0 || fd.mismatch > 0 {
return 1
}
return 0
}
func buildWorkers(ws []WorkerConfig, cc *chaosCtx) []*worker {
workers := make([]*worker, 0, len(ws))
for _, wc := range ws {
opts := optionsFromConn(wc.Host, wc.Port, wc.Password, wc.DB, wc.TLS)
w := &worker{cfg: wc, rdb: redis.NewClient(opts), cc: cc}
if wc.Role == "bloater" && wc.ValueBytes > 0 {
w.bloatBuf = make([]byte, wc.ValueBytes)
rand.Read(w.bloatBuf)
}
workers = append(workers, w)
}
return workers
}
func waitWorkers(wg *sync.WaitGroup, max time.Duration) {
done := make(chan struct{})
go func() { wg.Wait(); close(done) }()
select {
case <-done:
case <-time.After(max):
}
}
func containerIdent() string {
if v := os.Getenv("ZEROPS_Number"); v != "" {
return "vk-" + v
}
if v := os.Getenv("HOSTNAME"); v != "" {
return v
}
return "local-0"
}
// ---- final report -------------------------------------------------------
type finalDur struct {
acked, verified, missing, mismatch int64
missingSeqs []int64
serverCounter, localAcked int64
}
func finalSweep(workers []*worker, cc *chaosCtx) finalDur {
var verifier *worker
for _, w := range workers {
if w.cfg.Role == "reader" || w.cfg.Role == "readwriter" {
verifier = w
break
}
}
if verifier == nil && len(workers) > 0 {
verifier = workers[0]
}
if verifier == nil {
return finalDur{}
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
entries := cc.ledger.snapshot()
var verified, missing, mismatch int64
var missingSeqs []int64
for _, e := range entries {
got, err := verifier.rdb.Get(ctx, ledgerKey(e.seq)).Result()
if err == redis.Nil || err != nil {
missing++
missingSeqs = append(missingSeqs, e.seq)
continue
}
if got == e.value {
verified++
} else {
mismatch++
}
}
fd := finalDur{
acked: int64(len(entries)),
verified: verified,
missing: missing,
mismatch: mismatch,
missingSeqs: missingSeqs,
localAcked: cc.counter.Load(),
}
if v, err := verifier.rdb.Get(ctx, chaosCounterKey).Int64(); err == nil {
fd.serverCounter = v
}
return fd
}
// ---- helpers ------------------------------------------------------------
func appendBoundedDur(s []time.Duration, v time.Duration, max int) []time.Duration {
s = append(s, v)
if len(s) > max {
s = s[len(s)-max:]
}
return s
}
func appendBoundedFloat(s []float64, v float64, max int) []float64 {
s = append(s, v)
if len(s) > max {
s = s[len(s)-max:]
}
return s
}
func percentileLat(values []time.Duration, pct int) time.Duration {
if len(values) == 0 {
return 0
}
cp := make([]time.Duration, len(values))
copy(cp, values)
sort.Slice(cp, func(i, j int) bool { return cp[i] < cp[j] })
idx := pct * len(cp) / 100
if idx >= len(cp) {
idx = len(cp) - 1
}
return cp[idx]
}
func maxFloat(s []float64) float64 {
m := 0.0
for _, v := range s {
if v > m {
m = v
}
}
return m
}
func maxDur(a, b time.Duration) time.Duration {
if a > b {
return a
}
return b
}
func trim(s string, n int) string {
if len(s) <= n {
return s
}
return s[:n] + "…"
}