package main import ( "flag" "fmt" "net" "os" "os/signal" "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 } 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 { specs []*frameSpec txStats []*txStats rxStats []*rxStats streams []lossWindow txFDs []int rxFDs []int probeSpec *frameSpec probeTxFD int probeRxFD int statFD int cable *cableStats // Guards everything the sampler touches. The counters are read on their own // clock and drawn on another, and the two must not read them at once: // sampleDrops consumes what it reads, so a second caller would see a gap. mu sync.Mutex win *rateWindow drops uint64 base counterSet // The completed receive bucket the display draws its rate from, refreshed by // the sampler because the buckets are keyed by the mac's clock and staleness // has to be judged against the wall. rateFrames uint64 rateBytes uint64 epoch int64 epochAt time.Time nic atomic.Uint64 poller *nicPoller } // Everything the display reads, taken at one instant, so a pair of these // describes both the rates and the errors over the span between them. type counterSet struct { t time.Time s sample drops uint64 nic uint64 } func (d *direction) capture() counterSet { d.sampleDrops() s := d.snapshot() return counterSet{t: time.Now(), s: s, drops: d.drops, nic: d.nic.Load()} } // What someone testing a cable is asking, rather than how each failure happened // to be noticed. type errs struct { lost uint64 corrupt uint64 link uint64 internal uint64 } func (e errs) total() uint64 { return e.lost + e.corrupt + e.link + e.internal } func (e errs) add(o errs) errs { return errs{ lost: e.lost + o.lost, corrupt: e.corrupt + o.corrupt, link: e.link + o.link, internal: e.internal + o.internal, } } // A ring of one bucket per drawn frame, spanning rateWindowSpan. Rates come // from the gap between adjacent buckets and errors from the ends of the ring, // so both slide forward every frame instead of stepping once a second. type rateWindow struct { buf []counterSet idx int filled bool } func newRateWindow(n int) *rateWindow { return &rateWindow{buf: make([]counterSet, n)} } func (w *rateWindow) push(c counterSet) { w.buf[w.idx] = c w.idx++ if w.idx == len(w.buf) { w.idx = 0 w.filled = true } } func (w *rateWindow) count() int { if w.filled { return len(w.buf) } return w.idx } // Indexed oldest first, so a partly filled ring reads the same as a full one. func (w *rateWindow) at(i int) counterSet { if w.filled { i += w.idx } return w.buf[i%len(w.buf)] } // How long the newest epoch may sit still before the wire is taken to have gone // quiet. A stamp only advances when a frame arrives, so a frozen epoch means no // traffic rather than an unchanged rate. const rateStale = 100 * time.Millisecond // The newest epoch is still filling, since frames received in it may not have // been drained yet, so the rate is read from the one before it. func (d *direction) readRateBucket(now time.Time) { var newest int64 for _, r := range d.rxStats { if e := r.newest.Load(); e > newest { newest = e } } if newest > d.epoch { d.epoch, d.epochAt = newest, now } d.rateFrames, d.rateBytes = 0, 0 if d.epoch == 0 || now.Sub(d.epochAt) > rateStale { return } for _, r := range d.rxStats { f, b := r.bucket(d.epoch - 1) d.rateFrames += f d.rateBytes += b } } type sample struct { rxFrames, rxBytes uint64 lost, late uint64 crcErr, badMagic uint64 badHdr uint64 badLen uint64 txErrs 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) return endpoint{name: name, idx: ifi.Index, mac: mac, mtu: ifi.MTU}, nil } func (d *direction) snapshot() sample { var s sample for _, t := range d.txStats { s.txErrs += t.errs.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.badHdr += r.badHdr.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 and the rolling error window are about now // rather than since the reset, so they keep running; the origin goes into the // ring so the newest bucket never sits behind it. func (d *direction) reset() { d.mu.Lock() d.base = d.capture() d.win.push(d.base) d.mu.Unlock() 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{ {group: "NOW", title: "bits/s", width: 9, right: true}, {group: "NOW", title: "packets/s", width: 9, right: true}, {group: "NOW", title: "lost", width: 7, right: true}, {group: "NOW", title: "corrupt", width: 7, right: true}, {group: "NOW", title: "link", width: 7, right: true}, {group: "NOW", title: "internal", width: 8, right: true}, {group: "OVERALL", title: "elapsed", width: 9, right: true}, {group: "OVERALL", title: "packets", width: 9, right: true}, {group: "OVERALL", title: "bytes", width: 9, right: true}, {group: "OVERALL", title: "metres", width: 6, right: true}, {group: "OVERALL", title: "lost", width: 9, right: true}, {group: "OVERALL", title: "corrupt", width: 9, right: true}, {group: "OVERALL", title: "link", width: 9, right: true}, {group: "OVERALL", title: "internal", width: 9, right: true}, } // Shared by the console table and the framebuffer so both show the same // figures. type view struct { rxPPS float64 rxGbps float64 rxFrames, rxBytes uint64 since errs window errs cable cableView } func errsBetween(b, n counterSet) errs { return errs{ lost: n.s.lost - b.s.lost, // Four ways of noticing one thing: a payload that does not match its // checksum, a header that does not match its own, a header that is not // ours, and a length that cannot be. corrupt: (n.s.crcErr - b.s.crcErr) + (n.s.badHdr - b.s.badHdr) + (n.s.badMagic - b.s.badMagic) + (n.s.badLen - b.s.badLen), // What the hardware reported. Nothing the host declined to send is here, // so this one going red means the cable. link: (n.nic - b.nic) + (n.s.rxErrs - b.s.rxErrs), // Ours rather than the cable's. A late frame is unreachable while each // stream has a flow rule to its own queue, which is exactly why it is // worth counting. internal: (n.drops - b.drops) + (n.s.late - b.s.late) + (n.s.txErrs - b.s.txErrs), } } func (d *direction) counters(now counterSet) view { return view{ rxFrames: now.s.rxFrames - d.base.s.rxFrames, rxBytes: now.s.rxBytes - d.base.s.rxBytes, cable: d.cable.view(), since: errsBetween(d.base, now), } } func totalView(views []view) view { var t view for _, v := range views { t.rxPPS += v.rxPPS t.rxGbps += v.rxGbps t.rxFrames += v.rxFrames t.rxBytes += v.rxBytes t.since = t.since.add(v.since) t.window = t.window.add(v.window) } return t } func (d *direction) sample() { d.mu.Lock() d.win.push(d.capture()) d.readRateBucket(time.Now()) d.mu.Unlock() } // Draws what the sampler last put in the ring rather than reading the counters // again, so the display never participates in the measurement. func (d *direction) displayView() view { d.mu.Lock() n := d.win.count() if n == 0 { d.mu.Unlock() return view{cable: d.cable.view()} } v := d.counters(d.win.at(n - 1)) if n >= 2 { v.window = errsBetween(d.win.at(0), d.win.at(n-1)) } v.rxPPS = float64(d.rateFrames) / rateBucketSecs v.rxGbps = gbps(d.rateBytes, d.rateFrames, rateBucketSecs) d.mu.Unlock() return v } // The same figures the panel draws, in the same order: the last second as rates // and error flags, then everything since the reset. func totalRow(elapsed time.Duration, v view, target float64, length string) []string { return []string{ rateCell(v.rxGbps*1e9, target*1e9), scaleSI(v.rxPPS), flagCell(v.window.lost), flagCell(v.window.corrupt), flagCell(v.window.link), flagCell(v.window.internal), scaleTime(elapsed), scaleCount(v.rxFrames), scaleCount(v.rxBytes), length, statusCell(v.since.lost), statusCell(v.since.corrupt), statusCell(v.since.link), statusCell(v.since.internal), } } // 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.poller.prime() d.reset() } func buildDirection(label string, tx, rx endpoint) (*direction, error) { // Built before the windows, since each window judges sequence numbers against // the frontier its own sender publishes. txs := make([]*txStats, numStreams) for i := range txs { txs[i] = &txStats{} } d := &direction{ txStats: txs, streams: newLossWindows(txs), cable: newCableStats(), } // Held open for the life of the run: the stats ioctl is issued five times a // second and reopening a socket for each one is pure overhead. statFD, err := unix.Socket(unix.AF_INET, unix.SOCK_DGRAM, 0) if err != nil { return nil, fmt.Errorf("%s stats socket: %w", label, err) } d.statFD = statFD d.poller, err = newNICPoller(statFD, tx.name, rx.name, &d.nic) if err != nil { return nil, fmt.Errorf("%s: %w", label, err) } d.win = newRateWindow(int(rateWindowSpan/sampleInterval) + 1) for i := 0; i < numStreams; i++ { et := uint16(etherBase + i) d.specs = append(d.specs, newFrameSpec(rx.mac, tx.mac, et, frameSizes)) fd, err := openTxSocket(tx.idx) if err != nil { return nil, fmt.Errorf("%s tx socket: %w", label, err) } d.txFDs = append(d.txFDs, fd) fd, err = openRxSocket(rx.idx, et) if err != nil { return nil, fmt.Errorf("%s rx socket for 0x%04x: %w", label, et, err) } // The mac already stamps every frame for the probe's sake, so this only // asks for the stamp to be delivered. if err := enableRxTimestamps(fd); err != nil { return nil, fmt.Errorf("%s rx timestamps 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, 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, 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, done *atomic.Bool, rxReady *sync.WaitGroup, startTx <-chan struct{}) { for i, fd := range d.txFDs { w := &txWorker{ fd: fd, stream: uint16(i), spec: d.specs[i], batch: batchSize, stats: d.txStats[i], startTx: startTx, } wg.Add(1) go func() { defer wg.Done() w.run(done) }() } for i, fd := range d.rxFDs { w := &rxWorker{ fd: fd, batch: batchSize, stream: uint16(i), spec: d.specs[i], stats: d.rxStats[i], loss: &d.streams[i], ready: rxReady, } wg.Add(1) go func() { defer wg.Done() w.run(done) }() } sender := &probeSender{fd: d.probeTxFD, spec: d.probeSpec, stats: d.cable} wg.Add(1) go func() { defer wg.Done() sender.run(done, startTx) }() receiver := &probeReceiver{fd: d.probeRxFD, stats: d.cable, ready: rxReady} wg.Add(1) go func() { defer wg.Done() receiver.run(done) }() wg.Add(1) go func() { defer wg.Done() d.poller.run(done, startTx) }() } 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) unix.Close(d.statFD) } const ( numStreams = 7 batchSize = 64 probeEther uint16 = etherBase + numStreams // A constant rather than the negotiated speed, since this has to come up // with no cable in the port and nothing to negotiate. linkSpeed = 10.0 ) // The mac appends the fcs, so 60 and 1514 here are the smallest and largest // standard frames, 64 and 1518 on the wire. var frameSizes = []int{60, 128, 256, 512, 1024, 1280, 1514} func main() { // 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") 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, *nsPerM); err != nil { panic(err) } } const ( reportInterval = time.Second // Deliberately not tied to the refresh: letting a slow or blocked draw set // the sampling clock would stretch the window it reports. sampleInterval = 16 * time.Millisecond // How far back the shown errors reach. The rate is not taken from this ring // but from the receive buckets, which are keyed by the mac's clock. rateWindowSpan = time.Second ) // One sampler for both directions, so they are read back to back on one clock // rather than drifting apart on two. type sampler struct { dirs []*direction } func (s *sampler) run(done *atomic.Bool, startTx <-chan struct{}) { <-startTx tick := time.NewTicker(sampleInterval) defer tick.Stop() for !done.Load() { <-tick.C for _, d := range s.dirs { d.sample() } } } func run(aName, bName string, nsPerM float64) error { 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 frameSizes { 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} ethertypes := make([]uint16, numStreams) for i := range ethertypes { ethertypes[i] = uint16(etherBase + i) } if err := reportChecks("HOST SETTINGS", configureSystem(ifnames, ethertypes)); err != nil { return err } 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]) 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("%d", e.mtu), }) } fmt.Println(renderBox("LINKS", []string{"TAG", "INTERFACE", "MAC", "MTU"}, []bool{false, false, false, true}, linkRows)) fmt.Println() // One row carries both directions, so line rate is both links at once. target := linkSpeed * float64(len(dirs)) var done 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, &done, &rxReady, startTx) } samp := &sampler{dirs: dirs} wg.Add(1) go func() { defer wg.Done() samp.run(&done, 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() views := make([]view, len(dirs)) rows := make([]view, len(dirs)) stats := &streamTable{cols: intervalCols, headerEvery: 20} for { select { case <-sig: done.Store(true) wg.Wait() return nil case <-space: start = resetAll(dirs, stats) case <-disp.fb.flips: now := time.Now() 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() } // Empty until the probe has a stamp from each direction, so the // panel shows nothing there rather than a placeholder. cable := "" if m, ok := cableMetres(views, nsPerM); ok { cable = fmt.Sprintf("%.1f", m) } if err := disp.render(totalView(views), now.Sub(start), cable); err != nil { return err } case now := <-tick.C: elapsed := now.Sub(start) // Length needs both directions, so every row is sampled before any of // them is printed. for i, d := range dirs { rows[i] = d.displayView() } length := "-" if m, ok := cableMetres(rows, nsPerM); ok { length = fmt.Sprintf("%.1f", m) } for _, line := range stats.emit(totalRow(elapsed, totalView(rows), target, length)) { fmt.Println(line) } } } }