package main import ( "flag" "fmt" "math" "net" "os" "os/signal" "strconv" "strings" "sync" "sync/atomic" "syscall" "time" "golang.org/x/sys/unix" ) const wireOverhead = 24 type endpoint struct { name string tag string idx int mac [6]byte mtu int speed float64 } func (e endpoint) macString() string { return fmt.Sprintf("%02x:%02x:%02x:%02x:%02x:%02x", e.mac[0], e.mac[1], e.mac[2], e.mac[3], e.mac[4], e.mac[5]) } type direction struct { label string short string tx endpoint rx endpoint specs []*frameSpec txStats []*txStats rxStats []*rxStats streams []lossWindow txFDs []int rxFDs []int probeSpec *frameSpec probeTxFD int probeRxFD int cable *cableStats prevConsole sample win *rateWindow est *rateEstimators heldFrames heldValue heldSent heldValue drops uint64 errBase sample dropBase uint64 nicRaw uint64 nicNow uint64 nicBase uint64 recent errs recentBase sample recentDrops uint64 recentNic uint64 } // Smoothing has to be steady against high-frequency noise yet still chase a // real change quickly, with bounded state. So the gain is not fixed: the // innovation is compared against a running estimate of the noise itself (mean // absolute deviation, as in TCP's rtt/rttvar), and only an innovation that // stands out above that noise is chased hard. const ( estAlphaCalm = 0.015 estAlphaSnap = 0.45 estMADBeta = 0.05 estNoiseK = 3.0 ) type rateEstimator struct { minStep float64 relStep float64 est float64 mad float64 shown float64 n int } // While the rate window is still growing it already averages everything there // is, so smoothing it again would only average in the startup ramp twice. func (e *rateEstimator) update(x float64, windowFull bool) { if !windowFull { e.est, e.shown, e.mad, e.n = x, x, 0, 0 return } if e.n == 0 { e.est, e.shown, e.n = x, x, 1 return } err := x - e.est abs := math.Abs(err) if e.n == 1 { e.mad, e.n = abs, 2 } else { e.mad += (abs - e.mad) * estMADBeta } a := estAlphaCalm if e.mad > 0 { if excess := abs/(estNoiseK*e.mad) - 1; excess > 0 { a = estAlphaCalm + (estAlphaSnap-estAlphaCalm)*math.Min(excess, 1) } } e.est += err * a // A deadband on top, so the drawn text only changes when the estimate has // actually moved rather than on every frame. if math.Abs(e.est-e.shown) > math.Max(e.minStep, e.relStep*math.Abs(e.est)) { e.shown = e.est } } func (e *rateEstimator) value() float64 { return e.shown } type rateEstimators struct { txGbps, rxGbps rateEstimator txPPS, rxPPS rateEstimator } func newRateEstimators() *rateEstimators { return &rateEstimators{ txGbps: rateEstimator{minStep: 0.02, relStep: 0.001}, rxGbps: rateEstimator{minStep: 0.02, relStep: 0.001}, txPPS: rateEstimator{minStep: 2000, relStep: 0.002}, rxPPS: rateEstimator{minStep: 2000, relStep: 0.002}, } } // Monotonic totals climb by tens of thousands per frame, which is unreadable // churn at 60Hz, so the drawn value is held and refreshed a few times a second. type heldValue struct { v uint64 at time.Time } func (h *heldValue) get(now time.Time, cur uint64) uint64 { if now.Sub(h.at) >= totalsHold { h.v, h.at = cur, now } return h.v } type rateSample struct { t time.Time txFrames, txBytes, rxFrames, rxBytes uint64 } // Late is counted but left out of the total, since a reordered frame arrived. type errs struct { lost, late uint64 crc, badMagic uint64 badLen uint64 kdrop, link uint64 } func (e errs) total() uint64 { return e.lost + e.crc + e.badMagic + e.badLen + e.kdrop + e.link } func (e errs) add(o errs) errs { return errs{ lost: e.lost + o.lost, late: e.late + o.late, crc: e.crc + o.crc, badMagic: e.badMagic + o.badMagic, badLen: e.badLen + o.badLen, kdrop: e.kdrop + o.kdrop, link: e.link + o.link, } } // A sliding window: the counters are sampled every frame and the rate is taken // across the whole window, so the figure moves every frame while still being // measured over a long enough span to be steady. type rateWindow struct { samples []rateSample idx int filled bool } func newRateWindow(n int) *rateWindow { return &rateWindow{samples: make([]rateSample, n)} } func (w *rateWindow) push(s rateSample) { w.samples[w.idx] = s w.idx++ if w.idx == len(w.samples) { w.idx = 0 w.filled = true } } func (w *rateWindow) span() (oldest, newest rateSample, ok bool) { if !w.filled && w.idx < 2 { return oldest, newest, false } n := w.idx - 1 if n < 0 { n = len(w.samples) - 1 } o := 0 if w.filled { o = w.idx } return w.samples[o], w.samples[n], true } type sample struct { txFrames, txBytes uint64 rxFrames, rxBytes uint64 lost, late uint64 crcErr, badMagic uint64 badLen uint64 txErrs, txShort uint64 rxErrs uint64 } func lookupEndpoint(name string) (endpoint, error) { ifi, err := net.InterfaceByName(name) if err != nil { return endpoint{}, err } if len(ifi.HardwareAddr) != 6 { return endpoint{}, fmt.Errorf("%s: expected 6-byte MAC, got %q", name, ifi.HardwareAddr) } var mac [6]byte copy(mac[:], ifi.HardwareAddr) e := endpoint{name: name, idx: ifi.Index, mac: mac, mtu: ifi.MTU} if v, ok := readUint("/sys/class/net/" + name + "/speed"); ok { e.speed = float64(v) / 1000 } return e, nil } func parseSizes(s string) ([]int, error) { var out []int for _, f := range strings.Split(s, ",") { f = strings.TrimSpace(f) if f == "" { continue } v, err := strconv.Atoi(f) if err != nil { return nil, fmt.Errorf("bad size %q: %w", f, err) } if v < minFrame { return nil, fmt.Errorf("size %d below minimum %d", v, minFrame) } out = append(out, v) } if len(out) == 0 { return nil, fmt.Errorf("no sizes given") } return out, nil } func (d *direction) snapshot() sample { var s sample for _, t := range d.txStats { s.txFrames += t.frames.Load() s.txBytes += t.bytes.Load() s.txErrs += t.errs.Load() s.txShort += t.short.Load() } for _, r := range d.rxStats { s.rxFrames += r.frames.Load() s.rxBytes += r.bytes.Load() s.crcErr += r.crcErr.Load() s.badMagic += r.badMagic.Load() s.badLen += r.badLen.Load() s.rxErrs += r.rxErrs.Load() } for i := range d.streams { s.lost += d.streams[i].lost.Load() s.late += d.streams[i].late.Load() } return s } // Counters keep climbing in the workers, so resetting just moves the origin // everything is measured from. Rates are deliberately left running, since they // are instantaneous and would only blink to zero and back. func (d *direction) reset() { d.sampleDrops() d.errBase = d.snapshot() d.dropBase = d.drops d.heldFrames = heldValue{} d.heldSent = heldValue{} d.nicBase = d.nicNow d.cable.reset() } // Returns the new start time, so the uptime shown alongside the totals counts // from the reset rather than from launch. func resetAll(dirs []*direction, stats *streamTable) time.Time { for _, d := range dirs { d.reset() } stats.sinceHeader = 0 fmt.Println(stats.rule("counters reset")) return time.Now() } func (d *direction) sampleDrops() { for _, fd := range d.rxFDs { d.drops += packetDrops(fd) } } func gbps(bytes, frames uint64, secs float64) float64 { return float64((bytes+frames*wireOverhead)*8) / secs / 1e9 } var intervalCols = []colSpec{ {title: "UPTIME", width: 9, right: true}, {title: "DIR", width: 5}, {title: "TX pps", width: 9, right: true}, {title: "TX Gb/s", width: 7, right: true}, {title: "RX pps", width: 9, right: true}, {title: "RX Gb/s", width: 7, right: true}, {title: "LOST", width: 11, right: true}, {title: "LATE", width: 9, right: true}, {title: "CRC", width: 7, right: true}, {title: "BADMAG", width: 7, right: true}, {title: "KDROP", width: 11, right: true}, {title: "LINK", width: 13, right: true}, {title: "ERRORS", width: 13, right: true}, {title: "MIN ns", width: 9, right: true}, {title: "LEN m", width: 6, right: true}, } // One interval's numbers, shared by the console table and the framebuffer so // both always show the same figures. type view struct { txPPS, rxPPS float64 txGbps, rxGbps float64 rxFrames, rxGot uint64 since errs window errs cable cableView } // Cumulative fields, which need no rate window and are identical for both the // console and the display. func (d *direction) counters(now sample) view { b := d.errBase return view{ rxFrames: now.rxFrames - b.rxFrames, rxGot: now.rxBytes - b.rxBytes, cable: d.cable.view(), since: errs{ lost: now.lost - b.lost, late: now.late - b.late, crc: now.crcErr - b.crcErr, badMagic: now.badMagic - b.badMagic, badLen: now.badLen - b.badLen, kdrop: d.drops - d.dropBase, // A frame the stack refused and a frame the driver dropped are the same // failure seen from either side of the ring, and never the same frame // twice: a send that fails never reaches the driver to be dropped. link: d.nicNow - d.nicBase + (now.txErrs - b.txErrs) + (now.rxErrs - b.rxErrs), }, } } func totalView(views []view) view { var t view for _, v := range views { t.txPPS += v.txPPS t.rxPPS += v.rxPPS t.txGbps += v.txGbps t.rxGbps += v.rxGbps t.rxFrames += v.rxFrames t.rxGot += v.rxGot t.since = t.since.add(v.since) t.window = t.window.add(v.window) } return t } func (d *direction) view(prev *sample, secs float64) view { now := d.snapshot() p := *prev *prev = now d.sampleDrops() v := d.counters(now) txF := now.txFrames - p.txFrames rxF := now.rxFrames - p.rxFrames v.txPPS = float64(txF) / secs v.rxPPS = float64(rxF) / secs v.txGbps = gbps(now.txBytes-p.txBytes, txF, secs) v.rxGbps = gbps(now.rxBytes-p.rxBytes, rxF, secs) return v } func (d *direction) displayView(t time.Time) view { now := d.snapshot() d.sampleDrops() d.win.push(rateSample{t, now.txFrames, now.txBytes, now.rxFrames, now.rxBytes}) v := d.counters(now) v.window = d.recent v.rxFrames = d.heldFrames.get(t, v.rxFrames) v.rxGot = d.heldSent.get(t, v.rxGot) o, n, ok := d.win.span() if !ok { return v } secs := n.t.Sub(o.t).Seconds() if secs <= 0 { return v } txF := n.txFrames - o.txFrames rxF := n.rxFrames - o.rxFrames full := d.win.filled d.est.txPPS.update(float64(txF)/secs, full) d.est.rxPPS.update(float64(rxF)/secs, full) d.est.txGbps.update(gbps(n.txBytes-o.txBytes, txF, secs), full) d.est.rxGbps.update(gbps(n.rxBytes-o.rxBytes, rxF, secs), full) v.txPPS = d.est.txPPS.value() v.rxPPS = d.est.rxPPS.value() v.txGbps = d.est.txGbps.value() v.rxGbps = d.est.rxGbps.value() return v } func (d *direction) row(elapsed time.Duration, v view, target float64, length string) []string { return []string{ uptime(elapsed), paint(d.short, cCyan), commas(uint64(v.txPPS)), rateCell(v.txGbps, target), commas(uint64(v.rxPPS)), rateCell(v.rxGbps, target), statusCell(v.since.lost), statusCell(v.since.late), statusCell(v.since.crc), statusCell(v.since.badMagic), statusCell(v.since.kdrop), statusCell(v.since.link), statusCell(v.since.total()), paint(v.cable.minText(), cCyan), paint(length, cCyan), } } // Sampled once a second, since these are sysfs reads. The recent errors roll // here rather than over the rate window because the nic counters only move at // this rate, and a shorter span would alias them into a flicker. func (d *direction) sampleNIC() { d.accumulateNIC() now := d.snapshot() b := d.recentBase d.recent = errs{ lost: now.lost - b.lost, late: now.late - b.late, crc: now.crcErr - b.crcErr, badMagic: now.badMagic - b.badMagic, badLen: now.badLen - b.badLen, kdrop: d.drops - d.recentDrops, link: d.nicNow - d.recentNic + (now.txErrs - b.txErrs) + (now.rxErrs - b.rxErrs), } d.recentBase, d.recentDrops, d.recentNic = now, d.drops, d.nicNow } func (d *direction) readNICTotal() uint64 { tx := readNIC(d.tx.name) rx := readNIC(d.rx.name) return tx.tx + rx.rx + tx.carrierDown } // Sysfs nic counters run from boot and restart from zero whenever the driver // resets its statistics, so only their forward motion is accumulated. Taking // raw differences instead charges a boot's worth of errors to the first sample // and turns a reset into a near-2^64 underflow. func (d *direction) accumulateNIC() { raw := d.readNICTotal() if raw > d.nicRaw { d.nicNow += raw - d.nicRaw } d.nicRaw = raw } // Whatever the interfaces counted before now is not ours, and no interval has // elapsed yet, so every baseline starts here and nothing is reported until the // first one completes. func (d *direction) primeCounters() { d.nicRaw = d.readNICTotal() d.reset() d.recentBase, d.recentDrops, d.recentNic = d.snapshot(), d.drops, d.nicNow } func buildDirection(label string, tx, rx endpoint, sizes []int, cfg config) (*direction, error) { d := &direction{ label: label, short: tx.tag + "→" + rx.tag, tx: tx, rx: rx, streams: newLossWindows(cfg.streams), cable: newCableStats(), } for i := 0; i < cfg.streams; i++ { et := uint16(etherBase + i) d.specs = append(d.specs, newFrameSpec(rx.mac, tx.mac, et, sizes)) fd, err := openTxSocket(tx.idx) if err != nil { return nil, fmt.Errorf("%s tx socket: %w", label, err) } d.txFDs = append(d.txFDs, fd) d.txStats = append(d.txStats, &txStats{}) fd, err = openRxSocket(rx.idx, et) if err != nil { return nil, fmt.Errorf("%s rx socket for 0x%04x: %w", label, et, err) } d.rxFDs = append(d.rxFDs, fd) d.rxStats = append(d.rxStats, &rxStats{}) } // Deliberately given no flow rule: a few frames a second does not need a // queue of its own, and the stamps are taken at the wire either way. d.probeSpec = newFrameSpec(rx.mac, tx.mac, cfg.probeEther, []int{probeSize}) fd, err := openTxSocket(tx.idx) if err != nil { return nil, fmt.Errorf("%s probe tx socket: %w", label, err) } if err := enableTxTimestamps(fd); err != nil { return nil, fmt.Errorf("%s probe tx timestamps: %w", label, err) } d.probeTxFD = fd fd, err = openRxSocket(rx.idx, cfg.probeEther) if err != nil { return nil, fmt.Errorf("%s probe rx socket: %w", label, err) } if err := enableRxTimestamps(fd); err != nil { return nil, fmt.Errorf("%s probe rx timestamps: %w", label, err) } d.probeRxFD = fd return d, nil } func (d *direction) start(wg *sync.WaitGroup, doneTx, doneRx *atomic.Bool, cfg config, rxReady *sync.WaitGroup, startTx <-chan struct{}) { for i, fd := range d.txFDs { w := &txWorker{ fd: fd, stream: uint16(i), spec: d.specs[i], batch: cfg.batch, stats: d.txStats[i], startTx: startTx, } wg.Add(1) go func() { defer wg.Done() w.run(doneTx) }() } for i, fd := range d.rxFDs { w := &rxWorker{ fd: fd, batch: cfg.batch, spec: d.specs[i], stats: d.rxStats[i], streams: d.streams, ready: rxReady, } wg.Add(1) go func() { defer wg.Done() w.run(doneRx) }() } sender := &probeSender{fd: d.probeTxFD, spec: d.probeSpec, stats: d.cable} wg.Add(1) go func() { defer wg.Done() sender.run(doneTx, startTx) }() receiver := &probeReceiver{fd: d.probeRxFD, stats: d.cable, ready: rxReady} wg.Add(1) go func() { defer wg.Done() receiver.run(doneRx) }() } func (d *direction) close() { for _, fd := range d.txFDs { unix.Close(fd) } for _, fd := range d.rxFDs { unix.Close(fd) } unix.Close(d.probeTxFD) unix.Close(d.probeRxFD) } type config struct { streams int batch int probeEther uint16 nsPerM float64 } func main() { var ( // The names the kernel gives the only two ports built into it, since as // PID 1 there is no udev to rename them and no command line to pass. aName = flag.String("a", "eth0", "first interface") bName = flag.String("b", "eth1", "second interface") sizesArg = flag.String("sizes", "64,128,256,512,1024,1280,1514", "frame sizes in bytes, excluding FCS, cycled per packet") streams = flag.Int("streams", 7, "independent streams per direction, capped by rx rings; each gets its own ethertype, steered by a flow rule to its own rx queue") batch = flag.Int("batch", 64, "frames per sendmmsg/recvmmsg call") nsPerM = flag.Float64("ns-per-m", 4.85, "mean of both directions, per metre of cable") ) flag.Parse() // Nothing here is recoverable by the time it reaches this point, and as PID 1 // a plain exit would panic the kernel anyway with less to show for it. if err := run(*aName, *bName, *sizesArg, *streams, *batch, *nsPerM); err != nil { panic(err) } } const ( reportInterval = time.Second // Redraw fast so the panel feels live, but measure rates over a much longer // window than a frame, since a frame's worth of a bursty sender is noise. displayInterval = 16 * time.Millisecond // Only long enough to take the edge off one frame's sample; the estimator // does the real smoothing, so this stays small and bounded. rateWindowSpan = 250 * time.Millisecond totalsHold = 50 * time.Millisecond ) func run(aName, bName, sizesArg string, nStreams, batch int, nsPerM float64) error { sizes, err := parseSizes(sizesArg) if err != nil { return err } if err := reportChecks("BOOT", bootstrap()); err != nil { return err } a, err := lookupEndpoint(aName) if err != nil { return err } b, err := lookupEndpoint(bName) if err != nil { return err } for _, e := range []endpoint{a, b} { for _, s := range sizes { if s > e.mtu+ethHdrLen { return fmt.Errorf("size %d exceeds %s MTU %d (max frame %d)", s, e.name, e.mtu, e.mtu+ethHdrLen) } } } a.tag, b.tag = "A", "B" ifnames := []string{a.name, b.name} if nStreams < 1 { return fmt.Errorf("need at least one stream") } ethertypes := make([]uint16, nStreams) for i := range ethertypes { ethertypes[i] = uint16(etherBase + i) } if err := reportChecks("HOST SETTINGS", configureSystem(ifnames, ethertypes)); err != nil { return err } cfg := config{ streams: nStreams, batch: batch, probeEther: uint16(etherBase + nStreams), nsPerM: nsPerM, } var dirs []*direction for _, p := range [][2]endpoint{{a, b}, {b, a}} { d, err := buildDirection(p[0].name+"->"+p[1].name, p[0], p[1], sizes, cfg) if err != nil { return err } dirs = append(dirs, d) } defer func() { for _, d := range dirs { d.close() } }() var linkRows [][]string for _, e := range []endpoint{a, b} { linkRows = append(linkRows, []string{ paint(e.tag, cCyan), e.name, e.macString(), fmt.Sprintf("%.0f Gb/s", e.speed), fmt.Sprintf("%d", e.mtu), }) } fmt.Println(renderBox("LINKS", []string{"TAG", "INTERFACE", "MAC", "SPEED", "MTU"}, []bool{false, false, false, true, true}, linkRows)) target := a.speed if target <= 0 { target = 10 } sizeStrs := make([]string, len(sizes)) for i, s := range sizes { sizeStrs[i] = fmt.Sprintf("%d", s) } fmt.Println(renderBox("CONFIG", []string{"SETTING", "VALUE"}, []bool{false, false}, [][]string{ {"frame sizes", strings.Join(sizeStrs, " ")}, {"streams", fmt.Sprintf("%d per direction, ethertypes 0x%04x-0x%04x", nStreams, ethertypes[0], ethertypes[len(ethertypes)-1])}, {"probe", fmt.Sprintf("ethertype 0x%04x every %s", cfg.probeEther, probeInterval)}, {"batch", fmt.Sprintf("%d frames per syscall", batch)}, {"calibration", fmt.Sprintf("%g ns/m, zero taken from the shortest delay seen so far", nsPerM)}, {"buffers", fmt.Sprintf("sndbuf %s, rcvbuf %s", humanBytes(uint64(sockBufSize(dirs[0].txFDs[0], unix.SO_SNDBUF))), humanBytes(uint64(sockBufSize(dirs[0].rxFDs[0], unix.SO_RCVBUF))))}, })) fmt.Println() var doneTx, doneRx atomic.Bool var wg sync.WaitGroup var rxReady sync.WaitGroup startTx := make(chan struct{}) for _, d := range dirs { rxReady.Add(len(d.rxFDs) + 1) } for _, d := range dirs { d.start(&wg, &doneTx, &doneRx, cfg, &rxReady, startTx) } rxReady.Wait() sig := make(chan os.Signal, 1) signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM) space, restoreTerm := watchSpace() defer restoreTerm() disp, err := newDisplay() if err != nil { return fmt.Errorf("display: %w", err) } defer disp.close() touch, err := watchTouch(disp.fb.pw, disp.fb.ph) if err != nil { return fmt.Errorf("touchscreen: %w", err) } for _, d := range dirs { d.primeCounters() } start := time.Now() close(startTx) tick := time.NewTicker(reportInterval) defer tick.Stop() frame := time.NewTicker(displayInterval) defer frame.Stop() last := time.Now() views := make([]view, len(dirs)) rows := make([]view, len(dirs)) for _, d := range dirs { d.win = newRateWindow(int(rateWindowSpan/displayInterval) + 1) d.est = newRateEstimators() } stats := &streamTable{cols: intervalCols, headerEvery: 20} for { select { case <-sig: doneTx.Store(true) doneRx.Store(true) wg.Wait() return nil case <-space: start = resetAll(dirs, stats) case now := <-frame.C: px, py, down := touch.get() x, y := disp.fb.fromPanel(px, py) if disp.holdReset(x, y, down, now) { start = resetAll(dirs, stats) } for i, d := range dirs { views[i] = d.displayView(now) } cable := "-" if m, ok := cfg.cableMetres(views); ok { cable = fmt.Sprintf("%.1f m", m) } disp.render(totalView(views), now.Sub(start), target*float64(len(dirs)), cable) case now := <-tick.C: secs := now.Sub(last).Seconds() last = now elapsed := now.Sub(start) // Length needs both directions, so every row is sampled before any of // them is printed. for i, d := range dirs { d.sampleNIC() rows[i] = d.view(&d.prevConsole, secs) } length := "-" if m, ok := cfg.cableMetres(rows); ok { length = fmt.Sprintf("%.1f", m) } for i, d := range dirs { for _, line := range stats.emit(d.row(elapsed, rows[i], target, length)) { fmt.Println(line) } } } } }