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

898 lines
28 KiB
Go

package main
import (
"context"
"encoding/json"
"fmt"
"log"
"net/http"
"os"
"os/signal"
"sort"
"strconv"
"sync"
"syscall"
"time"
)
// ---------------------------------------------------------------------------
// Distributed collector — HTTP dashboard that ingests events from workers
// ---------------------------------------------------------------------------
//
// Endpoints:
// POST /events ingest a wireBatch from a worker container
// GET /api/state JSON snapshot for the live dashboard
// GET /api/events?cursor= JSON event log slice past the cursor
// GET / the dashboard HTML page
//
// Data model:
// workers keyed by "<containerID>/<workerName>" — same uiStats shape the
// TUI uses. A 1-second tick goroutine computes throughput deltas and p95
// over a sliding window of recent latencies.
type collectorWorker struct {
uiStats
containerID string
workerName string
}
type collectorState struct {
mu sync.Mutex
startAt time.Time
// keyed by "<container>/<worker>"; iteration order preserved via slice
keys []string
workers map[string]*collectorWorker
events []eventLine
nextCursor int64
inconsistencies []eventLine
nextInconsCursor int64
// global series for the top-of-dashboard charts
p95Series []float64
okPerSecSeries []float64
failPerSecSeries []float64
missingPerSecSeries []float64
globalLat []time.Duration
prevTotalOps int64
prevTotalFail int64
prevTotalMissing int64
}
func newCollectorState() *collectorState {
return &collectorState{
startAt: time.Now(),
workers: make(map[string]*collectorWorker),
}
}
// ---- ingest -------------------------------------------------------------
func (c *collectorState) ingest(b wireBatch) {
c.mu.Lock()
defer c.mu.Unlock()
for _, ev := range b.Events {
switch ev.Type {
case "op":
c.applyOp(b.ContainerID, ev)
case "sub":
c.applySub(b.ContainerID, ev)
case "ledger":
c.applyLedger(b.ContainerID, ev)
case "info":
c.appendEvent(ev.Level, ev.At, ev.Text)
}
}
}
func (c *collectorState) workerByKey(containerID, workerName string) *collectorWorker {
key := containerID + "/" + workerName
w, ok := c.workers[key]
if !ok {
w = &collectorWorker{containerID: containerID, workerName: workerName}
c.workers[key] = w
c.keys = append(c.keys, key)
sort.Strings(c.keys)
}
return w
}
func (c *collectorState) applyOp(containerID string, ev wireEvent) {
w := c.workerByKey(containerID, ev.Worker)
if w.role == "" {
w.role = roleFromOp(ev.Op)
}
w.ops++
if !ev.OK {
// Surface every individual failure in the dashboard event log so
// the operator can see what's actually breaking, not just the
// summarized OUTAGE START/END boundaries.
c.appendEvent("err", ev.At,
fmt.Sprintf("[%s/%s] %s FAIL: %s",
containerID, ev.Worker, ev.Op, trim(ev.Err, 200)))
w.fail++
w.consecFail++
if !w.inOutage {
w.inOutage = true
w.outageStart = ev.At
w.outageOpsLost = 1
c.appendEvent("err", ev.At,
fmt.Sprintf("[%s/%s] OUTAGE START", containerID, ev.Worker))
} else {
w.outageOpsLost++
}
return
}
w.ok++
if w.inOutage {
w.outages = append(w.outages, outageRec{
start: w.outageStart, end: ev.At, opsLost: w.outageOpsLost,
})
w.inOutage = false
w.consecFail = 0
w.outageOpsLost = 0
}
if ev.DurUS > 0 {
dur := time.Duration(ev.DurUS) * time.Microsecond
w.recentLat = appendBoundedDur(w.recentLat, dur, maxRecentLat)
c.globalLat = appendBoundedDur(c.globalLat, dur, maxGlobalLat)
}
}
func (c *collectorState) applySub(containerID string, ev wireEvent) {
w := c.workerByKey(containerID, ev.Worker)
if w.role == "" {
w.role = "subscriber"
}
w.msgs++
if w.lastSeq != 0 && ev.Seq > w.lastSeq+1 {
gap := ev.Seq - w.lastSeq - 1
w.gaps += gap
}
if !w.lastSeqAt.IsZero() {
if d := ev.At.Sub(w.lastSeqAt); d > w.longestGap {
w.longestGap = d
}
}
w.lastSeq = ev.Seq
w.lastSeqAt = ev.At
}
func (c *collectorState) applyLedger(containerID string, ev wireEvent) {
w := c.workerByKey(containerID, ev.Worker)
if w.role == "" {
w.role = "reader"
}
switch ev.Kind {
case "ok":
w.verified++
case "missing":
w.missing++
c.appendInconsistency("err", ev.At,
fmt.Sprintf("[%s/%s] LEDGER MISSING seq=%d %s",
containerID, ev.Worker, ev.Seq, ev.Detail))
case "mismatch":
w.mismatch++
c.appendInconsistency("err", ev.At,
fmt.Sprintf("[%s/%s] LEDGER MISMATCH seq=%d %s",
containerID, ev.Worker, ev.Seq, ev.Detail))
}
}
// roleFromOp infers a worker's role from the first op name we see, since
// workers don't currently transmit role with each event. Best-effort only.
func roleFromOp(op string) string {
switch op {
case "WRITE":
return "writer"
case "RW":
return "readwriter"
case "GET", "PING":
return "reader"
case "PUB":
return "publisher"
case "SUBSCRIBE":
return "subscriber"
case "BLOAT":
return "bloater"
}
return ""
}
func (c *collectorState) appendEvent(level string, at time.Time, text string) {
if at.IsZero() {
at = time.Now()
}
if level == "" {
level = "info"
}
if len(c.events) >= maxEvents {
copy(c.events, c.events[1:])
c.events = c.events[:len(c.events)-1]
}
c.events = append(c.events, eventLine{at: at, level: level, text: text})
c.nextCursor++
}
func (c *collectorState) appendInconsistency(level string, at time.Time, text string) {
if at.IsZero() {
at = time.Now()
}
if level == "" {
level = "info"
}
if len(c.inconsistencies) >= maxEvents {
copy(c.inconsistencies, c.inconsistencies[1:])
c.inconsistencies = c.inconsistencies[:len(c.inconsistencies)-1]
}
c.inconsistencies = append(c.inconsistencies, eventLine{at: at, level: level, text: text})
c.nextInconsCursor++
}
// ---- tick (1Hz) ---------------------------------------------------------
func (c *collectorState) tick() {
c.mu.Lock()
defer c.mu.Unlock()
if len(c.globalLat) > 0 {
cp := make([]time.Duration, len(c.globalLat))
copy(cp, c.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)
c.p95Series = appendBoundedFloat(c.p95Series, p95ms, maxSeriesPoints)
} else {
c.p95Series = appendBoundedFloat(c.p95Series, 0, maxSeriesPoints)
}
var totalOps, totalFail, totalMissing int64
for _, w := range c.workers {
totalOps += w.ops + w.msgs
totalFail += w.fail
totalMissing += w.missing
}
okDelta := (totalOps - totalFail) - (c.prevTotalOps - c.prevTotalFail)
if okDelta < 0 {
okDelta = 0
}
failDelta := totalFail - c.prevTotalFail
if failDelta < 0 {
failDelta = 0
}
missingDelta := totalMissing - c.prevTotalMissing
if missingDelta < 0 {
missingDelta = 0
}
c.prevTotalOps = totalOps
c.prevTotalFail = totalFail
c.prevTotalMissing = totalMissing
c.okPerSecSeries = appendBoundedFloat(c.okPerSecSeries, float64(okDelta), maxSeriesPoints)
c.failPerSecSeries = appendBoundedFloat(c.failPerSecSeries, float64(failDelta), maxSeriesPoints)
c.missingPerSecSeries = appendBoundedFloat(c.missingPerSecSeries, float64(missingDelta), maxSeriesPoints)
for _, w := range c.workers {
current := w.ops + w.msgs
delta := current - w.prevOps
if delta < 0 {
delta = 0
}
w.prevOps = current
w.opsHist = appendBoundedFloat(w.opsHist, float64(delta), maxSparkSamples)
}
}
// ---- snapshot view ------------------------------------------------------
type stateView struct {
StartAt time.Time `json:"start_at"`
Elapsed string `json:"elapsed"`
Containers int `json:"containers"`
InOutage int `json:"in_outage"`
Workers []workerView `json:"workers"`
Charts chartsView `json:"charts"`
}
type workerView struct {
Container string `json:"container"`
Name string `json:"name"`
Role string `json:"role"`
Ops int64 `json:"ops"`
OK int64 `json:"ok"`
Fail int64 `json:"fail"`
UptimePct float64 `json:"uptime_pct"`
P50MS float64 `json:"p50_ms"`
P95MS float64 `json:"p95_ms"`
OpsHist []float64 `json:"ops_hist"`
InOutage bool `json:"in_outage"`
Outages int `json:"outages"`
LongestMS float64 `json:"longest_outage_ms"`
// reader-only
Verified int64 `json:"verified,omitempty"`
Missing int64 `json:"missing,omitempty"`
Mismatch int64 `json:"mismatch,omitempty"`
// subscriber-only
Msgs int64 `json:"msgs,omitempty"`
Gaps int64 `json:"gaps,omitempty"`
LongestGapMS float64 `json:"longest_gap_ms,omitempty"`
}
type chartsView struct {
P95 []float64 `json:"p95"`
OkPerSec []float64 `json:"ok_per_sec"`
FailPerSec []float64 `json:"fail_per_sec"`
MissingPerSec []float64 `json:"missing_per_sec"`
P95Peak float64 `json:"p95_peak"`
OkPeak float64 `json:"ok_peak"`
FailPeak float64 `json:"fail_peak"`
MissingPeak float64 `json:"missing_peak"`
}
func (c *collectorState) snapshot() stateView {
c.mu.Lock()
defer c.mu.Unlock()
containers := map[string]struct{}{}
for _, k := range c.keys {
w := c.workers[k]
containers[w.containerID] = struct{}{}
}
view := stateView{
StartAt: c.startAt,
Elapsed: time.Since(c.startAt).Round(time.Second).String(),
Containers: len(containers),
Charts: chartsView{
P95: append([]float64(nil), c.p95Series...),
OkPerSec: append([]float64(nil), c.okPerSecSeries...),
FailPerSec: append([]float64(nil), c.failPerSecSeries...),
MissingPerSec: append([]float64(nil), c.missingPerSecSeries...),
P95Peak: maxFloat(c.p95Series),
OkPeak: maxFloat(c.okPerSecSeries),
FailPeak: maxFloat(c.failPerSecSeries),
MissingPeak: maxFloat(c.missingPerSecSeries),
},
}
for _, k := range c.keys {
w := c.workers[k]
uptime := 100.0
if w.ops > 0 {
uptime = float64(w.ok) / float64(w.ops) * 100
}
var longest time.Duration
for _, o := range w.outages {
if d := o.end.Sub(o.start); d > longest {
longest = d
}
}
if w.inOutage {
view.InOutage++
}
wv := workerView{
Container: w.containerID,
Name: w.workerName,
Role: w.role,
Ops: w.ops,
OK: w.ok,
Fail: w.fail,
UptimePct: uptime,
P50MS: float64(percentileLat(w.recentLat, 50)) / float64(time.Millisecond),
P95MS: float64(percentileLat(w.recentLat, 95)) / float64(time.Millisecond),
OpsHist: append([]float64(nil), w.opsHist...),
InOutage: w.inOutage,
Outages: len(w.outages),
LongestMS: float64(longest) / float64(time.Millisecond),
}
switch w.role {
case "reader":
wv.Verified = w.verified
wv.Missing = w.missing
wv.Mismatch = w.mismatch
case "subscriber":
wv.Msgs = w.msgs
wv.Gaps = w.gaps
wv.LongestGapMS = float64(w.longestGap) / float64(time.Millisecond)
}
view.Workers = append(view.Workers, wv)
}
return view
}
type eventOut struct {
ID int64 `json:"id"`
At string `json:"at"`
Level string `json:"level"`
Text string `json:"text"`
}
func (c *collectorState) eventsSince(cursor int64) ([]eventOut, int64) {
c.mu.Lock()
defer c.mu.Unlock()
// nextCursor counts events ever inserted; events slice may have dropped older.
// Compute base id of the first slot in the current slice:
baseID := c.nextCursor - int64(len(c.events))
if cursor < baseID {
cursor = baseID
}
startIdx := int(cursor - baseID)
if startIdx < 0 {
startIdx = 0
}
if startIdx > len(c.events) {
startIdx = len(c.events)
}
out := make([]eventOut, 0, len(c.events)-startIdx)
for i := startIdx; i < len(c.events); i++ {
e := c.events[i]
out = append(out, eventOut{
ID: baseID + int64(i),
At: e.at.Format("15:04:05.000"),
Level: e.level,
Text: e.text,
})
}
return out, c.nextCursor
}
func (c *collectorState) inconsistenciesSince(cursor int64) ([]eventOut, int64) {
c.mu.Lock()
defer c.mu.Unlock()
baseID := c.nextInconsCursor - int64(len(c.inconsistencies))
if cursor < baseID {
cursor = baseID
}
startIdx := int(cursor - baseID)
if startIdx < 0 {
startIdx = 0
}
if startIdx > len(c.inconsistencies) {
startIdx = len(c.inconsistencies)
}
out := make([]eventOut, 0, len(c.inconsistencies)-startIdx)
for i := startIdx; i < len(c.inconsistencies); i++ {
e := c.inconsistencies[i]
out = append(out, eventOut{
ID: baseID + int64(i),
At: e.at.Format("15:04:05.000"),
Level: e.level,
Text: e.text,
})
}
return out, c.nextInconsCursor
}
// ---- HTTP server --------------------------------------------------------
func runCollector(port int) int {
state := newCollectorState()
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
go func() {
t := time.NewTicker(tickInterval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
state.tick()
}
}
}()
mux := http.NewServeMux()
mux.HandleFunc("/events", state.handleIngest)
mux.HandleFunc("/api/state", state.handleState)
mux.HandleFunc("/api/events", state.handleEvents)
mux.HandleFunc("/api/inconsistencies", state.handleInconsistencies)
mux.HandleFunc("/healthz", func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
})
mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/" {
http.NotFound(w, r)
return
}
w.Header().Set("Content-Type", "text/html; charset=utf-8")
_, _ = w.Write([]byte(dashboardHTML))
})
srv := &http.Server{
Addr: fmt.Sprintf(":%d", port),
Handler: mux,
ReadHeaderTimeout: 5 * time.Second,
}
go func() {
log.Printf("collector listening on :%d", port)
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
log.Fatalf("collector: %v", err)
}
}()
<-ctx.Done()
log.Printf("collector shutting down…")
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
_ = srv.Shutdown(shutdownCtx)
return 0
}
func (c *collectorState) handleIngest(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
http.Error(w, "POST only", http.StatusMethodNotAllowed)
return
}
var b wireBatch
if err := json.NewDecoder(http.MaxBytesReader(w, r.Body, 1<<20)).Decode(&b); err != nil {
http.Error(w, "bad request: "+err.Error(), http.StatusBadRequest)
return
}
c.ingest(b)
w.WriteHeader(http.StatusNoContent)
}
func (c *collectorState) handleState(w http.ResponseWriter, _ *http.Request) {
view := c.snapshot()
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(view)
}
func (c *collectorState) handleEvents(w http.ResponseWriter, r *http.Request) {
cursor := int64(0)
if v := r.URL.Query().Get("cursor"); v != "" {
if n, err := strconv.ParseInt(v, 10, 64); err == nil {
cursor = n
}
}
events, next := c.eventsSince(cursor)
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(map[string]any{
"events": events,
"next_cursor": next,
})
}
func (c *collectorState) handleInconsistencies(w http.ResponseWriter, r *http.Request) {
cursor := int64(0)
if v := r.URL.Query().Get("cursor"); v != "" {
if n, err := strconv.ParseInt(v, 10, 64); err == nil {
cursor = n
}
}
events, next := c.inconsistenciesSince(cursor)
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(map[string]any{
"events": events,
"next_cursor": next,
})
}
// ---- embedded dashboard HTML -------------------------------------------
const dashboardHTML = `<!DOCTYPE html>
<html lang="en">
<head>
<meta charset="utf-8">
<title>valkey-ha chaos</title>
<meta name="viewport" content="width=device-width,initial-scale=1">
<style>
:root {
--bg: #0c0c10;
--panel: #15151c;
--border: #2a2a35;
--text: #e6e6ea;
--dim: #6c6c7a;
--accent: #d75faf;
--ok: #5fd75f;
--warn: #ffaf5f;
--err: #ff5f5f;
--chart: #af87ff;
--chart-fail: #ff5f5f;
}
* { box-sizing: border-box; }
body { font-family: ui-monospace, 'JetBrains Mono', Menlo, Consolas, monospace;
background: var(--bg); color: var(--text); margin: 0; padding: 16px;
font-size: 13px; line-height: 1.45; }
h1 { color: var(--accent); margin: 0 0 4px; font-size: 18px; }
.header-meta { color: var(--dim); margin-bottom: 14px; }
.badge { padding: 1px 6px; border-radius: 3px; font-weight: 600; }
.badge.ok { background: #1c3a1c; color: var(--ok); }
.badge.err { background: #401818; color: var(--err); }
.badge.warn { background: #402a18; color: var(--warn); }
.grid { display: grid; grid-template-columns: 1fr 1fr 1fr 1fr; gap: 14px; margin-bottom: 16px; }
@media (max-width: 1400px) { .grid { grid-template-columns: 1fr 1fr; } }
@media (max-width: 700px) { .grid { grid-template-columns: 1fr; } }
.panel { background: var(--panel); border: 1px solid var(--border);
border-radius: 4px; padding: 10px 12px; }
.panel h3 { margin: 0 0 6px; color: #5fafff; font-size: 12px;
text-transform: uppercase; letter-spacing: 0.06em; }
.panel .meta { color: var(--dim); font-size: 11px; margin-bottom: 4px; }
table { border-collapse: collapse; width: 100%; font-size: 12px; }
th, td { padding: 4px 8px; text-align: left; vertical-align: middle;
border-bottom: 1px solid var(--border); }
th { color: #5fafff; font-weight: 600; text-transform: uppercase;
font-size: 10px; letter-spacing: 0.05em; }
td.num { text-align: right; font-variant-numeric: tabular-nums; }
tr.outage td:first-child { border-left: 3px solid var(--err); }
.container-tag { color: var(--accent); }
.events { background: var(--panel); border: 1px solid var(--border);
border-radius: 4px; padding: 8px 10px; max-height: 50vh;
overflow-y: auto; font-size: 12px; }
.events .ev { padding: 1px 0; white-space: pre-wrap; word-break: break-word; }
.events .ev .ts { color: var(--dim); margin-right: 6px; }
.events .ev.err { color: var(--err); }
.events .ev.warn { color: var(--warn); }
.events .ev.info { color: var(--text); }
svg { display: block; width: 100%; height: 80px; }
polyline { fill: none; stroke: var(--chart); stroke-width: 1.4; }
.area { fill: var(--chart); fill-opacity: 0.18; stroke: none; }
.spark { font-family: ui-monospace, monospace; color: var(--chart); }
.dim { color: var(--dim); }
.filter { margin-bottom: 8px; }
.filter input { background: var(--panel); color: var(--text);
border: 1px solid var(--border); padding: 4px 8px;
border-radius: 3px; font: inherit; }
.footer { color: var(--dim); margin-top: 12px; font-size: 11px; }
</style>
</head>
<body>
<h1>valkey-ha chaos</h1>
<div class="header-meta" id="header">connecting…</div>
<div class="grid">
<div class="panel">
<h3>p95 latency (rolling 2 min)</h3>
<div class="meta" id="lat-meta">—</div>
<svg id="chart-lat" viewBox="0 0 100 30" preserveAspectRatio="none"></svg>
</div>
<div class="panel">
<h3>ok ops/sec (rolling 2 min)</h3>
<div class="meta" id="tp-meta">—</div>
<svg id="chart-tp" viewBox="0 0 100 30" preserveAspectRatio="none"></svg>
</div>
<div class="panel">
<h3>fail ops/sec (rolling 2 min)</h3>
<div class="meta" id="fail-meta">—</div>
<svg id="chart-fail" viewBox="0 0 100 30" preserveAspectRatio="none"></svg>
</div>
<div class="panel">
<h3>missing ledger/sec (rolling 2 min)</h3>
<div class="meta" id="miss-meta">—</div>
<svg id="chart-miss" viewBox="0 0 100 30" preserveAspectRatio="none"></svg>
</div>
</div>
<div class="panel" style="margin-bottom:16px">
<h3>workers</h3>
<div class="filter">
<input id="filter" placeholder="filter by container or worker name" />
</div>
<table id="workers">
<thead><tr>
<th>CONTAINER</th><th>NAME</th><th>ROLE</th>
<th class="num">OPS</th><th class="num">OK</th><th class="num">FAIL</th>
<th class="num">UPTIME</th><th class="num">P50</th><th class="num">P95</th>
<th>SPARK</th><th class="num">OUTAGES</th><th>STATUS</th><th>NOTES</th>
</tr></thead>
<tbody></tbody>
</table>
</div>
<div>
<h3 style="color:#5fafff;font-size:12px;text-transform:uppercase;
letter-spacing:0.06em;margin:0 0 6px">client errors</h3>
<div class="events" id="events"></div>
</div>
<div style="margin-top:16px">
<h3 style="color:#5fafff;font-size:12px;text-transform:uppercase;
letter-spacing:0.06em;margin:0 0 6px">ledger inconsistencies</h3>
<div class="events" id="inconsistencies"></div>
</div>
<div class="footer" id="footer">collector dashboard · poll 1s</div>
<script>
const sparkChars = [' ','▁','▂','▃','▄','▅','▆','▇','█'];
function spark(values, width) {
if (!values || !values.length) return '';
const max = Math.max(...values);
if (max <= 0) return ' '.repeat(width);
let out = '';
for (let col = 0; col < width; col++) {
const idx = Math.min(values.length - 1, Math.floor(col * values.length / width));
const v = values[idx];
const ratio = v / max;
const i = Math.max(0, Math.min(8, Math.floor(ratio * 8)));
out += sparkChars[i];
}
return out;
}
function drawChart(svgID, values, color) {
const stroke = color || 'var(--chart)';
const svg = document.getElementById(svgID);
svg.innerHTML = '';
if (!values || !values.length) return;
const max = Math.max(1, ...values);
const w = 100, h = 30;
const step = w / Math.max(1, values.length - 1);
let path = '';
for (let i = 0; i < values.length; i++) {
const x = i * step;
const y = h - (values[i] / max) * h;
path += (i === 0 ? 'M' : 'L') + x.toFixed(2) + ',' + y.toFixed(2) + ' ';
}
const ns = 'http://www.w3.org/2000/svg';
const area = document.createElementNS(ns, 'path');
area.setAttribute('fill', stroke);
area.setAttribute('fill-opacity', '0.18');
area.setAttribute('stroke', 'none');
area.setAttribute('d', path + 'L' + w + ',' + h + ' L0,' + h + ' Z');
svg.appendChild(area);
const line = document.createElementNS(ns, 'path');
line.setAttribute('d', path);
line.setAttribute('fill', 'none');
line.setAttribute('stroke', stroke);
line.setAttribute('stroke-width', '1.2');
svg.appendChild(line);
}
function fmt(n) { return n == null ? '-' : n.toLocaleString(); }
function fmtMS(v) {
if (v == null || v === 0) return '-';
if (v < 1) return (v * 1000).toFixed(0) + 'µs';
if (v < 1000) return v.toFixed(1) + 'ms';
return (v / 1000).toFixed(2) + 's';
}
let cursor = 0;
let inconsCursor = 0;
let lastWorkers = [];
let filterStr = '';
document.getElementById('filter').addEventListener('input', e => {
filterStr = e.target.value.toLowerCase();
renderTable(lastWorkers);
});
function renderTable(workers) {
const tbody = document.querySelector('#workers tbody');
tbody.innerHTML = '';
for (const w of workers) {
if (filterStr) {
const hay = (w.container + ' ' + w.name + ' ' + w.role).toLowerCase();
if (!hay.includes(filterStr)) continue;
}
const tr = document.createElement('tr');
if (w.in_outage) tr.classList.add('outage');
const status = w.in_outage
? '<span class="badge err">OUTAGE</span>'
: (w.role === 'subscriber' && w.msgs === 0
? '<span class="badge warn">NO MSGS</span>'
: '<span class="badge ok">OK</span>');
let notes = '';
if (w.role === 'reader') {
notes = '<span class="dim">v=' + fmt(w.verified) +
' miss=' + fmt(w.missing) +
' mis=' + fmt(w.mismatch) + '</span>';
} else if (w.role === 'subscriber') {
notes = '<span class="dim">msgs=' + fmt(w.msgs) +
' gaps=' + fmt(w.gaps) +
' longest_gap=' + fmtMS(w.longest_gap_ms) + '</span>';
}
tr.innerHTML =
'<td><span class="container-tag">' + w.container + '</span></td>' +
'<td>' + w.name + '</td>' +
'<td>' + (w.role || '?') + '</td>' +
'<td class="num">' + fmt(w.ops) + '</td>' +
'<td class="num">' + fmt(w.ok) + '</td>' +
'<td class="num">' + fmt(w.fail) + '</td>' +
'<td class="num">' + (w.uptime_pct != null ? w.uptime_pct.toFixed(2) + '%' : '-') + '</td>' +
'<td class="num">' + fmtMS(w.p50_ms) + '</td>' +
'<td class="num">' + fmtMS(w.p95_ms) + '</td>' +
'<td><span class="spark">' + spark(w.ops_hist, 14) + '</span></td>' +
'<td class="num">' + fmt(w.outages) + '</td>' +
'<td>' + status + '</td>' +
'<td>' + notes + '</td>';
tbody.appendChild(tr);
}
}
async function tickState() {
try {
const s = await (await fetch('/api/state')).json();
const tag = s.in_outage > 0
? '<span class="badge err">' + s.in_outage + ' worker(s) in outage</span>'
: '<span class="badge ok">healthy</span>';
document.getElementById('header').innerHTML =
s.containers + ' container(s) · ' + s.workers.length + ' worker(s) · elapsed ' +
s.elapsed + ' · ' + tag;
document.getElementById('lat-meta').textContent =
'peak ' + s.charts.p95_peak.toFixed(1) + 'ms · last ' +
(s.charts.p95.length ? s.charts.p95[s.charts.p95.length-1].toFixed(1) + 'ms' : '-');
document.getElementById('tp-meta').textContent =
'peak ' + Math.round(s.charts.ok_peak) + ' · last ' +
(s.charts.ok_per_sec.length ? Math.round(s.charts.ok_per_sec[s.charts.ok_per_sec.length-1]) : '-');
document.getElementById('fail-meta').textContent =
'peak ' + Math.round(s.charts.fail_peak) + ' · last ' +
(s.charts.fail_per_sec.length ? Math.round(s.charts.fail_per_sec[s.charts.fail_per_sec.length-1]) : '-');
document.getElementById('miss-meta').textContent =
'peak ' + Math.round(s.charts.missing_peak) + ' · last ' +
(s.charts.missing_per_sec.length ? Math.round(s.charts.missing_per_sec[s.charts.missing_per_sec.length-1]) : '-');
drawChart('chart-lat', s.charts.p95);
drawChart('chart-tp', s.charts.ok_per_sec);
drawChart('chart-fail', s.charts.fail_per_sec, 'var(--chart-fail)');
drawChart('chart-miss', s.charts.missing_per_sec, 'var(--chart-fail)');
lastWorkers = s.workers;
renderTable(s.workers);
} catch (e) {
document.getElementById('header').textContent = 'fetch error: ' + e;
}
}
async function tickEvents() {
try {
const r = await (await fetch('/api/events?cursor=' + cursor)).json();
cursor = r.next_cursor;
const div = document.getElementById('events');
const wasAtBottom = div.scrollHeight - div.scrollTop - div.clientHeight < 30;
for (const e of r.events) {
const d = document.createElement('div');
d.className = 'ev ' + (e.level || 'info');
const ts = document.createElement('span');
ts.className = 'ts';
ts.textContent = e.at;
d.appendChild(ts);
d.appendChild(document.createTextNode(e.text));
div.appendChild(d);
}
while (div.children.length > 1000) div.removeChild(div.firstChild);
if (wasAtBottom) div.scrollTop = div.scrollHeight;
} catch (e) { /* noop */ }
}
async function tickInconsistencies() {
try {
const r = await (await fetch('/api/inconsistencies?cursor=' + inconsCursor)).json();
inconsCursor = r.next_cursor;
const div = document.getElementById('inconsistencies');
const wasAtBottom = div.scrollHeight - div.scrollTop - div.clientHeight < 30;
for (const e of r.events) {
const d = document.createElement('div');
d.className = 'ev ' + (e.level || 'info');
const ts = document.createElement('span');
ts.className = 'ts';
ts.textContent = e.at;
d.appendChild(ts);
d.appendChild(document.createTextNode(e.text));
div.appendChild(d);
}
while (div.children.length > 1000) div.removeChild(div.firstChild);
if (wasAtBottom) div.scrollTop = div.scrollHeight;
} catch (e) { /* noop */ }
}
setInterval(tickState, 1000);
setInterval(tickEvents, 1000);
setInterval(tickInconsistencies, 1000);
tickState();
tickEvents();
tickInconsistencies();
</script>
</body>
</html>
`