Files
cabletest/main.go
T

783 lines
20 KiB
Go
Raw Normal View History

2026-07-25 17:24:29 -07:00
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
nicNow uint64
nicBase 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
}
func (e *rateEstimator) update(x float64) {
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
}
// 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
}
2026-07-25 19:29:00 -07:00
// 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() {
2026-07-25 19:29:00 -07:00
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()
2026-07-25 19:29:00 -07:00
}
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
txFrames, txSent uint64
rxFrames, rxGot uint64
lost, late uint64
crc, badMagic uint64
kdrop, link uint64
errors uint64
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 {
2026-07-25 19:29:00 -07:00
b := d.errBase
v := view{
txFrames: now.txFrames - b.txFrames,
txSent: now.txBytes - b.txBytes,
rxFrames: now.rxFrames - b.rxFrames,
rxGot: now.rxBytes - b.rxBytes,
lost: now.lost - b.lost,
late: now.late - b.late,
crc: now.crcErr - b.crcErr,
badMagic: now.badMagic - b.badMagic,
kdrop: d.drops - d.dropBase,
cable: d.cable.view(),
}
// 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.
v.link = d.nicNow - d.nicBase +
(now.txErrs - b.txErrs) + (now.rxErrs - b.rxErrs)
v.errors = v.lost + v.crc + v.badMagic + (now.badLen - b.badLen) + v.kdrop + v.link
return v
}
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.txFrames = d.heldFrames.get(t, v.txFrames)
v.txSent = d.heldSent.get(t, v.txSent)
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
d.est.txPPS.update(float64(txF) / secs)
d.est.rxPPS.update(float64(rxF) / secs)
d.est.txGbps.update(gbps(n.txBytes-o.txBytes, txF, secs))
d.est.rxGbps.update(gbps(n.rxBytes-o.rxBytes, rxF, secs))
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.lost),
statusCell(v.late),
statusCell(v.crc),
statusCell(v.badMagic),
statusCell(v.kdrop),
statusCell(v.link),
statusCell(v.errors),
paint(v.cable.minText(), cCyan),
paint(length, cCyan),
}
}
// Read once a second rather than per frame, since these are sysfs files; the
// display uses whatever the last sample left behind.
func (d *direction) sampleNIC() {
tx := readNIC(d.tx.name)
rx := readNIC(d.rx.name)
d.nicNow = tx.tx + rx.rx + tx.carrierDown
}
func buildDirection(label string, tx, rx endpoint, patIdx int, 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(patIdx, 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{})
}
// The probe carries its own ethertype so it lands on a socket of its own, but
// it is left unsteered: it is a few frames a second and does not need a queue
// to itself, and the stamps are taken at the wire either way.
d.probeSpec = newFrameSpec(patIdx, 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
// Whatever the interfaces have counted before now is not ours.
d.sampleNIC()
d.nicBase = d.nicNow
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
zeroNS float64
nsPerM float64
}
2026-07-25 17:24:29 -07:00
func main() {
var (
aName = flag.String("a", "", "first interface")
bName = flag.String("b", "", "second interface")
sizesArg = flag.String("sizes", "64,128,256,512,1024,1280,1514", "frame sizes in bytes, excluding FCS, cycled per packet")
patArg = flag.String("pattern", "prbs", "payload pattern")
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")
duplex = flag.Bool("duplex", true, "run both directions simultaneously")
zeroNS = flag.Float64("zero-ns", 4327.5, "both directions summed at zero cable length; belongs to the media adapters, recalibrate when they change")
nsPerM = flag.Float64("ns-per-m", 10.909, "both directions summed, per metre of cable")
)
flag.Parse()
if err := run(*aName, *bName, *sizesArg, *patArg,
*streams, *batch, *duplex, *zeroNS, *nsPerM); err != nil {
fmt.Fprintln(os.Stderr, "error:", err)
os.Exit(1)
}
}
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, patArg string,
nStreams, batch int, duplex bool, zeroNS, nsPerM float64) error {
if aName == "" || bName == "" {
return fmt.Errorf("both -a and -b are required")
}
sizes, err := parseSizes(sizesArg)
if err != nil {
return err
}
patIdx, err := patternIndex(patArg)
if 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)
}
var fatal []string
var tuneRows [][]string
for _, r := range configureSystem(ifnames, ethertypes) {
tuneRows = append(tuneRows, []string{r.item, r.status(), r.detail()})
if r.fatal {
fatal = append(fatal, r.item)
}
}
fmt.Println(renderBox("HOST SETTINGS",
[]string{"CHECK", "STATUS", "DETAIL"},
[]bool{false, false, false}, tuneRows))
if len(fatal) > 0 {
return fmt.Errorf("cannot test with %s in this state", strings.Join(fatal, ", "))
}
cfg := config{
streams: nStreams,
batch: batch,
probeEther: uint16(etherBase + nStreams),
zeroNS: zeroNS,
nsPerM: nsPerM,
}
var dirs []*direction
d0, err := buildDirection(a.name+"->"+b.name, a, b, patIdx, sizes, cfg)
if err != nil {
return err
}
dirs = append(dirs, d0)
if duplex {
d1, err := buildDirection(b.name+"->"+a.name, b, a, patIdx, sizes, cfg)
if err != nil {
return err
}
dirs = append(dirs, d1)
}
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("TEST",
[]string{"SETTING", "VALUE"},
[]bool{false, false}, [][]string{
{"pattern", patterns[patIdx].name},
{"frame sizes", strings.Join(sizeStrs, " ")},
{"streams", fmt.Sprintf("%d per direction, ethertypes 0x%04x-0x%04x, one rx queue each",
nStreams, ethertypes[0], ethertypes[len(ethertypes)-1])},
{"batch", fmt.Sprintf("%d frames per syscall", batch)},
{"payload verify", "crc32c on every frame"},
{"cable probe", fmt.Sprintf("%d-byte frame on ethertype 0x%04x every %s, hardware stamped at both macs",
probeSize, cfg.probeEther, probeInterval)},
{"duplex", fmt.Sprintf("%v", duplex)},
{"socket buffers", fmt.Sprintf("sndbuf %s, rcvbuf %s (granted)",
humanBytes(uint64(sockBufSize(dirs[0].txFDs[0], unix.SO_SNDBUF))),
humanBytes(uint64(sockBufSize(dirs[0].rxFDs[0], unix.SO_RCVBUF))))},
}))
2026-07-25 19:29:00 -07:00
fmt.Println(paint("rates are per interval; error counts are cumulative, press space to reset them", cDim))
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)
2026-07-25 19:29:00 -07:00
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.w, disp.fb.h)
if err != nil {
return fmt.Errorf("touchscreen: %w", err)
}
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
2026-07-25 19:29:00 -07:00
case <-space:
start = resetAll(dirs, stats)
case now := <-frame.C:
if x, y, down := touch.get(); disp.holdReset(x, y, down, now) {
start = resetAll(dirs, stats)
}
for i, d := range dirs {
views[i] = d.displayView(now)
}
disp.render(dirs, views, now.Sub(start), target, cfg.cableText(views))
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)
}
}
}
}
2026-07-25 17:24:29 -07:00
}