From 6a43ba022c43a5e11e6464e22cdd3cba8e588af4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Mat=C4=9Bj=20Pavl=C3=AD=C4=8Dek?= Date: Thu, 11 Jun 2026 10:52:17 +0200 Subject: [PATCH] valkey chaos --- valkey/chaos.go | 708 +--------------------------- valkey/config-sea-ha.yaml | 89 ++++ valkey/go.mod | 23 - valkey/go.sum | 50 -- valkey/main.go | 381 ++++++--------- valkey/zerops-chaos-sea-import.yaml | 17 + valkey/zerops.yaml | 24 + 7 files changed, 287 insertions(+), 1005 deletions(-) create mode 100644 valkey/config-sea-ha.yaml create mode 100644 valkey/zerops-chaos-sea-import.yaml diff --git a/valkey/chaos.go b/valkey/chaos.go index 504fed0..d552e56 100644 --- a/valkey/chaos.go +++ b/valkey/chaos.go @@ -19,9 +19,6 @@ import ( "syscall" "time" - "github.com/charmbracelet/bubbles/viewport" - tea "github.com/charmbracelet/bubbletea" - "github.com/charmbracelet/lipgloss" "github.com/redis/go-redis/v9" ) @@ -129,13 +126,9 @@ type outageRec struct { // ---- 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. +// 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) @@ -145,16 +138,6 @@ type EventSink interface { 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 { @@ -564,7 +547,7 @@ func (w *worker) runSubscriber(ctx context.Context) { } } -// ---- TUI message types -------------------------------------------------- +// ---- event message types (worker → sink) ------------------------------- type opResultMsg struct { worker string @@ -593,9 +576,7 @@ type eventMsg struct { text string } -type tickMsg time.Time - -// ---- TUI model ---------------------------------------------------------- +// ---- shared stats types (mirrored server-side by the collector) --------- type uiStats struct { role string @@ -629,591 +610,6 @@ type eventLine struct { 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 { @@ -1369,84 +765,6 @@ func finalSweep(workers []*worker, cc *chaosCtx) finalDur { 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 { @@ -1496,22 +814,6 @@ func maxDur(a, b time.Duration) time.Duration { 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 diff --git a/valkey/config-sea-ha.yaml b/valkey/config-sea-ha.yaml new file mode 100644 index 0000000..070e09c --- /dev/null +++ b/valkey/config-sea-ha.yaml @@ -0,0 +1,89 @@ +# Chaos workload for sea1, run against the HA valkeys. +# +# Ports (per Zerops valkey:ha): 6379 = master plaintext, 7000 = replica +# plaintext. NOTE: the TLS ports (6380/7001) require mutual TLS (a client +# certificate) on these migration valkeys — the tester only does server-side +# trust (--tls-ca), so we use the plaintext ports with password auth, the same +# path the seeding used. +# +# Auth: per-service password via ${valkeyhaN_password} (wired in zerops.yaml +# setup `chaosworker-sea`). +# +# Targeting rationale: +# valkeyha3 (2 GiB) — full mix incl. write-heavy roles. The durability +# ledger (vk:chaos:ledger:*) is written with NO TTL and +# grows over time, so writes go only where there's RAM +# headroom. +# valkeyha1/2 (0.5 GiB) — pub/sub only: exercises the master ingress and +# pub/sub path with zero key growth (and no cross-host +# read-after-write misses, since readers verify keys the +# same container wrote). +workers: + # ---- valkeyha3: full mix -------------------------------------------------- + - name: rw-ha3 + role: readwriter + host: valkeyha3.zerops + port: 6379 + password: ${valkeyha3_password} + interval: 10ms + - name: wr-ha3 + role: writer + host: valkeyha3.zerops + port: 6379 + password: ${valkeyha3_password} + interval: 5ms + - name: rom-ha3 + role: reader + host: valkeyha3.zerops + port: 6379 + password: ${valkeyha3_password} + interval: 5ms + - name: ror-ha3 + role: reader + host: valkeyha3.zerops + port: 7000 + password: ${valkeyha3_password} + interval: 3ms + - name: pub-ha3 + role: publisher + host: valkeyha3.zerops + port: 6379 + password: ${valkeyha3_password} + channel: "vk:chaos:ha3" + interval: 5ms + - name: sub-ha3 + role: subscriber + host: valkeyha3.zerops + port: 6379 + password: ${valkeyha3_password} + channel: "vk:chaos:ha3" + + # ---- valkeyha1: pub/sub only --------------------------------------------- + - name: pub-ha1 + role: publisher + host: valkeyha1.zerops + port: 6379 + password: ${valkeyha1_password} + channel: "vk:chaos:ha1" + interval: 5ms + - name: sub-ha1 + role: subscriber + host: valkeyha1.zerops + port: 6379 + password: ${valkeyha1_password} + channel: "vk:chaos:ha1" + + # ---- valkeyha2: pub/sub only --------------------------------------------- + - name: pub-ha2 + role: publisher + host: valkeyha2.zerops + port: 6379 + password: ${valkeyha2_password} + channel: "vk:chaos:ha2" + interval: 5ms + - name: sub-ha2 + role: subscriber + host: valkeyha2.zerops + port: 6379 + password: ${valkeyha2_password} + channel: "vk:chaos:ha2" diff --git a/valkey/go.mod b/valkey/go.mod index c71a613..99886d8 100644 --- a/valkey/go.mod +++ b/valkey/go.mod @@ -3,34 +3,11 @@ module valkey-test go 1.24.2 require ( - github.com/charmbracelet/bubbles v1.0.0 - github.com/charmbracelet/bubbletea v1.3.10 - github.com/charmbracelet/lipgloss v1.1.0 github.com/redis/go-redis/v9 v9.7.0 gopkg.in/yaml.v3 v3.0.1 ) require ( - github.com/aymanbagabas/go-osc52/v2 v2.0.1 // indirect github.com/cespare/xxhash/v2 v2.2.0 // indirect - github.com/charmbracelet/colorprofile v0.4.1 // indirect - github.com/charmbracelet/x/ansi v0.11.6 // indirect - github.com/charmbracelet/x/cellbuf v0.0.15 // indirect - github.com/charmbracelet/x/term v0.2.2 // indirect - github.com/clipperhouse/displaywidth v0.9.0 // indirect - github.com/clipperhouse/stringish v0.1.1 // indirect - github.com/clipperhouse/uax29/v2 v2.5.0 // indirect github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect - github.com/erikgeiser/coninput v0.0.0-20211004153227-1c3628e74d0f // indirect - github.com/lucasb-eyer/go-colorful v1.3.0 // indirect - github.com/mattn/go-isatty v0.0.20 // indirect - github.com/mattn/go-localereader v0.0.1 // indirect - github.com/mattn/go-runewidth v0.0.19 // indirect - github.com/muesli/ansi v0.0.0-20230316100256-276c6243b2f6 // indirect - github.com/muesli/cancelreader v0.2.2 // indirect - github.com/muesli/termenv v0.16.0 // indirect - github.com/rivo/uniseg v0.4.7 // indirect - github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e // indirect - golang.org/x/sys v0.38.0 // indirect - golang.org/x/text v0.3.8 // indirect ) diff --git a/valkey/go.sum b/valkey/go.sum index 93efbcc..caca62d 100644 --- a/valkey/go.sum +++ b/valkey/go.sum @@ -1,63 +1,13 @@ -github.com/aymanbagabas/go-osc52/v2 v2.0.1 h1:HwpRHbFMcZLEVr42D4p7XBqjyuxQH5SMiErDT4WkJ2k= -github.com/aymanbagabas/go-osc52/v2 v2.0.1/go.mod h1:uYgXzlJ7ZpABp8OJ+exZzJJhRNQ2ASbcXHWsFqH8hp8= github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= github.com/cespare/xxhash/v2 v2.2.0 h1:DC2CZ1Ep5Y4k3ZQ899DldepgrayRUGE6BBZ/cd9Cj44= github.com/cespare/xxhash/v2 v2.2.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= -github.com/charmbracelet/bubbles v1.0.0 h1:12J8/ak/uCZEMQ6KU7pcfwceyjLlWsDLAxB5fXonfvc= -github.com/charmbracelet/bubbles v1.0.0/go.mod h1:9d/Zd5GdnauMI5ivUIVisuEm3ave1XwXtD1ckyV6r3E= -github.com/charmbracelet/bubbletea v1.3.10 h1:otUDHWMMzQSB0Pkc87rm691KZ3SWa4KUlvF9nRvCICw= -github.com/charmbracelet/bubbletea v1.3.10/go.mod h1:ORQfo0fk8U+po9VaNvnV95UPWA1BitP1E0N6xJPlHr4= -github.com/charmbracelet/colorprofile v0.4.1 h1:a1lO03qTrSIRaK8c3JRxJDZOvhvIeSco3ej+ngLk1kk= -github.com/charmbracelet/colorprofile v0.4.1/go.mod h1:U1d9Dljmdf9DLegaJ0nGZNJvoXAhayhmidOdcBwAvKk= -github.com/charmbracelet/lipgloss v1.1.0 h1:vYXsiLHVkK7fp74RkV7b2kq9+zDLoEU4MZoFqR/noCY= -github.com/charmbracelet/lipgloss v1.1.0/go.mod h1:/6Q8FR2o+kj8rz4Dq0zQc3vYf7X+B0binUUBwA0aL30= -github.com/charmbracelet/x/ansi v0.11.6 h1:GhV21SiDz/45W9AnV2R61xZMRri5NlLnl6CVF7ihZW8= -github.com/charmbracelet/x/ansi v0.11.6/go.mod h1:2JNYLgQUsyqaiLovhU2Rv/pb8r6ydXKS3NIttu3VGZQ= -github.com/charmbracelet/x/cellbuf v0.0.15 h1:ur3pZy0o6z/R7EylET877CBxaiE1Sp1GMxoFPAIztPI= -github.com/charmbracelet/x/cellbuf v0.0.15/go.mod h1:J1YVbR7MUuEGIFPCaaZ96KDl5NoS0DAWkskup+mOY+Q= -github.com/charmbracelet/x/term v0.2.2 h1:xVRT/S2ZcKdhhOuSP4t5cLi5o+JxklsoEObBSgfgZRk= -github.com/charmbracelet/x/term v0.2.2/go.mod h1:kF8CY5RddLWrsgVwpw4kAa6TESp6EB5y3uxGLeCqzAI= -github.com/clipperhouse/displaywidth v0.9.0 h1:Qb4KOhYwRiN3viMv1v/3cTBlz3AcAZX3+y9OLhMtAtA= -github.com/clipperhouse/displaywidth v0.9.0/go.mod h1:aCAAqTlh4GIVkhQnJpbL0T/WfcrJXHcj8C0yjYcjOZA= -github.com/clipperhouse/stringish v0.1.1 h1:+NSqMOr3GR6k1FdRhhnXrLfztGzuG+VuFDfatpWHKCs= -github.com/clipperhouse/stringish v0.1.1/go.mod h1:v/WhFtE1q0ovMta2+m+UbpZ+2/HEXNWYXQgCt4hdOzA= -github.com/clipperhouse/uax29/v2 v2.5.0 h1:x7T0T4eTHDONxFJsL94uKNKPHrclyFI0lm7+w94cO8U= -github.com/clipperhouse/uax29/v2 v2.5.0/go.mod h1:Wn1g7MK6OoeDT0vL+Q0SQLDz/KpfsVRgg6W7ihQeh4g= github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f h1:lO4WD4F/rVNCu3HqELle0jiPLLBs70cWOduZpkS1E78= github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f/go.mod h1:cuUVRXasLTGF7a8hSLbxyZXjz+1KgoB3wDUb6vlszIc= -github.com/erikgeiser/coninput v0.0.0-20211004153227-1c3628e74d0f h1:Y/CXytFA4m6baUTXGLOoWe4PQhGxaX0KpnayAqC48p4= -github.com/erikgeiser/coninput v0.0.0-20211004153227-1c3628e74d0f/go.mod h1:vw97MGsxSvLiUE2X8qFplwetxpGLQrlU1Q9AUEIzCaM= -github.com/lucasb-eyer/go-colorful v1.3.0 h1:2/yBRLdWBZKrf7gB40FoiKfAWYQ0lqNcbuQwVHXptag= -github.com/lucasb-eyer/go-colorful v1.3.0/go.mod h1:R4dSotOR9KMtayYi1e77YzuveK+i7ruzyGqttikkLy0= -github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= -github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= -github.com/mattn/go-localereader v0.0.1 h1:ygSAOl7ZXTx4RdPYinUpg6W99U8jWvWi9Ye2JC/oIi4= -github.com/mattn/go-localereader v0.0.1/go.mod h1:8fBrzywKY7BI3czFoHkuzRoWE9C+EiG4R1k4Cjx5p88= -github.com/mattn/go-runewidth v0.0.19 h1:v++JhqYnZuu5jSKrk9RbgF5v4CGUjqRfBm05byFGLdw= -github.com/mattn/go-runewidth v0.0.19/go.mod h1:XBkDxAl56ILZc9knddidhrOlY5R/pDhgLpndooCuJAs= -github.com/muesli/ansi v0.0.0-20230316100256-276c6243b2f6 h1:ZK8zHtRHOkbHy6Mmr5D264iyp3TiX5OmNcI5cIARiQI= -github.com/muesli/ansi v0.0.0-20230316100256-276c6243b2f6/go.mod h1:CJlz5H+gyd6CUWT45Oy4q24RdLyn7Md9Vj2/ldJBSIo= -github.com/muesli/cancelreader v0.2.2 h1:3I4Kt4BQjOR54NavqnDogx/MIoWBFa0StPA8ELUXHmA= -github.com/muesli/cancelreader v0.2.2/go.mod h1:3XuTXfFS2VjM+HTLZY9Ak0l6eUKfijIfMUZ4EgX0QYo= -github.com/muesli/termenv v0.16.0 h1:S5AlUN9dENB57rsbnkPyfdGuWIlkmzJjbFf0Tf5FWUc= -github.com/muesli/termenv v0.16.0/go.mod h1:ZRfOIKPFDYQoDFF4Olj7/QJbW60Ol/kL1pU3VfY/Cnk= github.com/redis/go-redis/v9 v9.7.0 h1:HhLSs+B6O021gwzl+locl0zEDnyNkxMtf/Z3NNBMa9E= github.com/redis/go-redis/v9 v9.7.0/go.mod h1:f6zhXITC7JUJIlPEiBOTXxJgPLdZcA93GewI7inzyWw= -github.com/rivo/uniseg v0.4.7 h1:WUdvkW8uEhrYfLC4ZzdpI2ztxP1I582+49Oc5Mq64VQ= -github.com/rivo/uniseg v0.4.7/go.mod h1:FN3SvrM+Zdj16jyLfmOkMNblXMcoc8DfTHruCPUcx88= -github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e h1:JVG44RsyaB9T2KIHavMF/ppJZNG9ZpyihvCd0w101no= -github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e/go.mod h1:RbqR21r5mrJuqunuUZ/Dhy/avygyECGrLceyNeo4LiM= -golang.org/x/exp v0.0.0-20231006140011-7918f672742d h1:jtJma62tbqLibJ5sFQz8bKtEM8rJBtfilJ2qTU199MI= -golang.org/x/exp v0.0.0-20231006140011-7918f672742d/go.mod h1:ldy0pHrwJyGW56pPQzzkH36rKxoZW1tw7ZJpeKx+hdo= -golang.org/x/sys v0.0.0-20210809222454-d867a43fc93e/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.38.0 h1:3yZWxaJjBmCWXqhN1qh02AkOnCQ1poK6oF+a7xWL6Gc= -golang.org/x/sys v0.38.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= -golang.org/x/text v0.3.8 h1:nAL+RVCQ9uMn3vJZbV+MRnydTJFPf8qqY42YiA6MrqY= -golang.org/x/text v0.3.8/go.mod h1:E6s5w1FMmriuDzIBO73fBruAKo1PCIq6d2Q6DHfQ8WQ= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= diff --git a/valkey/main.go b/valkey/main.go index 40e77f3..086af85 100644 --- a/valkey/main.go +++ b/valkey/main.go @@ -139,17 +139,6 @@ func applyFlagOverrides(cfg *Config, host string, port int, password string, tls return nil } -func (c *Config) primary() (*redis.Options, error) { - if c.Host != "" && c.Port != 0 { - return optionsFromConn(c.Host, c.Port, c.Password, c.DB, c.TLS), nil - } - if len(c.Workers) > 0 { - w := c.Workers[0] - return optionsFromConn(w.Host, w.Port, w.Password, w.DB, w.TLS), nil - } - return nil, fmt.Errorf("config has neither top-level host nor workers") -} - func optionsFromConn(host string, port int, password string, db int, useTLS bool) *redis.Options { opts := &redis.Options{ Addr: fmt.Sprintf("%s:%d", host, port), @@ -219,20 +208,19 @@ func envInt(name string, def int) int { } func main() { - configPath := flag.String("config", "config.yaml", "path to YAML config file") - continuous := flag.Bool("continuous", false, "single-worker rolling smoke loop") - chaosLocal := flag.Bool("chaos", false, "multi-worker chaos TUI (local, in-process)") - workerMode := flag.Bool("worker", false, "distributed worker — runs the chaos workload and ships events to --collector-url") - collectorMode := flag.Bool("collector", false, "distributed collector — HTTP dashboard that ingests events from workers") + configPath := flag.String("config", "config.yaml", "path to YAML config file (chaos --worker/--collector)") + workerMode := flag.Bool("worker", false, "chaos worker — runs the config.yaml workload and ships events to --collector-url") + collectorMode := flag.Bool("collector", false, "chaos collector — HTTP dashboard that ingests events from workers") collectorURL := flag.String("collector-url", os.Getenv("COLLECTOR_URL"), "collector base URL (worker mode); env: COLLECTOR_URL") port := flag.Int("collector-port", envInt("PORT", 8080), "HTTP port (collector mode); env: PORT") - interval := flag.Duration("interval", time.Second, "tick interval for --continuous") - hostOverride := flag.String("host", "", "override host for all workers (local testing)") + hostOverride := flag.String("host", "", "override host for the cli tool (--seed/--verify/etc and the default test suite)") portOverride := flag.Int("port", 0, "override port for all workers (local testing)") passwordOverride := flag.String("password", "", "override password for all workers (local testing)") tlsEnable := flag.Bool("tls", false, "enable TLS for all workers using system root CAs (local testing)") tlsCAPath := flag.String("tls-ca", "", "path to TLS CA cert PEM; enables TLS for all workers (local testing)") - seedConn := flag.String("conn", "", "Valkey connection string for --seed/--verify/--flush/--inspect (e.g. redis://default:pw@valkey1:6379 or rediss://...:6380 for TLS); overrides --host/--port/--password/--tls") + seedConn := flag.String("conn", "", "Valkey connection string for the cli tool (default test suite, --watch, --seed/--verify/--flush/--inspect) — e.g. redis://default:pw@valkey1:6379 or rediss://...:6380 for TLS; overrides --host/--port/--password/--tls") + watchMode := flag.Bool("watch", false, "PING the target on --interval until Ctrl+C, printing each outage's start/recovery/duration and a summary on exit") + watchInterval := flag.Duration("interval", 500*time.Millisecond, "probe interval for --watch") seedMode := flag.Bool("seed", false, "bulk-load one Valkey with deterministic data (local; use --host valkey.zerops over zcli VPN)") verifyMode := flag.Bool("verify", false, "re-read seed:* keys and check they're 1:1 with the seed (same --seed-* flags)") flushMode := flag.Bool("flush", false, "FLUSHDB the selected --db (clears the whole database)") @@ -246,54 +234,19 @@ func main() { seedTimeout := flag.Duration("timeout", 5*time.Second, "dial/read/write timeout in --seed/--verify/--flush/--inspect mode (raise for slow VPN links)") flag.Parse() + // ── Regime 1: chaos ────────────────────────────────────────────────── + // Driven solely by config.yaml. loadConfig applies the env templating that + // lets one config file drive the distributed fleet: ${var} expansion pulls + // secrets from the container env, and $ZEROPS_Number gives each replica a + // distinct DB + pub/sub channel suffix. No connection-flag overrides here. switch { case *collectorMode: os.Exit(runCollector(*port)) - case *seedMode, *verifyMode, *flushMode, *inspectMode: - opts, err := buildConnOptions(*seedConn, *hostOverride, *portOverride, *passwordOverride, *tlsEnable, *tlsCAPath) - if err != nil { - log.Fatalf("connection error: %v", err) - } - // --db overrides whatever the conn string / flags resolved to, so the - // same target can be pointed at any logical database without editing - // the conn URL. -1 means "leave as-is". - if *dbOverride >= 0 { - opts.DB = *dbOverride - } - // Generous timeouts: the zcli VPN adds latency and can stall a dial - // well past go-redis's 5s default mid-run. PoolTimeout > DialTimeout - // so a slow dial doesn't surface as pool exhaustion first. - opts.DialTimeout = *seedTimeout - opts.ReadTimeout = *seedTimeout - opts.WriteTimeout = *seedTimeout - opts.PoolTimeout = *seedTimeout + 5*time.Second - switch { - case *inspectMode: - os.Exit(runInspect(opts)) - case *flushMode: - os.Exit(runFlush(opts)) - case *verifyMode: - os.Exit(runVerify(opts, *seedSeed, *seedMB*1024*1024, *seedValueBytes, *seedBatch)) - default: - os.Exit(runSeed(opts, *seedSeed, *seedMB*1024*1024, *seedValueBytes, *seedBatch, *seedThrottleMB*1024*1024)) - } - } - - cfg, err := loadConfig(*configPath) - if err != nil { - log.Fatalf("config error: %v", err) - } - if err := applyFlagOverrides(cfg, *hostOverride, *portOverride, *passwordOverride, *tlsEnable, *tlsCAPath); err != nil { - log.Fatalf("flag override error: %v", err) - } - - switch { - case *chaosLocal: - if len(cfg.Workers) == 0 { - log.Fatalf("--chaos requires workers: list in %s", *configPath) - } - os.Exit(runChaos(cfg.Workers)) case *workerMode: + cfg, err := loadConfig(*configPath) + if err != nil { + log.Fatalf("config error: %v", err) + } if len(cfg.Workers) == 0 { log.Fatalf("--worker requires workers: list in %s", *configPath) } @@ -303,20 +256,146 @@ func main() { os.Exit(runWorker(cfg.Workers, *collectorURL)) } - primary, err := cfg.primary() + // ── Regime 2: cli tool ─────────────────────────────────────────────── + // One-time commands against a single Valkey, resolved from --conn or the + // split --host/--port/--password/--tls[-ca] flags. Covers the default test + // suite plus --seed/--verify/--flush/--inspect. + opts, err := buildConnOptions(*seedConn, *hostOverride, *portOverride, *passwordOverride, *tlsEnable, *tlsCAPath) if err != nil { - log.Fatalf("config error: %v", err) + log.Fatalf("connection error: %v", err) } - fmt.Printf("redis config: addr=%s db=%d tls=%t password=%s\n", - primary.Addr, primary.DB, primary.TLSConfig != nil, redactPassword(primary.Password)) - rdb := redis.NewClient(primary) + // --db overrides whatever the conn string / flags resolved to, so the same + // target can be pointed at any logical database without editing the conn + // URL. -1 means "leave as-is". + if *dbOverride >= 0 { + opts.DB = *dbOverride + } + // Generous timeouts: the zcli VPN adds latency and can stall a dial well + // past go-redis's 5s default mid-run. PoolTimeout > DialTimeout so a slow + // dial doesn't surface as pool exhaustion first. + opts.DialTimeout = *seedTimeout + opts.ReadTimeout = *seedTimeout + opts.WriteTimeout = *seedTimeout + opts.PoolTimeout = *seedTimeout + 5*time.Second + + switch { + case *inspectMode: + os.Exit(runInspect(opts)) + case *flushMode: + os.Exit(runFlush(opts)) + case *verifyMode: + os.Exit(runVerify(opts, *seedSeed, *seedMB*1024*1024, *seedValueBytes, *seedBatch)) + case *seedMode: + os.Exit(runSeed(opts, *seedSeed, *seedMB*1024*1024, *seedValueBytes, *seedBatch, *seedThrottleMB*1024*1024)) + case *watchMode: + os.Exit(runWatch(opts, *watchInterval)) + default: + os.Exit(runBasicTests(opts)) + } +} + +// runWatch PINGs the target every interval until SIGINT/SIGTERM, reporting the +// up→down and down→up edges so you can measure failover/outage windows. The +// per-probe timeout is opts.ReadTimeout (set from --timeout); lower it for +// tighter outage-edge resolution. Detection granularity ≈ max(interval, timeout). +func runWatch(opts *redis.Options, interval time.Duration) int { + const tsLayout = "15:04:05.000" + fmt.Printf("watch: addr=%s db=%d tls=%t password=%s interval=%s timeout=%s (Ctrl+C to stop)\n", + opts.Addr, opts.DB, opts.TLSConfig != nil, redactPassword(opts.Password), interval, opts.ReadTimeout) + + rdb := redis.NewClient(opts) defer rdb.Close() - if *continuous { - runContinuous(rdb, *interval) - return + ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer stop() + + ticker := time.NewTicker(interval) + defer ticker.Stop() + + probe := func() error { + pctx, cancel := context.WithTimeout(ctx, opts.ReadTimeout) + defer cancel() + return rdb.Ping(pctx).Err() } + var ( + started = time.Now() + up = true + first = true + outageStart time.Time + lastBeat = time.Now() + outages int + totalDown time.Duration + longest time.Duration + ) + // endOutage closes the current outage, accumulates it, and prints recovery. + endOutage := func(at time.Time, recovered bool) { + d := at.Sub(outageStart) + totalDown += d + if d > longest { + longest = d + } + if recovered { + fmt.Printf("%s RECOVERED — outage lasted %s\n", at.Format(tsLayout), d.Round(time.Millisecond)) + } else { + fmt.Printf("%s OUTAGE ONGOING at exit — %s so far\n", at.Format(tsLayout), d.Round(time.Millisecond)) + } + } + + for { + select { + case <-ctx.Done(): + now := time.Now() + if !up { + endOutage(now, false) + } + fmt.Printf("\nstopped after %s: outages=%d total_down=%s longest=%s\n", + time.Since(started).Round(time.Millisecond), outages, + totalDown.Round(time.Millisecond), longest.Round(time.Millisecond)) + return 0 + case <-ticker.C: + } + + err := probe() + // SIGINT can land mid-probe; treat a cancelled probe as shutdown, not an outage. + if ctx.Err() != nil { + continue + } + nowUp := err == nil + now := time.Now() + + switch { + case first: + first = false + up = nowUp + if nowUp { + fmt.Printf("%s up\n", now.Format(tsLayout)) + } else { + outageStart, outages = now, outages+1 + fmt.Printf("%s DOWN at start — %s\n", now.Format(tsLayout), trim(err.Error(), 120)) + } + case up && !nowUp: + outageStart, outages = now, outages+1 + fmt.Printf("%s OUTAGE START — %s\n", now.Format(tsLayout), trim(err.Error(), 120)) + case !up && nowUp: + endOutage(now, true) + case up && nowUp && now.Sub(lastBeat) >= 10*time.Second: + fmt.Printf("%s … up (elapsed %s, outages %d)\n", + now.Format(tsLayout), time.Since(started).Round(time.Second), outages) + lastBeat = now + } + up = nowUp + } +} + +// runBasicTests runs the one-shot smoke suite against a single Valkey resolved +// from --conn or the split connection flags. +func runBasicTests(opts *redis.Options) int { + fmt.Printf("redis config: addr=%s db=%d tls=%t password=%s\n", + opts.Addr, opts.DB, opts.TLSConfig != nil, redactPassword(opts.Password)) + rdb := redis.NewClient(opts) + defer rdb.Close() + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) defer cancel() @@ -352,8 +431,9 @@ func main() { fmt.Printf("\n%d/%d tests passed\n", len(tests)-failed, len(tests)) if failed > 0 { - os.Exit(1) + return 1 } + return 0 } func testPing(ctx context.Context, r *redis.Client) error { @@ -560,160 +640,3 @@ func testCleanup(ctx context.Context, r *redis.Client) error { } return r.Del(ctx, keys...).Err() } - -func runContinuous(r *redis.Client, interval time.Duration) { - stopCtx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) - defer stop() - - subCtx, cancelSub := context.WithCancel(stopCtx) - defer cancelSub() - sub := r.Subscribe(subCtx, "vk:cont:channel") - defer sub.Close() - if _, err := sub.Receive(subCtx); err != nil { - log.Fatalf("subscribe: %v", err) - } - subCh := sub.Channel() - - ops := []struct { - name string - fn func(context.Context, *redis.Client, int64) error - }{ - {"PING", opPing}, - {"SET/GET", opSetGet}, - {"INCR", opIncr}, - {"LIST", opList}, - {"INFO", opInfo}, - {"PUB/SUB", func(ctx context.Context, r *redis.Client, i int64) error { - return opPubSub(ctx, r, i, subCh) - }}, - } - - var cycle, okCount, failCount int64 - ticker := time.NewTicker(interval) - defer ticker.Stop() - - fmt.Printf("continuous mode: addr=%s interval=%s (Ctrl+C to stop)\n", r.Options().Addr, interval) - - for { - select { - case <-stopCtx.Done(): - fmt.Printf("\nstopped: cycles=%d ok=%d failed=%d\n", cycle, okCount, failCount) - return - case <-ticker.C: - } - cycle++ - - results := make([]string, len(ops)) - var failures []string - cycleStart := time.Now() - for i, op := range ops { - ctx, cancel := context.WithTimeout(stopCtx, interval) - start := time.Now() - err := op.fn(ctx, r, cycle) - cancel() - dur := time.Since(start) - if err != nil { - results[i] = fmt.Sprintf("%s=FAIL", op.name) - failures = append(failures, fmt.Sprintf(" %s (%s): %v", op.name, dur, err)) - } else { - results[i] = fmt.Sprintf("%s=ok(%s)", op.name, dur.Round(time.Microsecond)) - } - } - - if len(failures) == 0 { - okCount++ - } else { - failCount++ - } - - fmt.Printf("[%s] #%d %s | total=%s\n", - time.Now().Format("15:04:05"), cycle, joinResults(results), time.Since(cycleStart).Round(time.Microsecond)) - for _, f := range failures { - fmt.Println(f) - } - } -} - -func joinResults(parts []string) string { - out := "" - for i, p := range parts { - if i > 0 { - out += " " - } - out += p - } - return out -} - -func opPing(ctx context.Context, r *redis.Client, _ int64) error { - pong, err := r.Ping(ctx).Result() - if err != nil { - return err - } - if pong != "PONG" { - return fmt.Errorf("expected PONG, got %q", pong) - } - return nil -} - -func opSetGet(ctx context.Context, r *redis.Client, i int64) error { - val := fmt.Sprintf("v-%d", i) - if err := r.Set(ctx, "vk:cont:str", val, 10*time.Second).Err(); err != nil { - return err - } - got, err := r.Get(ctx, "vk:cont:str").Result() - if err != nil { - return err - } - if got != val { - return fmt.Errorf("set %q, got %q", val, got) - } - return nil -} - -func opIncr(ctx context.Context, r *redis.Client, _ int64) error { - _, err := r.Incr(ctx, "vk:cont:counter").Result() - return err -} - -func opList(ctx context.Context, r *redis.Client, i int64) error { - if err := r.RPush(ctx, "vk:cont:list", fmt.Sprintf("item-%d", i)).Err(); err != nil { - return err - } - if err := r.LTrim(ctx, "vk:cont:list", -10, -1).Err(); err != nil { - return err - } - _, err := r.LRange(ctx, "vk:cont:list", 0, -1).Result() - return err -} - -func opInfo(ctx context.Context, r *redis.Client, _ int64) error { - info, err := r.Info(ctx, "server").Result() - if err != nil { - return err - } - if len(info) == 0 { - return fmt.Errorf("empty INFO") - } - return nil -} - -func opPubSub(ctx context.Context, r *redis.Client, i int64, ch <-chan *redis.Message) error { - payload := fmt.Sprintf("ping-%d", i) - if err := r.Publish(ctx, "vk:cont:channel", payload).Err(); err != nil { - return err - } - for { - select { - case msg := <-ch: - if msg == nil { - return fmt.Errorf("channel closed") - } - if msg.Payload == payload { - return nil - } - case <-ctx.Done(): - return fmt.Errorf("timeout waiting for %q", payload) - } - } -} diff --git a/valkey/zerops-chaos-sea-import.yaml b/valkey/zerops-chaos-sea-import.yaml new file mode 100644 index 0000000..d26c9f0 --- /dev/null +++ b/valkey/zerops-chaos-sea-import.yaml @@ -0,0 +1,17 @@ +# Chaos infra for sea1 (valkeys already exist). Adds the dashboard + a small +# worker fleet. Worker fleet kept modest (max 3) for a LIGHT project. +# +# Import: zcli project service-import zerops-chaos-sea-import.yaml -P +# Deploy: zcli push -P --service-id --setup chaoscollector +# zcli push -P --service-id --setup chaosworker-sea +services: + - hostname: chaoscollector + type: go@1 + enableSubdomainAccess: true + minContainers: 1 + maxContainers: 1 + + - hostname: chaosworker + type: go@1 + minContainers: 1 + maxContainers: 3 diff --git a/valkey/zerops.yaml b/valkey/zerops.yaml index 3819be4..14357de 100644 --- a/valkey/zerops.yaml +++ b/valkey/zerops.yaml @@ -86,3 +86,27 @@ zerops: VALKEY_VALKEYHA1_PASSWORD: ${valkeyha1_password} VALKEY_VALKEYHA2_PASSWORD: ${valkeyha2_password} VALKEY_VALKEYHA3_PASSWORD: ${valkeyha3_password} + + # sea1 chaos worker against the HA valkeys, per-service auth over plaintext. + # The valkey TLS ports require mutual TLS (client cert), which the tester + # doesn't present, so the workers use the plaintext ports (config-sea-ha.yaml) + # with passwords. VALKEY_PASSWORD is intentionally NOT set so the per-worker + # ${valkeyhaN_password} values survive (applyZeropsOverrides only clobbers + # passwords when VALKEY_PASSWORD is present). + - setup: chaosworker-sea + build: + base: ubuntu/go@1 + buildCommands: + - go build -o app . + deployFiles: + - app + - config-sea-ha.yaml + cache: true + run: + base: ubuntu@latest + start: ./app --worker --config config-sea-ha.yaml + envVariables: + COLLECTOR_URL: http://chaoscollector.zerops:8080 + valkeyha1_password: ${valkeyha1_password} + valkeyha2_password: ${valkeyha2_password} + valkeyha3_password: ${valkeyha3_password}