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