Files
2026-05-17 21:01:07 +02:00

1521 lines
37 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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] + "…"
}