mirror of
https://github.com/slackhq/nebula.git
synced 2026-08-15 08:36:57 +02:00
913a37cfee
Device loses io.ReadWriteCloser + NewMultiQueueReader in favor of Queues(n), which returns up to n tio.Queue objects; platforms without multiqueue hand back their single queue and the interface sizes its reader routines to what it actually got. Queue.Read returns a batch of borrowed packets (single-element for every current backend) so a future backend can deliver more than one packet per syscall without another interface change. The Linux poll/eventfd machinery moves out of tun_linux.go into the new overlay/tio package: nonblocking fds, a shared shutdown eventfd owned by the queue set, and pollfd arrays built on the stack so concurrent writers parked in blockOnWrite no longer share Revents storage. Other platforms wrap their existing one-datagram Read/Write in a singleQueue adapter that owns a private scratch buffer, so multiqueue-by-sharing devices (user, disabled) no longer race concurrent readers on one buffer. This is the tun-interface subset of better-tun-interface-ordering, extracted at 18dc13b with none of the GSO/GRO offload mechanics and no udp/sendmmsg changes. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
141 lines
3.1 KiB
Go
141 lines
3.1 KiB
Go
package overlay
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"log/slog"
|
|
"net/netip"
|
|
"strings"
|
|
|
|
"github.com/rcrowley/go-metrics"
|
|
"github.com/slackhq/nebula/iputil"
|
|
"github.com/slackhq/nebula/overlay/tio"
|
|
"github.com/slackhq/nebula/routing"
|
|
)
|
|
|
|
type disabledTun struct {
|
|
read chan []byte
|
|
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
|
|
}
|
|
|
|
// Read hands the next queued packet to a reader, copying it into b. Reads
|
|
// from concurrent queues are safe: the channel receive serializes them and
|
|
// each queue copies into its own private scratch buffer.
|
|
func (t *disabledTun) Read(b []byte) (int, error) {
|
|
r, ok := <-t.read
|
|
if !ok {
|
|
return 0, io.EOF
|
|
}
|
|
|
|
t.tx.Inc(1)
|
|
if t.l.Enabled(context.Background(), slog.LevelDebug) {
|
|
t.l.Debug("Write payload", "raw", prettyPacket(r))
|
|
}
|
|
|
|
return copy(b, r), nil
|
|
}
|
|
|
|
func newDisabledTun(vpnNetworks []netip.Prefix, queueLen int, metricsEnabled bool, l *slog.Logger) *disabledTun {
|
|
tun := &disabledTun{
|
|
vpnNetworks: vpnNetworks,
|
|
read: make(chan []byte, queueLen),
|
|
l: l,
|
|
}
|
|
|
|
if metricsEnabled {
|
|
tun.tx = metrics.GetOrRegisterCounter("messages.tx.message", nil)
|
|
tun.rx = metrics.GetOrRegisterCounter("messages.rx.message", nil)
|
|
} else {
|
|
tun.tx = &metrics.NilCounter{}
|
|
tun.rx = &metrics.NilCounter{}
|
|
}
|
|
|
|
return tun
|
|
}
|
|
|
|
func (*disabledTun) Activate() error {
|
|
return nil
|
|
}
|
|
|
|
func (*disabledTun) RoutesFor(addr netip.Addr) routing.Gateways {
|
|
return routing.Gateways{}
|
|
}
|
|
|
|
func (t *disabledTun) Networks() []netip.Prefix {
|
|
return t.vpnNetworks
|
|
}
|
|
|
|
func (*disabledTun) Name() string {
|
|
return "disabled"
|
|
}
|
|
|
|
func (t *disabledTun) handleICMPEchoRequest(b []byte) bool {
|
|
out := make([]byte, len(b))
|
|
out = iputil.CreateICMPEchoResponse(b, out)
|
|
if out == nil {
|
|
return false
|
|
}
|
|
|
|
// attempt to write it, but don't block
|
|
select {
|
|
case t.read <- out:
|
|
default:
|
|
t.l.Debug("tun_disabled: dropped ICMP Echo Reply response")
|
|
}
|
|
|
|
return true
|
|
}
|
|
|
|
func (t *disabledTun) Write(b []byte) (int, error) {
|
|
t.rx.Inc(1)
|
|
|
|
// Check for ICMP Echo Request before spending time doing the full parsing
|
|
if t.handleICMPEchoRequest(b) {
|
|
if t.l.Enabled(context.Background(), slog.LevelDebug) {
|
|
t.l.Debug("Disabled tun responded to ICMP Echo Request", "raw", prettyPacket(b))
|
|
}
|
|
} else if t.l.Enabled(context.Background(), slog.LevelDebug) {
|
|
t.l.Debug("Disabled tun received unexpected payload", "raw", prettyPacket(b))
|
|
}
|
|
return len(b), nil
|
|
}
|
|
|
|
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, nil
|
|
}
|
|
|
|
func (t *disabledTun) Close() error {
|
|
if t.read != nil {
|
|
close(t.read)
|
|
t.read = nil
|
|
}
|
|
return nil
|
|
}
|
|
|
|
type prettyPacket []byte
|
|
|
|
func (p prettyPacket) String() string {
|
|
var s strings.Builder
|
|
|
|
for i, b := range p {
|
|
if i > 0 && i%8 == 0 {
|
|
s.WriteString(" ")
|
|
}
|
|
s.WriteString(fmt.Sprintf("%02x ", b))
|
|
}
|
|
|
|
return s.String()
|
|
}
|