diff --git a/main.go b/main.go index 71abd82..e66604d 100644 --- a/main.go +++ b/main.go @@ -99,28 +99,6 @@ func parseSizes(s string) ([]int, error) { return out, nil } -func parseCPUs(s string) ([]int, error) { - if strings.TrimSpace(s) == "" { - return nil, nil - } - var out []int - for _, f := range strings.Split(s, ",") { - v, err := strconv.Atoi(strings.TrimSpace(f)) - if err != nil { - return nil, fmt.Errorf("bad cpu %q: %w", f, err) - } - out = append(out, v) - } - return out, nil -} - -func cpuAt(cpus []int, i int) int { - if len(cpus) == 0 { - return -1 - } - return cpus[i%len(cpus)] -} - func (d *direction) snapshot() sample { var s sample for _, t := range d.txStats { @@ -252,7 +230,6 @@ func (d *direction) start(wg *sync.WaitGroup, doneTx, doneRx *atomic.Bool, cfg c stream: uint16(i), spec: d.specs[i], batch: cfg.batch, - cpu: cpuAt(cfg.txCPUs, i), stats: d.txStats[i], startTx: startTx, } @@ -266,7 +243,6 @@ func (d *direction) start(wg *sync.WaitGroup, doneTx, doneRx *atomic.Bool, cfg c w := &rxWorker{ fd: fd, batch: cfg.batch, - cpu: cpuAt(cfg.rxCPUs, i), spec: d.specs[i], stats: d.rxStats[i], streams: d.streams, @@ -293,8 +269,6 @@ func (d *direction) close() { type config struct { streams int batch int - txCPUs []int - rxCPUs []int } func main() { @@ -306,12 +280,10 @@ func main() { 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") - txCPUs = flag.String("txcpus", "", "comma-separated CPUs to pin tx workers to") - rxCPUs = flag.String("rxcpus", "", "comma-separated CPUs to pin rx workers to") ) flag.Parse() - if err := run(*aName, *bName, *sizesArg, *patArg, *txCPUs, *rxCPUs, + if err := run(*aName, *bName, *sizesArg, *patArg, *streams, *batch, *duplex); err != nil { fmt.Fprintln(os.Stderr, "error:", err) os.Exit(1) @@ -320,7 +292,7 @@ func main() { const reportInterval = 500 * time.Millisecond -func run(aName, bName, sizesArg, patArg, txCPUsArg, rxCPUsArg string, +func run(aName, bName, sizesArg, patArg string, nStreams, batch int, duplex bool) error { if aName == "" || bName == "" { @@ -334,14 +306,6 @@ func run(aName, bName, sizesArg, patArg, txCPUsArg, rxCPUsArg string, if err != nil { return err } - txCPUs, err := parseCPUs(txCPUsArg) - if err != nil { - return err - } - rxCPUs, err := parseCPUs(rxCPUsArg) - if err != nil { - return err - } a, err := lookupEndpoint(aName) if err != nil { return err @@ -387,8 +351,6 @@ func run(aName, bName, sizesArg, patArg, txCPUsArg, rxCPUsArg string, cfg := config{ streams: nStreams, batch: batch, - txCPUs: txCPUs, - rxCPUs: rxCPUs, } var dirs []*direction diff --git a/rx.go b/rx.go index 33511e6..c4f52fe 100644 --- a/rx.go +++ b/rx.go @@ -21,7 +21,6 @@ type rxStats struct { type rxWorker struct { fd int batch int - cpu int spec *frameSpec stats *rxStats streams []lossWindow @@ -30,8 +29,6 @@ type rxWorker struct { } func (w *rxWorker) run(done *atomic.Bool) { - pinTo(w.cpu) - bufs := make([][]byte, w.batch) for i := range bufs { bufs[i] = make([]byte, maxFrame) diff --git a/sock.go b/sock.go index aee6ab1..0122f3b 100644 --- a/sock.go +++ b/sock.go @@ -2,7 +2,6 @@ package main import ( "fmt" - "runtime" "unsafe" "golang.org/x/sys/unix" @@ -115,17 +114,6 @@ func newMmsghdrs(bufs [][]byte) ([]mmsghdr, []unix.Iovec) { return hdrs, iovs } -func pinTo(cpu int) error { - if cpu < 0 { - return nil - } - runtime.LockOSThread() - var set unix.CPUSet - set.Zero() - set.Set(cpu) - return unix.SchedSetaffinity(0, &set) -} - func packetDrops(fd int) uint64 { st, err := unix.GetsockoptTpacketStats(fd, unix.SOL_PACKET, unix.PACKET_STATISTICS) if err != nil { diff --git a/tx.go b/tx.go index d2cff31..4622a1b 100644 --- a/tx.go +++ b/tx.go @@ -19,14 +19,11 @@ type txWorker struct { stream uint16 spec *frameSpec batch int - cpu int stats *txStats startTx <-chan struct{} } func (w *txWorker) run(done *atomic.Bool) { - pinTo(w.cpu) - bufs := make([][]byte, w.batch) for i := range bufs { bufs[i] = make([]byte, w.spec.maxSize)