From adf71d14585b3fc3e7c9ddcb9cdadea40aa7aa8e Mon Sep 17 00:00:00 2001 From: JackDoan Date: Tue, 14 Jul 2026 10:56:21 -0500 Subject: [PATCH] simplify making new Queues --- control_lifecycle_test.go | 17 +++++++------ inside.go | 4 +-- interface.go | 50 +++++++++++++++++++------------------ overlay/device.go | 10 +++++--- overlay/overlaytest/noop.go | 13 ++-------- overlay/tun_android.go | 12 ++------- overlay/tun_darwin.go | 12 ++------- overlay/tun_disabled.go | 25 ++++++------------- overlay/tun_freebsd.go | 12 ++------- overlay/tun_ios.go | 12 ++------- overlay/tun_linux.go | 24 +++++++++++------- overlay/tun_linux_test.go | 17 ++++++------- overlay/tun_netbsd.go | 12 ++------- overlay/tun_openbsd.go | 12 ++------- overlay/tun_tester.go | 12 ++------- overlay/tun_windows.go | 12 ++------- overlay/user.go | 19 +++----------- overlay/user_test.go | 26 +++++++------------ 18 files changed, 105 insertions(+), 196 deletions(-) diff --git a/control_lifecycle_test.go b/control_lifecycle_test.go index 0890baeb..3e942cb6 100644 --- a/control_lifecycle_test.go +++ b/control_lifecycle_test.go @@ -51,12 +51,8 @@ func (d *fakeDevice) Activate() error { return nil } func (d *fakeDevice) Networks() []netip.Prefix { return nil } func (d *fakeDevice) Name() string { return "fake" } func (d *fakeDevice) RoutesFor(netip.Addr) routing.Gateways { return nil } -func (d *fakeDevice) SupportsMultiqueue() bool { return false } -func (d *fakeDevice) NewMultiQueueReader() error { - return errors.New("unsupported") -} -func (d *fakeDevice) Readers() []tio.Queue { return []tio.Queue{d} } +func (d *fakeDevice) Queues(int) ([]tio.Queue, error) { return []tio.Queue{d}, nil } // newReadyControl hand-builds the minimum Control that Main would have // produced right before Start, including the construction token NewInterface @@ -82,7 +78,6 @@ func newReadyControl(t *testing.T) (*Control, *fakeDevice, *fakeConn) { inside: dev, outside: conn, writers: []udp.Conn{conn}, - readers: make([]tio.Queue, 1), batchers: make([]batch.RxBatcher, 1), routines: 1, hostMap: newHostMap(l), @@ -164,7 +159,14 @@ type multiqueueDevice struct { *fakeDevice } -func (d *multiqueueDevice) SupportsMultiqueue() bool { return true } +// Queues claims multiqueue support but fails to open the second queue, +// exercising the activation error path. +func (d *multiqueueDevice) Queues(n int) ([]tio.Queue, error) { + if n > 1 { + return nil, errors.New("second queue failed to open") + } + return d.fakeDevice.Queues(n) +} func TestControl_StartMultiqueueFailureReleases(t *testing.T) { dev := &multiqueueDevice{fakeDevice: newFakeDevice()} @@ -175,7 +177,6 @@ func TestControl_StartMultiqueueFailureReleases(t *testing.T) { inside: dev, outside: conn, writers: []udp.Conn{conn}, - readers: make([]tio.Queue, 2), batchers: make([]batch.RxBatcher, 2), routines: 2, l: test.NewLogger(), diff --git a/inside.go b/inside.go index ed331ebd..acf19ba4 100644 --- a/inside.go +++ b/inside.go @@ -57,7 +57,7 @@ func (f *Interface) consumeInsidePacket(pkt tio.Packet, fwPacket *firewall.Packe // kernel as one giant blob; segment first so the loopback // path sees one IP datagram per Write. err := tio.SegmentSuperpacket(pkt, func(seg []byte) error { - _, werr := f.readers[q].Write(seg) + _, werr := f.queues[q].Write(seg) return werr }) if err != nil { @@ -273,7 +273,7 @@ func (f *Interface) rejectInside(packet []byte, out []byte, q int) { return } - _, err := f.readers[q].Write(out) + _, err := f.queues[q].Write(out) if err != nil { f.l.Error("Failed to write to tun", "error", err) } diff --git a/interface.go b/interface.go index 256849b3..ab7c19df 100644 --- a/interface.go +++ b/interface.go @@ -119,8 +119,8 @@ type Interface struct { ctx context.Context writers []udp.Conn - readers []tio.Queue - // batchers is one per tun queue, wrapping readers[i]. + queues []tio.Queue + // batchers is one per tun queue, wrapping queues[i]. // decryptToTun sends plaintext into the batch.RxBatcher; // listenOut calls its Flush at the end of each UDP recvmmsg batch. batchers []batch.RxBatcher @@ -222,7 +222,6 @@ func NewInterface(ctx context.Context, c *InterfaceConfig) (*Interface, error) { routines: c.routines, version: c.version, writers: make([]udp.Conn, c.routines), - readers: make([]tio.Queue, c.routines), batchers: make([]batch.RxBatcher, c.routines), myVpnNetworks: cs.myVpnNetworks, myVpnNetworksTable: cs.myVpnNetworksTable, @@ -276,36 +275,39 @@ func (f *Interface) activate() error { "boringcrypto", boringEnabled(), ) - if f.routines > 1 { - if !f.inside.SupportsMultiqueue() || !f.outside.SupportsMultipleReaders() { - f.routines = 1 - f.l.Warn("routines is not supported on this platform, falling back to a single routine") - } + if f.routines > 1 && !f.outside.SupportsMultipleReaders() { + f.routines = 1 + f.l.Warn("multiple udp readers are not supported on this platform, falling back to a single routine") } + // Prepare the tun queues. A device that can't open that many hands back + // fewer (a single queue on platforms without multiqueue support) and we + // size the reader routines to what we actually got. + queues, err := f.inside.Queues(f.routines) + if err != nil { + return err + } + if len(queues) < f.routines { + f.l.Warn("tun multiqueue is not supported on this platform, falling back to fewer routines", + "requested", f.routines, "opened", len(queues)) + f.routines = len(queues) + } + f.queues = queues + metrics.GetOrRegisterGauge("routines", nil).Update(int64(f.routines)) - // Prepare n tun queues - for i := 0; i < f.routines; i++ { - if i > 0 { - if err = f.inside.NewMultiQueueReader(); err != nil { - return err - } - } - } - f.readers = f.inside.Readers() - for i := range f.readers { - caps := tio.QueueCapabilities(f.readers[i]) + for i := range f.queues { + caps := tio.QueueCapabilities(f.queues[i]) if caps.TSO || caps.USO { // Multi-lane: TCP gets coalesced when TSO is on, UDP when USO // is on, everything else (and either lane disabled) falls // through to passthrough so non-IP / non-TCP-UDP traffic still // reaches the TUN. arena := batch.NewArena(batch.DefaultMultiArenaCap) - f.batchers[i] = batch.NewMultiCoalescer(f.readers[i], f.l, arena, caps.TSO, caps.USO) + f.batchers[i] = batch.NewMultiCoalescer(f.queues[i], f.l, arena, caps.TSO, caps.USO) } else { arena := batch.NewArena(batch.DefaultPassthroughArenaCap) - f.batchers[i] = batch.NewPassthrough(f.readers[i], arena) + f.batchers[i] = batch.NewPassthrough(f.queues[i], arena) } } @@ -329,7 +331,7 @@ func (f *Interface) run() { // Launch n queues to read packets from tun dev for i := 0; i < f.routines; i++ { f.wg.Go(func() { - f.listenIn(f.readers[i], i) + f.listenIn(f.queues[i], i) }) } @@ -392,7 +394,7 @@ func (f *Interface) listenOut(i int) { f.l.Debug("underlay reader is done", "reader", i) } -func (f *Interface) listenIn(reader tio.Queue, i int) { +func (f *Interface) listenIn(queue tio.Queue, i int) { // Pinning this thread (and goroutine) to a single CPU keeps every sendmmsg from this goroutine going through the // same TX ring on the nic, so the wire sees per-flow order. Skip entirely when tun.pin_threads is false. if f.pinThreads { @@ -423,7 +425,7 @@ func (f *Interface) listenIn(reader tio.Queue, i int) { conntrackCache := firewall.NewConntrackCacheTicker(f.ctx, f.l, f.conntrackCacheTimeout) for { - pkts, err := reader.Read() + pkts, err := queue.Read() if err != nil { // Same shutdown noise handling as listenOut if !f.closed.Load() && f.ctx.Err() == nil { diff --git a/overlay/device.go b/overlay/device.go index 8044ee75..eca35c16 100644 --- a/overlay/device.go +++ b/overlay/device.go @@ -18,7 +18,11 @@ type Device interface { Networks() []netip.Prefix Name() string RoutesFor(netip.Addr) routing.Gateways - SupportsMultiqueue() bool - NewMultiQueueReader() error - Readers() []tio.Queue + // Queues returns the device's packet queues, opening additional ones as + // needed until there are n. Platforms without multiqueue support return + // their single queue regardless of n, so callers must size reader loops + // to len(result), not n; implementations never return more than n. An + // error means a queue that should have opened could not; the caller owns + // cleanup via Close. Called once, during interface activation. + Queues(n int) ([]tio.Queue, error) } diff --git a/overlay/overlaytest/noop.go b/overlay/overlaytest/noop.go index 6a39ab43..0268c9ec 100644 --- a/overlay/overlaytest/noop.go +++ b/overlay/overlaytest/noop.go @@ -3,7 +3,6 @@ package overlaytest import ( - "errors" "net/netip" "github.com/slackhq/nebula/overlay/tio" @@ -39,16 +38,8 @@ func (NoopTun) Write([]byte) (int, error) { return 0, nil } -func (NoopTun) SupportsMultiqueue() bool { - return false -} - -func (NoopTun) NewMultiQueueReader() error { - return errors.New("unsupported") -} - -func (NoopTun) Readers() []tio.Queue { - return []tio.Queue{NoopTun{}} +func (NoopTun) Queues(int) ([]tio.Queue, error) { + return []tio.Queue{NoopTun{}}, nil } func (NoopTun) Close() error { diff --git a/overlay/tun_android.go b/overlay/tun_android.go index c1ca1220..f7ab417a 100644 --- a/overlay/tun_android.go +++ b/overlay/tun_android.go @@ -97,14 +97,6 @@ func (t *tun) Name() string { return "android" } -func (t *tun) SupportsMultiqueue() bool { - return false -} - -func (t *tun) NewMultiQueueReader() error { - return fmt.Errorf("TODO: multiqueue not implemented for android") -} - -func (t *tun) Readers() []tio.Queue { - return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)} +func (t *tun) Queues(int) ([]tio.Queue, error) { + return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)}, nil } diff --git a/overlay/tun_darwin.go b/overlay/tun_darwin.go index b966330e..368182a9 100644 --- a/overlay/tun_darwin.go +++ b/overlay/tun_darwin.go @@ -606,14 +606,6 @@ func (t *tun) Name() string { return t.Device } -func (t *tun) SupportsMultiqueue() bool { - return false -} - -func (t *tun) NewMultiQueueReader() error { - return fmt.Errorf("TODO: multiqueue not implemented for darwin") -} - -func (t *tun) Readers() []tio.Queue { - return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)} +func (t *tun) Queues(int) ([]tio.Queue, error) { + return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)}, nil } diff --git a/overlay/tun_disabled.go b/overlay/tun_disabled.go index 870b318e..82204ad6 100644 --- a/overlay/tun_disabled.go +++ b/overlay/tun_disabled.go @@ -19,10 +19,9 @@ type disabledTun struct { vpnNetworks []netip.Prefix // Track these metrics since we don't have the tun device to do it for us - tx metrics.Counter - rx metrics.Counter - l *slog.Logger - numReaders int + tx metrics.Counter + rx metrics.Counter + l *slog.Logger } // Read hands the next queued packet to a reader, copying it into b. Reads @@ -47,7 +46,6 @@ func newDisabledTun(vpnNetworks []netip.Prefix, queueLen int, metricsEnabled boo vpnNetworks: vpnNetworks, read: make(chan []byte, queueLen), l: l, - numReaders: 1, } if metricsEnabled { @@ -108,23 +106,14 @@ func (t *disabledTun) Write(b []byte) (int, error) { return len(b), nil } -func (t *disabledTun) SupportsMultiqueue() bool { - return true -} - -func (t *disabledTun) NewMultiQueueReader() error { - t.numReaders++ - return nil -} - -func (t *disabledTun) Readers() []tio.Queue { - out := make([]tio.Queue, t.numReaders) - for i := range t.numReaders { +func (t *disabledTun) Queues(n int) ([]tio.Queue, error) { + out := make([]tio.Queue, n) + for i := range out { // NoClose: the shared channel and metrics are owned by the // disabledTun; Close on the device tears them down once for everybody. out[i] = tio.NewSingleQueueNoClose(t, defaultBatchBufSize) } - return out + return out, nil } func (t *disabledTun) Close() error { diff --git a/overlay/tun_freebsd.go b/overlay/tun_freebsd.go index 914d70f8..e6479001 100644 --- a/overlay/tun_freebsd.go +++ b/overlay/tun_freebsd.go @@ -560,12 +560,8 @@ func (t *tun) Name() string { return t.Device } -func (t *tun) SupportsMultiqueue() bool { - return false -} - -func (t *tun) NewMultiQueueReader() error { - return fmt.Errorf("TODO: multiqueue not implemented for freebsd") +func (t *tun) Queues(int) ([]tio.Queue, error) { + return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)}, nil } func (t *tun) addRoutes(logErrors bool) error { @@ -592,10 +588,6 @@ func (t *tun) addRoutes(logErrors bool) error { return nil } -func (t *tun) Readers() []tio.Queue { - return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)} -} - func (t *tun) removeRoutes(routes []Route) error { for _, r := range routes { if !r.Install { diff --git a/overlay/tun_ios.go b/overlay/tun_ios.go index be3b539b..56603b02 100644 --- a/overlay/tun_ios.go +++ b/overlay/tun_ios.go @@ -160,14 +160,6 @@ func (t *tun) Name() string { return "iOS" } -func (t *tun) SupportsMultiqueue() bool { - return false -} - -func (t *tun) NewMultiQueueReader() error { - return fmt.Errorf("TODO: multiqueue not implemented for ios") -} - -func (t *tun) Readers() []tio.Queue { - return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)} +func (t *tun) Queues(int) ([]tio.Queue, error) { + return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)}, nil } diff --git a/overlay/tun_linux.go b/overlay/tun_linux.go index 1f643f88..386db52f 100644 --- a/overlay/tun_linux.go +++ b/overlay/tun_linux.go @@ -39,7 +39,7 @@ type tun struct { // the kernel: usoOffloadFlags when USO was accepted, tsoOffloadFlags on // the TSO-only fallback, or 0 when vnetHdr is off. TUNSETOFFLOAD is // device-wide (drivers/net/tun.c set_offload updates tun->set_features - // for the whole netdev), so NewMultiQueueReader must replay this exact + // for the whole netdev), so addQueue must replay this exact // mask on every added queue — issuing a narrower mask there would // silently downgrade offloads (e.g. disable USO) for all queues while // they still advertise the stale capability. @@ -174,7 +174,7 @@ func newTun(c *config.C, l *slog.Logger, vpnNetworks []netip.Prefix, multiqueue } vnetHdr := true // offloadFlags is the exact TUN_F_* mask the kernel accepted. We remember - // it (rather than a plain bool) so NewMultiQueueReader can replay the + // it (rather than a plain bool) so addQueue can replay the // identical device-wide mask on added queues instead of downgrading them. var offloadFlags uint name, err := tunSetIff(fd, nameStr, baseFlags|unix.IFF_VNET_HDR) @@ -349,11 +349,21 @@ func (t *tun) reload(c *config.C, initial bool) error { return nil } -func (t *tun) SupportsMultiqueue() bool { - return true +// Queues opens additional kernel multiqueue fds until the device has n +// queues, then returns them all. The first queue was opened by newTun; each +// extra fd replays the negotiated offload state (see addQueue). +func (t *tun) Queues(n int) ([]tio.Queue, error) { + for len(t.readers.Queues()) < n { + if err := t.addQueue(); err != nil { + return nil, err + } + } + return t.readers.Queues(), nil } -func (t *tun) NewMultiQueueReader() error { +// addQueue opens one more IFF_MULTI_QUEUE fd on the device and adds it to +// the queue set. +func (t *tun) addQueue() error { t.closeLock.Lock() defer t.closeLock.Unlock() @@ -837,10 +847,6 @@ func (t *tun) updateRoutes(r netlink.RouteUpdate) { t.routeTree.Store(newTree) } -func (t *tun) Readers() []tio.Queue { - return t.readers.Queues() -} - func (t *tun) Close() error { t.closeLock.Lock() defer t.closeLock.Unlock() diff --git a/overlay/tun_linux_test.go b/overlay/tun_linux_test.go index 401d5a07..2144c780 100644 --- a/overlay/tun_linux_test.go +++ b/overlay/tun_linux_test.go @@ -41,7 +41,7 @@ func TestTunAdvMSS(t *testing.T) { func TestOffloadUSOEnabled(t *testing.T) { // usoOffloadFlags must be a strict superset of tsoOffloadFlags. Otherwise // the TSO-only fallback (and the historic hardcoded-mask bug in - // NewMultiQueueReader) would not actually be a downgrade. + // addQueue) would not actually be a downgrade. if usoOffloadFlags&tsoOffloadFlags != tsoOffloadFlags { t.Fatalf("usoOffloadFlags (%#x) is not a superset of tsoOffloadFlags (%#x)", usoOffloadFlags, tsoOffloadFlags) } @@ -67,20 +67,19 @@ func TestOffloadUSOEnabled(t *testing.T) { } } -// TestNewMultiQueueReaderReplaysNegotiatedMask guards the device-wide -// TUNSETOFFLOAD downgrade bug: NewMultiQueueReader must issue the exact mask -// newTun negotiated (t.offloadFlags), not a hardcoded TSO-only mask. Because -// TUNSETOFFLOAD is per-netdev, a narrower mask on an added queue silently -// disables USO for every queue on a USO-capable kernel while the queues keep -// advertising it. +// TestAddQueueReplaysNegotiatedMask guards the device-wide TUNSETOFFLOAD +// downgrade bug: addQueue must issue the exact mask newTun negotiated +// (t.offloadFlags), not a hardcoded TSO-only mask. Because TUNSETOFFLOAD is +// per-netdev, a narrower mask on an added queue silently disables USO for +// every queue on a USO-capable kernel while the queues keep advertising it. // // A full multi-queue exercise needs /dev/net/tun and CAP_NET_ADMIN, which are // not available in CI/sandbox, so this asserts on the struct field that the // TUNSETOFFLOAD argument is read from. -func TestNewMultiQueueReaderReplaysNegotiatedMask(t *testing.T) { +func TestAddQueueReplaysNegotiatedMask(t *testing.T) { t.Run("uso-negotiated", func(t *testing.T) { tn := &tun{vnetHdr: true, offloadFlags: usoOffloadFlags} - // The ioctl argument in NewMultiQueueReader is uintptr(t.offloadFlags); + // The ioctl argument in addQueue is uintptr(t.offloadFlags); // it must equal the negotiated USO mask, and must NOT be the TSO-only // mask (the original bug). if tn.offloadFlags != usoOffloadFlags { diff --git a/overlay/tun_netbsd.go b/overlay/tun_netbsd.go index 9904a88a..97691543 100644 --- a/overlay/tun_netbsd.go +++ b/overlay/tun_netbsd.go @@ -68,10 +68,6 @@ type tun struct { fd int } -func (t *tun) Readers() []tio.Queue { - return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)} -} - var deviceNameRE = regexp.MustCompile(`^tun[0-9]+$`) func newTunFromFd(_ *config.C, _ *slog.Logger, _ int, _ []netip.Prefix) (*tun, error) { @@ -394,12 +390,8 @@ func (t *tun) Name() string { return t.Device } -func (t *tun) SupportsMultiqueue() bool { - return false -} - -func (t *tun) NewMultiQueueReader() error { - return fmt.Errorf("TODO: multiqueue not implemented for netbsd") +func (t *tun) Queues(int) ([]tio.Queue, error) { + return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)}, nil } func (t *tun) addRoutes(logErrors bool) error { diff --git a/overlay/tun_openbsd.go b/overlay/tun_openbsd.go index c4b1a302..23816b0d 100644 --- a/overlay/tun_openbsd.go +++ b/overlay/tun_openbsd.go @@ -369,12 +369,8 @@ func (t *tun) Name() string { return t.Device } -func (t *tun) SupportsMultiqueue() bool { - return false -} - -func (t *tun) NewMultiQueueReader() error { - return fmt.Errorf("TODO: multiqueue not implemented for openbsd") +func (t *tun) Queues(int) ([]tio.Queue, error) { + return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)}, nil } func (t *tun) addRoutes(logErrors bool) error { @@ -425,10 +421,6 @@ func (t *tun) deviceBytes() (o [16]byte) { return } -func (t *tun) Readers() []tio.Queue { - return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)} -} - func addRoute(prefix netip.Prefix, gateways []netip.Prefix) error { sock, err := unix.Socket(unix.AF_ROUTE, unix.SOCK_RAW, unix.AF_UNSPEC) if err != nil { diff --git a/overlay/tun_tester.go b/overlay/tun_tester.go index 34911527..4b2685e0 100644 --- a/overlay/tun_tester.go +++ b/overlay/tun_tester.go @@ -178,14 +178,6 @@ func (t *TestTun) Read(b []byte) (int, error) { return n, nil } -func (t *TestTun) Readers() []tio.Queue { - return []tio.Queue{tio.NewSingleQueue(t, udp.MTU)} -} - -func (t *TestTun) SupportsMultiqueue() bool { - return false -} - -func (t *TestTun) NewMultiQueueReader() error { - return fmt.Errorf("TODO: multiqueue not implemented") +func (t *TestTun) Queues(int) ([]tio.Queue, error) { + return []tio.Queue{tio.NewSingleQueue(t, udp.MTU)}, nil } diff --git a/overlay/tun_windows.go b/overlay/tun_windows.go index a5ae036e..6be85ffc 100644 --- a/overlay/tun_windows.go +++ b/overlay/tun_windows.go @@ -263,16 +263,8 @@ func (t *winTun) Write(b []byte) (int, error) { return t.tun.Write(b, 0) } -func (t *winTun) SupportsMultiqueue() bool { - return false -} - -func (t *winTun) NewMultiQueueReader() error { - return fmt.Errorf("TODO: multiqueue not implemented for windows") -} - -func (t *winTun) Readers() []tio.Queue { - return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)} +func (t *winTun) Queues(int) ([]tio.Queue, error) { + return []tio.Queue{tio.NewSingleQueue(t, defaultBatchBufSize)}, nil } func (t *winTun) Close() error { diff --git a/overlay/user.go b/overlay/user.go index 7693f653..2d775bde 100644 --- a/overlay/user.go +++ b/overlay/user.go @@ -24,13 +24,11 @@ func NewUserDevice(vpnNetworks []netip.Prefix) (Device, error) { outboundWriter: ow, inboundReader: ir, inboundWriter: iw, - numReaders: 1, }, nil } type UserDevice struct { vpnNetworks []netip.Prefix - numReaders int outboundReader *io.PipeReader outboundWriter *io.PipeWriter @@ -49,25 +47,16 @@ func (d *UserDevice) RoutesFor(ip netip.Addr) routing.Gateways { return routing.Gateways{routing.NewGateway(ip, 1)} } -func (d *UserDevice) SupportsMultiqueue() bool { - return true -} - -func (d *UserDevice) NewMultiQueueReader() error { - d.numReaders++ - return nil -} - -func (d *UserDevice) Readers() []tio.Queue { - out := make([]tio.Queue, d.numReaders) - for i := range d.numReaders { +func (d *UserDevice) Queues(n int) ([]tio.Queue, error) { + out := make([]tio.Queue, n) + for i := range out { // All queues share the underlying pipes (the io.Pipe serializes // concurrent callers) but each owns a private scratch buffer so // concurrent Reads across queues never alias. NoClose: the pipes are // owned by the UserDevice and torn down once by UserDevice.Close. out[i] = tio.NewSingleQueueNoClose(d, defaultBatchBufSize) } - return out + return out, nil } func (d *UserDevice) Pipe() (*io.PipeReader, *io.PipeWriter) { diff --git a/overlay/user_test.go b/overlay/user_test.go index f0904bb1..9e0e9c9c 100644 --- a/overlay/user_test.go +++ b/overlay/user_test.go @@ -24,29 +24,21 @@ func newTestUserDevice(t *testing.T) *UserDevice { return ud } -// TestUserDeviceReadersDistinctBuffers is the regression test for the -// multiqueue packet-corruption bug: Readers() used to hand the same -// *UserDevice (and therefore the same read scratch buffer) to every queue, so -// one reader's borrowed Packet.Bytes was overwritten by another reader's -// concurrent Read. Readers() must now return numReaders DISTINCT queue -// objects, each with its own backing buffer — verified behaviorally below by -// holding one queue's borrowed slice across the other queue's Read. +// TestUserDeviceReadersDistinctBuffers ensures each Queue is actually different func TestUserDeviceReadersDistinctBuffers(t *testing.T) { d := newTestUserDevice(t) - // One extra reader => two queues total. - if err := d.NewMultiQueueReader(); err != nil { - t.Fatalf("NewMultiQueueReader: %v", err) + readers, err := d.Queues(2) + if err != nil { + t.Fatalf("Queues: %v", err) } - - readers := d.Readers() if len(readers) != 2 { - t.Fatalf("Readers() returned %d queues, want 2", len(readers)) + t.Fatalf("Queues(2) returned %d queues, want 2", len(readers)) } // Distinct queue objects. if readers[0] == readers[1] { - t.Fatal("Readers() returned the same queue object twice") + t.Fatal("Queues(2) returned the same queue object twice") } // Drive one packet through each queue and confirm the borrowed bytes from @@ -99,10 +91,10 @@ func TestUserDeviceReadersDistinctBuffers(t *testing.T) { // and corrupted each other's returned slices. func TestUserDeviceReadersConcurrentRace(t *testing.T) { d := newTestUserDevice(t) - if err := d.NewMultiQueueReader(); err != nil { - t.Fatalf("NewMultiQueueReader: %v", err) + readers, err := d.Queues(2) + if err != nil { + t.Fatalf("Queues: %v", err) } - readers := d.Readers() _, ow := d.Pipe() const iterations = 200