1521 lines
37 KiB
Go
1521 lines
37 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/charmbracelet/bubbles/viewport"
|
||
tea "github.com/charmbracelet/bubbletea"
|
||
"github.com/charmbracelet/lipgloss"
|
||
"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 don't know whether their events go to an in-process Bubble Tea
|
||
// program (local TUI) or a remote collector over HTTP. Two implementations:
|
||
//
|
||
// teaSink wraps *tea.Program — used by --chaos
|
||
// httpSink batches and POSTs to a collector — used by --worker
|
||
//
|
||
// Both deliver the same opResultMsg / subRecvMsg / ledgerCheckMsg shapes.
|
||
|
||
type EventSink interface {
|
||
OpResult(opResultMsg)
|
||
SubRecv(subRecvMsg)
|
||
LedgerCheck(ledgerCheckMsg)
|
||
Event(eventMsg)
|
||
Close() error
|
||
}
|
||
|
||
// ---- tea sink (local TUI) ----------------------------------------------
|
||
|
||
type teaSink struct{ p *tea.Program }
|
||
|
||
func (s *teaSink) OpResult(m opResultMsg) { s.p.Send(m) }
|
||
func (s *teaSink) SubRecv(m subRecvMsg) { s.p.Send(m) }
|
||
func (s *teaSink) LedgerCheck(m ledgerCheckMsg) { s.p.Send(m) }
|
||
func (s *teaSink) Event(m eventMsg) { s.p.Send(m) }
|
||
func (s *teaSink) Close() error { return nil }
|
||
|
||
// ---- 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()})
|
||
}
|
||
}
|
||
}
|
||
|
||
// ---- TUI message types --------------------------------------------------
|
||
|
||
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
|
||
}
|
||
|
||
type tickMsg time.Time
|
||
|
||
// ---- TUI model ----------------------------------------------------------
|
||
|
||
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
|
||
}
|
||
|
||
type chaosModel struct {
|
||
workers []WorkerConfig
|
||
stats map[string]*uiStats
|
||
|
||
events []eventLine
|
||
eventVP viewport.Model
|
||
|
||
p95Series []float64
|
||
okPerSecSeries []float64
|
||
failPerSecSeries []float64
|
||
|
||
globalLat []time.Duration
|
||
prevTotalOps int64
|
||
prevTotalFail int64
|
||
|
||
width, height int
|
||
startAt time.Time
|
||
quitting bool
|
||
}
|
||
|
||
func newChaosModel(ws []WorkerConfig, startAt time.Time) *chaosModel {
|
||
stats := make(map[string]*uiStats, len(ws))
|
||
for _, w := range ws {
|
||
stats[w.Name] = &uiStats{
|
||
role: w.Role,
|
||
addr: fmt.Sprintf("%s:%d", w.Host, w.Port),
|
||
}
|
||
}
|
||
vp := viewport.New(80, 8)
|
||
vp.SetContent(dimStyle.Render("(no events yet)"))
|
||
return &chaosModel{
|
||
workers: ws,
|
||
stats: stats,
|
||
eventVP: vp,
|
||
startAt: startAt,
|
||
}
|
||
}
|
||
|
||
func (m *chaosModel) Init() tea.Cmd {
|
||
return tea.Tick(tickInterval, func(t time.Time) tea.Msg { return tickMsg(t) })
|
||
}
|
||
|
||
func (m *chaosModel) Update(msg tea.Msg) (tea.Model, tea.Cmd) {
|
||
switch msg := msg.(type) {
|
||
case tea.WindowSizeMsg:
|
||
m.width = msg.Width
|
||
m.height = msg.Height
|
||
m.recomputeLayout()
|
||
return m, nil
|
||
case tea.KeyMsg:
|
||
switch msg.String() {
|
||
case "q", "ctrl+c", "esc":
|
||
m.quitting = true
|
||
return m, tea.Quit
|
||
case "up", "k":
|
||
m.eventVP.LineUp(1)
|
||
case "down", "j":
|
||
m.eventVP.LineDown(1)
|
||
case "G", "end":
|
||
m.eventVP.GotoBottom()
|
||
case "g", "home":
|
||
m.eventVP.GotoTop()
|
||
case "pgup", "b":
|
||
m.eventVP.HalfViewUp()
|
||
case "pgdown", "f", " ":
|
||
m.eventVP.HalfViewDown()
|
||
}
|
||
return m, nil
|
||
case opResultMsg:
|
||
m.applyOpResult(msg)
|
||
return m, nil
|
||
case subRecvMsg:
|
||
m.applySubRecv(msg)
|
||
return m, nil
|
||
case ledgerCheckMsg:
|
||
m.applyLedgerCheck(msg)
|
||
return m, nil
|
||
case eventMsg:
|
||
m.appendEvent(msg.level, msg.at, msg.text)
|
||
return m, nil
|
||
case tickMsg:
|
||
m.onTick()
|
||
return m, tea.Tick(tickInterval, func(t time.Time) tea.Msg { return tickMsg(t) })
|
||
}
|
||
return m, nil
|
||
}
|
||
|
||
func (m *chaosModel) applyOpResult(msg opResultMsg) {
|
||
s, ok := m.stats[msg.worker]
|
||
if !ok {
|
||
return
|
||
}
|
||
s.ops++
|
||
if msg.err != nil {
|
||
s.fail++
|
||
s.consecFail++
|
||
if !s.inOutage {
|
||
s.inOutage = true
|
||
s.outageStart = msg.at
|
||
s.outageOpsLost = 1
|
||
m.appendEvent("err", msg.at,
|
||
fmt.Sprintf("[%s] OUTAGE START — %s: %s", msg.worker, msg.op, trim(msg.err.Error(), 80)))
|
||
} else {
|
||
s.outageOpsLost++
|
||
}
|
||
return
|
||
}
|
||
s.ok++
|
||
if s.inOutage {
|
||
d := msg.at.Sub(s.outageStart)
|
||
s.outages = append(s.outages, outageRec{
|
||
start: s.outageStart, end: msg.at, opsLost: s.outageOpsLost,
|
||
})
|
||
m.appendEvent("info", msg.at,
|
||
fmt.Sprintf("[%s] OUTAGE END — duration=%s ops_lost=%d",
|
||
msg.worker, d.Round(time.Millisecond), s.outageOpsLost))
|
||
s.inOutage = false
|
||
s.consecFail = 0
|
||
s.outageOpsLost = 0
|
||
}
|
||
if msg.dur > 0 {
|
||
s.recentLat = appendBoundedDur(s.recentLat, msg.dur, maxRecentLat)
|
||
m.globalLat = appendBoundedDur(m.globalLat, msg.dur, maxGlobalLat)
|
||
}
|
||
}
|
||
|
||
func (m *chaosModel) applySubRecv(msg subRecvMsg) {
|
||
s, ok := m.stats[msg.worker]
|
||
if !ok {
|
||
return
|
||
}
|
||
s.msgs++
|
||
if s.lastSeq != 0 && msg.seq > s.lastSeq+1 {
|
||
gap := msg.seq - s.lastSeq - 1
|
||
s.gaps += gap
|
||
m.appendEvent("warn", msg.at,
|
||
fmt.Sprintf("[%s] PUBSUB GAP last=%d now=%d missed=%d",
|
||
msg.worker, s.lastSeq, msg.seq, gap))
|
||
}
|
||
if !s.lastSeqAt.IsZero() {
|
||
if d := msg.at.Sub(s.lastSeqAt); d > s.longestGap {
|
||
s.longestGap = d
|
||
}
|
||
}
|
||
s.lastSeq = msg.seq
|
||
s.lastSeqAt = msg.at
|
||
}
|
||
|
||
func (m *chaosModel) applyLedgerCheck(msg ledgerCheckMsg) {
|
||
s, ok := m.stats[msg.worker]
|
||
if !ok {
|
||
return
|
||
}
|
||
switch msg.kind {
|
||
case "ok":
|
||
s.verified++
|
||
case "missing":
|
||
s.missing++
|
||
m.appendEvent("err", time.Now(),
|
||
fmt.Sprintf("[%s] LEDGER MISSING seq=%d %s", msg.worker, msg.seq, msg.detail))
|
||
case "mismatch":
|
||
s.mismatch++
|
||
m.appendEvent("err", time.Now(),
|
||
fmt.Sprintf("[%s] LEDGER MISMATCH seq=%d %s", msg.worker, msg.seq, msg.detail))
|
||
}
|
||
}
|
||
|
||
func (m *chaosModel) appendEvent(level string, at time.Time, text string) {
|
||
if len(m.events) >= maxEvents {
|
||
copy(m.events, m.events[1:])
|
||
m.events = m.events[:len(m.events)-1]
|
||
}
|
||
m.events = append(m.events, eventLine{at: at, level: level, text: text})
|
||
atBottom := m.eventVP.AtBottom()
|
||
m.eventVP.SetContent(m.renderEventContent())
|
||
if atBottom {
|
||
m.eventVP.GotoBottom()
|
||
}
|
||
}
|
||
|
||
func (m *chaosModel) onTick() {
|
||
if len(m.globalLat) > 0 {
|
||
cp := make([]time.Duration, len(m.globalLat))
|
||
copy(cp, m.globalLat)
|
||
sort.Slice(cp, func(i, j int) bool { return cp[i] < cp[j] })
|
||
idx := 95 * len(cp) / 100
|
||
if idx >= len(cp) {
|
||
idx = len(cp) - 1
|
||
}
|
||
p95ms := float64(cp[idx]) / float64(time.Millisecond)
|
||
m.p95Series = appendBoundedFloat(m.p95Series, p95ms, maxSeriesPoints)
|
||
} else {
|
||
m.p95Series = appendBoundedFloat(m.p95Series, 0, maxSeriesPoints)
|
||
}
|
||
|
||
var totalOps, totalFail int64
|
||
for _, s := range m.stats {
|
||
totalOps += s.ops + s.msgs
|
||
totalFail += s.fail
|
||
}
|
||
okDelta := (totalOps - totalFail) - (m.prevTotalOps - m.prevTotalFail)
|
||
if okDelta < 0 {
|
||
okDelta = 0
|
||
}
|
||
failDelta := totalFail - m.prevTotalFail
|
||
if failDelta < 0 {
|
||
failDelta = 0
|
||
}
|
||
m.prevTotalOps = totalOps
|
||
m.prevTotalFail = totalFail
|
||
m.okPerSecSeries = appendBoundedFloat(m.okPerSecSeries, float64(okDelta), maxSeriesPoints)
|
||
m.failPerSecSeries = appendBoundedFloat(m.failPerSecSeries, float64(failDelta), maxSeriesPoints)
|
||
|
||
for _, s := range m.stats {
|
||
current := s.ops + s.msgs
|
||
delta := current - s.prevOps
|
||
if delta < 0 {
|
||
delta = 0
|
||
}
|
||
s.prevOps = current
|
||
s.opsHist = appendBoundedFloat(s.opsHist, float64(delta), maxSparkSamples)
|
||
}
|
||
}
|
||
|
||
func (m *chaosModel) recomputeLayout() {
|
||
headerH := 1
|
||
tableTitleH := 1
|
||
tableRows := len(m.workers) + 2
|
||
chartTitleH := 1
|
||
chartH := 8
|
||
chartFooterH := 1
|
||
eventTitleH := 1
|
||
footerH := 1
|
||
used := headerH + tableTitleH + tableRows + chartTitleH + chartH + chartFooterH + eventTitleH + footerH + 2
|
||
eventH := m.height - used
|
||
if eventH < 4 {
|
||
eventH = 4
|
||
}
|
||
w := m.width - 2
|
||
if w < 20 {
|
||
w = 20
|
||
}
|
||
m.eventVP.Width = w
|
||
m.eventVP.Height = eventH
|
||
m.eventVP.SetContent(m.renderEventContent())
|
||
}
|
||
|
||
// ---- TUI rendering ------------------------------------------------------
|
||
|
||
var (
|
||
titleStyle = lipgloss.NewStyle().Bold(true).Foreground(lipgloss.Color("212"))
|
||
headStyle = lipgloss.NewStyle().Bold(true).Foreground(lipgloss.Color("39"))
|
||
dimStyle = lipgloss.NewStyle().Foreground(lipgloss.Color("241"))
|
||
okStyle = lipgloss.NewStyle().Foreground(lipgloss.Color("42"))
|
||
warnStyle = lipgloss.NewStyle().Foreground(lipgloss.Color("214"))
|
||
errStyle = lipgloss.NewStyle().Foreground(lipgloss.Color("196"))
|
||
chartStyle = lipgloss.NewStyle().Foreground(lipgloss.Color("141"))
|
||
)
|
||
|
||
func (m *chaosModel) View() string {
|
||
if m.quitting {
|
||
return "stopping…\n"
|
||
}
|
||
if m.width == 0 || m.height == 0 {
|
||
return "loading…\n"
|
||
}
|
||
if m.width < 90 || m.height < 25 {
|
||
return fmt.Sprintf("terminal too small: %dx%d (need 90x25 minimum). press q to quit.\n",
|
||
m.width, m.height)
|
||
}
|
||
|
||
header := m.renderHeader()
|
||
table := m.renderTable()
|
||
charts := m.renderCharts()
|
||
events := m.renderEvents()
|
||
footer := dimStyle.Render(" q quit · ↑/↓ scroll · pgup/pgdn half-page · g/G top/bottom")
|
||
|
||
return strings.Join([]string{header, "", table, "", charts, "", events, footer}, "\n")
|
||
}
|
||
|
||
func (m *chaosModel) renderHeader() string {
|
||
elapsed := time.Since(m.startAt).Round(time.Second)
|
||
var inOutage int
|
||
for _, s := range m.stats {
|
||
if s.inOutage {
|
||
inOutage++
|
||
}
|
||
}
|
||
tag := okStyle.Render("healthy")
|
||
if inOutage > 0 {
|
||
tag = errStyle.Render(fmt.Sprintf("%d worker(s) in outage", inOutage))
|
||
}
|
||
return titleStyle.Render(" valkey-ha chaos") +
|
||
dimStyle.Render(fmt.Sprintf(" · %d workers · elapsed %s · ", len(m.workers), elapsed)) +
|
||
tag
|
||
}
|
||
|
||
func (m *chaosModel) renderTable() string {
|
||
var b strings.Builder
|
||
b.WriteString(headStyle.Render(" workers") + "\n")
|
||
hdr := fmt.Sprintf(" %-14s %-11s %7s %7s %6s %7s %8s %8s %-12s %s",
|
||
"NAME", "ROLE", "OPS", "OK", "FAIL", "UPTIME", "P50", "P95", "OPS/SEC", "STATUS")
|
||
b.WriteString(dimStyle.Render(hdr) + "\n")
|
||
|
||
for _, w := range m.workers {
|
||
s := m.stats[w.Name]
|
||
switch s.role {
|
||
case "subscriber":
|
||
line := fmt.Sprintf(" %-14s %-11s msgs=%-7d gaps=%-3d longest_gap=%-7s %s",
|
||
w.Name, "subscriber", s.msgs, s.gaps,
|
||
shortDur(s.longestGap), m.statusBadge(s))
|
||
b.WriteString(line + "\n")
|
||
default:
|
||
uptime := 100.0
|
||
if s.ops > 0 {
|
||
uptime = float64(s.ok) / float64(s.ops) * 100
|
||
}
|
||
p50 := percentileLat(s.recentLat, 50)
|
||
p95 := percentileLat(s.recentLat, 95)
|
||
spark := chartStyle.Render(sparkline(s.opsHist, 12))
|
||
extra := ""
|
||
if s.role == "reader" {
|
||
extra = dimStyle.Render(fmt.Sprintf(" v=%d miss=%d mis=%d",
|
||
s.verified, s.missing, s.mismatch))
|
||
}
|
||
line := fmt.Sprintf(" %-14s %-11s %7d %7d %6d %6.2f%% %8s %8s %s %s%s",
|
||
w.Name, s.role, s.ops, s.ok, s.fail, uptime,
|
||
shortDur(p50), shortDur(p95),
|
||
spark, m.statusBadge(s), extra)
|
||
b.WriteString(line + "\n")
|
||
}
|
||
}
|
||
return strings.TrimRight(b.String(), "\n")
|
||
}
|
||
|
||
func (m *chaosModel) statusBadge(s *uiStats) string {
|
||
if s.inOutage {
|
||
return errStyle.Render("OUTAGE")
|
||
}
|
||
if s.role == "subscriber" && s.msgs == 0 && time.Since(m.startAt) > 5*time.Second {
|
||
return warnStyle.Render("NO MSGS")
|
||
}
|
||
return okStyle.Render("OK")
|
||
}
|
||
|
||
func (m *chaosModel) renderCharts() string {
|
||
chartH := 8
|
||
chartW := (m.width - 4) / 3 // 3 charts with 2× " " separators between
|
||
if chartW < 20 {
|
||
chartW = 20
|
||
}
|
||
|
||
p95Max := maxFloat(m.p95Series)
|
||
if p95Max < 1 {
|
||
p95Max = 1
|
||
}
|
||
p95Lines := renderBarChart(m.p95Series, chartW, chartH, p95Max)
|
||
p95Title := headStyle.Render(" p95 latency") +
|
||
dimStyle.Render(fmt.Sprintf(" peak %.1fms · last %s",
|
||
p95Max, lastFloatLabel(m.p95Series, "ms")))
|
||
p95Box := p95Title + "\n" + chartStyle.Render(strings.Join(p95Lines, "\n"))
|
||
|
||
tpMax := maxFloat(m.okPerSecSeries)
|
||
if tpMax < 5 {
|
||
tpMax = 5
|
||
}
|
||
tpLines := renderBarChart(m.okPerSecSeries, chartW, chartH, tpMax)
|
||
tpTitle := headStyle.Render(" ok ops/sec") +
|
||
dimStyle.Render(fmt.Sprintf(" peak %.0f · last %s",
|
||
tpMax, lastFloatLabel(m.okPerSecSeries, "")))
|
||
tpBox := tpTitle + "\n" + chartStyle.Render(strings.Join(tpLines, "\n"))
|
||
|
||
failMax := maxFloat(m.failPerSecSeries)
|
||
if failMax < 1 {
|
||
failMax = 1
|
||
}
|
||
failLines := renderBarChart(m.failPerSecSeries, chartW, chartH, failMax)
|
||
failTitle := headStyle.Render(" fail ops/sec") +
|
||
dimStyle.Render(fmt.Sprintf(" peak %.0f · last %s",
|
||
failMax, lastFloatLabel(m.failPerSecSeries, "")))
|
||
failBox := failTitle + "\n" + errStyle.Render(strings.Join(failLines, "\n"))
|
||
|
||
return lipgloss.JoinHorizontal(lipgloss.Top, p95Box, " ", tpBox, " ", failBox)
|
||
}
|
||
|
||
func (m *chaosModel) renderEvents() string {
|
||
title := headStyle.Render(" events") +
|
||
dimStyle.Render(fmt.Sprintf(" (%d, scroll with ↑/↓)", len(m.events)))
|
||
return title + "\n" + m.eventVP.View()
|
||
}
|
||
|
||
func (m *chaosModel) renderEventContent() string {
|
||
if len(m.events) == 0 {
|
||
return dimStyle.Render("(no events yet)")
|
||
}
|
||
var b strings.Builder
|
||
for _, e := range m.events {
|
||
ts := dimStyle.Render(e.at.Format("15:04:05.000"))
|
||
var styled string
|
||
switch e.level {
|
||
case "err":
|
||
styled = errStyle.Render(e.text)
|
||
case "warn":
|
||
styled = warnStyle.Render(e.text)
|
||
default:
|
||
styled = e.text
|
||
}
|
||
b.WriteString(ts + " " + styled + "\n")
|
||
}
|
||
return strings.TrimRight(b.String(), "\n")
|
||
}
|
||
|
||
// ---- chart primitives ---------------------------------------------------
|
||
|
||
var sparkChars = []rune(" ▁▂▃▄▅▆▇█")
|
||
|
||
func sparkline(values []float64, width int) string {
|
||
if width <= 0 {
|
||
return ""
|
||
}
|
||
if len(values) == 0 {
|
||
return strings.Repeat(" ", width)
|
||
}
|
||
maxV := 0.0
|
||
for _, v := range values {
|
||
if v > maxV {
|
||
maxV = v
|
||
}
|
||
}
|
||
var b strings.Builder
|
||
n := len(values)
|
||
for col := 0; col < width; col++ {
|
||
if maxV == 0 {
|
||
b.WriteRune(' ')
|
||
continue
|
||
}
|
||
idx := col * n / width
|
||
if idx >= n {
|
||
idx = n - 1
|
||
}
|
||
v := values[idx]
|
||
if v < 0 {
|
||
v = 0
|
||
}
|
||
if v > maxV {
|
||
v = maxV
|
||
}
|
||
i := int(v / maxV * float64(len(sparkChars)-1))
|
||
if i < 0 {
|
||
i = 0
|
||
}
|
||
if i >= len(sparkChars) {
|
||
i = len(sparkChars) - 1
|
||
}
|
||
b.WriteRune(sparkChars[i])
|
||
}
|
||
return b.String()
|
||
}
|
||
|
||
func renderBarChart(values []float64, width, height int, maxV float64) []string {
|
||
rows := make([][]rune, height)
|
||
for i := range rows {
|
||
rows[i] = []rune(strings.Repeat(" ", width))
|
||
}
|
||
flush := func() []string {
|
||
out := make([]string, height)
|
||
for i := range rows {
|
||
out[i] = string(rows[i])
|
||
}
|
||
return out
|
||
}
|
||
if width <= 0 || height <= 0 || maxV <= 0 || len(values) == 0 {
|
||
return flush()
|
||
}
|
||
n := len(values)
|
||
subSteps := []rune(" ▁▂▃▄▅▆▇")
|
||
for col := 0; col < width; col++ {
|
||
idx := col * n / width
|
||
if idx >= n {
|
||
idx = n - 1
|
||
}
|
||
v := values[idx]
|
||
if v < 0 {
|
||
v = 0
|
||
}
|
||
if v > maxV {
|
||
v = maxV
|
||
}
|
||
cells := int(v / maxV * float64(height) * 8)
|
||
full := cells / 8
|
||
rem := cells % 8
|
||
for r := 0; r < full && r < height; r++ {
|
||
rows[height-1-r][col] = '█'
|
||
}
|
||
if full < height && rem > 0 {
|
||
rows[height-1-full][col] = subSteps[rem]
|
||
}
|
||
}
|
||
return flush()
|
||
}
|
||
|
||
func lastFloatLabel(s []float64, unit string) string {
|
||
if len(s) == 0 {
|
||
return "-"
|
||
}
|
||
v := s[len(s)-1]
|
||
if unit == "ms" {
|
||
return fmt.Sprintf("%.1f%s", v, unit)
|
||
}
|
||
return fmt.Sprintf("%.0f%s", v, unit)
|
||
}
|
||
|
||
// ---- runChaos (local TUI) ----------------------------------------------
|
||
|
||
func runChaos(ws []WorkerConfig) int {
|
||
for _, w := range ws {
|
||
if err := validateWorker(w); err != nil {
|
||
fmt.Fprintf(os.Stderr, "worker %q: %v\n", w.Name, err)
|
||
return 2
|
||
}
|
||
}
|
||
|
||
logFile, ferr := os.Create("chaos-debug.log")
|
||
if ferr == nil {
|
||
defer logFile.Close()
|
||
log.SetOutput(logFile)
|
||
} else {
|
||
log.SetOutput(io.Discard)
|
||
}
|
||
defer log.SetOutput(os.Stderr)
|
||
|
||
workerCtx, cancel := context.WithCancel(context.Background())
|
||
|
||
cc := &chaosCtx{
|
||
ledger: newLedger(),
|
||
startAt: time.Now(),
|
||
}
|
||
|
||
workers := buildWorkers(ws, cc)
|
||
defer func() {
|
||
for _, w := range workers {
|
||
_ = w.rdb.Close()
|
||
}
|
||
}()
|
||
|
||
model := newChaosModel(ws, cc.startAt)
|
||
p := tea.NewProgram(model, tea.WithAltScreen())
|
||
cc.sink = &teaSink{p: p}
|
||
|
||
sigCh := make(chan os.Signal, 1)
|
||
signal.Notify(sigCh, os.Interrupt, syscall.SIGTERM)
|
||
go func() {
|
||
select {
|
||
case <-sigCh:
|
||
p.Quit()
|
||
case <-workerCtx.Done():
|
||
}
|
||
}()
|
||
|
||
var wg sync.WaitGroup
|
||
for _, w := range workers {
|
||
wg.Add(1)
|
||
go func(w *worker) {
|
||
defer wg.Done()
|
||
w.run(workerCtx)
|
||
}(w)
|
||
}
|
||
|
||
if _, err := p.Run(); err != nil {
|
||
fmt.Fprintf(os.Stderr, "tui error: %v\n", err)
|
||
}
|
||
|
||
cancel()
|
||
waitWorkers(&wg, 5*time.Second)
|
||
|
||
fmt.Println("\nrunning final durability sweep…")
|
||
fd := finalSweep(workers, cc)
|
||
printFinalReport(model, time.Since(cc.startAt), fd)
|
||
|
||
if fd.missing > 0 || fd.mismatch > 0 {
|
||
return 1
|
||
}
|
||
return 0
|
||
}
|
||
|
||
// ---- 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
|
||
}
|
||
|
||
func printFinalReport(m *chaosModel, dur time.Duration, fd finalDur) {
|
||
fmt.Println("=== valkey-ha chaos report ===")
|
||
fmt.Printf("duration: %s\n\n", dur.Round(time.Millisecond))
|
||
|
||
fmt.Println("per-worker:")
|
||
for _, w := range m.workers {
|
||
s := m.stats[w.Name]
|
||
switch s.role {
|
||
case "subscriber":
|
||
fmt.Printf(" %-14s [sub] msgs=%d gaps=%d longest_gap=%s\n",
|
||
w.Name, s.msgs, s.gaps, s.longestGap.Round(time.Millisecond))
|
||
default:
|
||
uptime := 100.0
|
||
if s.ops > 0 {
|
||
uptime = float64(s.ok) / float64(s.ops) * 100
|
||
}
|
||
var longest time.Duration
|
||
for _, o := range s.outages {
|
||
if d := o.end.Sub(o.start); d > longest {
|
||
longest = d
|
||
}
|
||
}
|
||
extra := ""
|
||
if s.role == "reader" {
|
||
extra = fmt.Sprintf(" verified=%d missing=%d mismatch=%d",
|
||
s.verified, s.missing, s.mismatch)
|
||
}
|
||
fmt.Printf(" %-14s [%-10s] ops=%d ok=%d fail=%d uptime=%.2f%% p50=%s p95=%s outages=%d longest=%s%s\n",
|
||
w.Name, s.role, s.ops, s.ok, s.fail, uptime,
|
||
shortDur(percentileLat(s.recentLat, 50)),
|
||
shortDur(percentileLat(s.recentLat, 95)),
|
||
len(s.outages), longest.Round(time.Millisecond),
|
||
extra)
|
||
}
|
||
}
|
||
|
||
fmt.Println("\ndurability:")
|
||
fmt.Printf(" ledger: acked=%d verified=%d missing=%d mismatch=%d\n",
|
||
fd.acked, fd.verified, fd.missing, fd.mismatch)
|
||
if len(fd.missingSeqs) > 0 {
|
||
n := 20
|
||
if len(fd.missingSeqs) < n {
|
||
n = len(fd.missingSeqs)
|
||
}
|
||
fmt.Printf(" missing seqs (first %d): %v\n", n, fd.missingSeqs[:n])
|
||
}
|
||
delta := fd.serverCounter - fd.localAcked
|
||
fmt.Printf(" counter: local_acked=%d server_value=%d delta=%d\n",
|
||
fd.localAcked, fd.serverCounter, delta)
|
||
|
||
type tagged struct {
|
||
w string
|
||
rec outageRec
|
||
}
|
||
var all []tagged
|
||
for _, w := range m.workers {
|
||
for _, o := range m.stats[w.Name].outages {
|
||
all = append(all, tagged{w: w.Name, rec: o})
|
||
}
|
||
}
|
||
sort.Slice(all, func(i, j int) bool { return all[i].rec.start.Before(all[j].rec.start) })
|
||
if len(all) > 0 {
|
||
fmt.Println("\noutages (chronological):")
|
||
for _, t := range all {
|
||
fmt.Printf(" %s -> %s (%s) worker=%s ops_lost=%d\n",
|
||
t.rec.start.Format("15:04:05.000"),
|
||
t.rec.end.Format("15:04:05.000"),
|
||
t.rec.end.Sub(t.rec.start).Round(time.Millisecond),
|
||
t.w, t.rec.opsLost)
|
||
}
|
||
}
|
||
|
||
fmt.Println("\nnotes:")
|
||
fmt.Println(" - Valkey 7.2 + Sentinel does NOT replicate pub/sub from master to replicas.")
|
||
fmt.Println(" A subscriber on the read-only VIP receiving 0 messages is expected.")
|
||
fmt.Println(" - debug log: ./chaos-debug.log")
|
||
}
|
||
|
||
// ---- 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 shortDur(d time.Duration) string {
|
||
if d <= 0 {
|
||
return "-"
|
||
}
|
||
if d < time.Microsecond {
|
||
return "0"
|
||
}
|
||
if d < time.Millisecond {
|
||
return fmt.Sprintf("%dµs", d.Microseconds())
|
||
}
|
||
if d < time.Second {
|
||
return fmt.Sprintf("%.1fms", float64(d)/float64(time.Millisecond))
|
||
}
|
||
return d.Round(time.Millisecond).String()
|
||
}
|
||
|
||
func trim(s string, n int) string {
|
||
if len(s) <= n {
|
||
return s
|
||
}
|
||
return s[:n] + "…"
|
||
}
|