Skip to content

Commit 468a862

Browse files
committed
tunnel: remove OutboundChan from vpn.Device, use callback instead
this drastically reduces scheduler overhead on high bandwidth networks
1 parent e022d76 commit 468a862

3 files changed

Lines changed: 58 additions & 58 deletions

File tree

application.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -120,6 +120,7 @@ func (a *Application) Init(ctx context.Context, tunDevice tun.Device) error {
120120
a.Dns = NewDNSService(a.Conf, a.Eventbus, a.ctx, a.logger)
121121
a.AuthStatus = service.NewAuthStatus(a.P2p, a.Conf, a.Eventbus)
122122
a.Tunnel = service.NewTunnel(a.P2p, vpnDevice, a.Conf)
123+
go vpnDevice.ReadTUNPackets(a.Tunnel.HandleReadPackets)
123124
a.SOCKS5, err = service.NewSOCKS5(a.P2p, a.Conf)
124125
if err != nil {
125126
return fmt.Errorf("failed to init socks5: %v", err)

service/tunnel.go

Lines changed: 37 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -24,26 +24,32 @@ const (
2424
)
2525

2626
type Tunnel struct {
27-
p2p P2p
28-
conf *config.Config
29-
device *vpn.Device
30-
logger *log.ZapEventLogger
31-
peersLock sync.RWMutex
32-
peerIDToPeer map[peer.ID]*VpnPeer
33-
netIPToPeer map[string]*VpnPeer
27+
p2p P2p
28+
conf *config.Config
29+
device *vpn.Device
30+
logger *log.ZapEventLogger
31+
32+
isClosed atomic.Bool
33+
peersLock sync.RWMutex
34+
peerIDToPeer map[peer.ID]*VpnPeer
35+
netIPToPeer map[string]*VpnPeer
36+
udpBroadcastAddr net.IP
3437
}
3538

3639
func NewTunnel(p2pService P2p, device *vpn.Device, conf *config.Config) *Tunnel {
40+
localIP, netMask := conf.VPNLocalIPMask()
41+
udpBroadcastAddr := vpn.GetIPv4BroadcastAddress(&net.IPNet{IP: localIP, Mask: netMask})
42+
3743
tunnel := &Tunnel{
38-
p2p: p2pService,
39-
conf: conf,
40-
device: device,
41-
logger: log.Logger("awl/service/tunnel"),
42-
peerIDToPeer: make(map[peer.ID]*VpnPeer),
43-
netIPToPeer: make(map[string]*VpnPeer),
44+
p2p: p2pService,
45+
conf: conf,
46+
device: device,
47+
logger: log.Logger("awl/service/tunnel"),
48+
peerIDToPeer: make(map[peer.ID]*VpnPeer),
49+
netIPToPeer: make(map[string]*VpnPeer),
50+
udpBroadcastAddr: udpBroadcastAddr,
4451
}
4552
tunnel.RefreshPeersList()
46-
go tunnel.backgroundReadPackets()
4753

4854
return tunnel
4955
}
@@ -160,6 +166,8 @@ func (t *Tunnel) Close() {
160166
t.peersLock.Lock()
161167
defer t.peersLock.Unlock()
162168

169+
t.isClosed.Store(true)
170+
163171
for _, vpnPeer := range t.peerIDToPeer {
164172
localIP := *vpnPeer.localIP.Load()
165173
vpnPeer.Close(t)
@@ -168,17 +176,24 @@ func (t *Tunnel) Close() {
168176
}
169177
}
170178

171-
func (t *Tunnel) backgroundReadPackets() {
172-
localIP, netMask := t.conf.VPNLocalIPMask()
173-
broadcastAddr := vpn.GetIPv4BroadcastAddress(&net.IPNet{IP: localIP, Mask: netMask})
179+
// HandleReadPackets for successfully handled packets it sets packet in slice as nil
180+
func (t *Tunnel) HandleReadPackets(packets []*vpn.Packet) {
181+
t.peersLock.RLock()
182+
defer t.peersLock.RUnlock()
183+
184+
if t.isClosed.Load() {
185+
return
186+
}
187+
188+
for i, packet := range packets {
189+
if packet == nil {
190+
continue
191+
}
174192

175-
// TODO: batch read
176-
for packet := range t.device.OutboundChan() {
177193
// TODO: ipv6 support
178-
if packet.Dst.Equal(broadcastAddr) || packet.Dst.Equal(net.IPv4bcast) {
194+
if packet.Dst.Equal(t.udpBroadcastAddr) || packet.Dst.Equal(net.IPv4bcast) {
179195
// udp broadcast
180196

181-
t.peersLock.RLock()
182197
for _, vpnPeer := range t.netIPToPeer {
183198
// TODO: replace with event-based check OnConnected/OnDisconnected to improve performance
184199
if !t.p2p.IsConnected(vpnPeer.peerID) {
@@ -195,26 +210,19 @@ func (t *Tunnel) backgroundReadPackets() {
195210
}
196211
}
197212

198-
t.device.PutTempPacket(packet)
199-
t.peersLock.RUnlock()
200-
201213
continue
202214
}
203215

204-
t.peersLock.RLock()
205216
vpnPeer, ok := t.netIPToPeer[string(packet.Dst)]
206217
if !ok {
207-
t.device.PutTempPacket(packet)
208-
t.peersLock.RUnlock()
209218
continue
210219
}
211220

212221
select {
213222
case vpnPeer.outboundCh <- packet:
223+
packets[i] = nil
214224
default:
215-
t.device.PutTempPacket(packet)
216225
}
217-
t.peersLock.RUnlock()
218226
}
219227
}
220228

vpn/vpn.go

Lines changed: 20 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -9,24 +9,22 @@ import (
99
"sync/atomic"
1010

1111
"github.com/ipfs/go-log/v2"
12+
"go.uber.org/zap"
1213
"golang.zx2c4.com/wireguard/tun"
1314
)
1415

1516
const (
1617
InterfaceMTU = 3500
17-
maxContentSize = InterfaceMTU * 2 // TODO: determine real size
18-
outboundChCap = 50
18+
maxContentSize = InterfaceMTU + 100
1919
// internal tun header. see offset in tun_darwin (4) and tun_linux (virtioNetHdrLen, currently 10)
2020
tunPacketOffset = 14
2121
)
2222

2323
type Device struct {
24-
tun tun.Device
25-
mtu int64
26-
localIP net.IP
27-
outboundCh chan *Packet
24+
tun tun.Device
25+
mtu int64
26+
localIP net.IP
2827

29-
closeCh chan struct{}
3028
packetsPool sync.Pool
3129
logger *log.ZapEventLogger
3230
}
@@ -49,19 +47,16 @@ func NewDevice(existingTun tun.Device, interfaceName string, localIP net.IP, ipM
4947
}
5048

5149
dev := &Device{
52-
tun: tunDevice,
53-
mtu: int64(realMtu),
54-
localIP: localIP,
55-
outboundCh: make(chan *Packet, outboundChCap),
50+
tun: tunDevice,
51+
mtu: int64(realMtu),
52+
localIP: localIP,
5653
packetsPool: sync.Pool{
5754
New: func() interface{} {
5855
return new(Packet)
5956
}},
60-
logger: log.Logger("awl/vpn"),
61-
closeCh: make(chan struct{}),
57+
logger: log.Logger("awl/vpn"),
6258
}
6359
go dev.tunEventsReader()
64-
go dev.tunPacketsReader()
6560

6661
return dev, nil
6762
}
@@ -97,12 +92,7 @@ func (d *Device) WritePacket(data *Packet, senderIP net.IP) error {
9792
return nil
9893
}
9994

100-
func (d *Device) OutboundChan() <-chan *Packet {
101-
return d.outboundCh
102-
}
103-
10495
func (d *Device) Close() error {
105-
close(d.closeCh)
10696
return d.tun.Close()
10797
}
10898

@@ -137,9 +127,7 @@ func (d *Device) tunEventsReader() {
137127
}
138128
}
139129

140-
func (d *Device) tunPacketsReader() {
141-
defer close(d.outboundCh)
142-
130+
func (d *Device) ReadTUNPackets(packetsHandler func([]*Packet)) {
143131
batchSize := d.tun.BatchSize()
144132
packets := make([]*Packet, batchSize)
145133
bufs := make([][]byte, batchSize)
@@ -167,16 +155,19 @@ func (d *Device) tunPacketsReader() {
167155
data.Packet = data.Buffer[tunPacketOffset : size+tunPacketOffset]
168156
okay := data.Parse()
169157
if !okay {
158+
d.logger.Error("Failed to parse packet",
159+
zap.ByteString("packet", data.Packet[:min(10, len(data.Packet))]),
160+
)
161+
162+
packets[i] = nil
163+
d.PutTempPacket(data)
170164
continue
171165
}
166+
}
172167

173-
select {
174-
case <-d.closeCh:
175-
return
176-
case d.outboundCh <- data:
177-
// ok
178-
}
179-
packets[i] = nil
168+
if packetsCount > 0 {
169+
// packetsHandler skips nil packets in slice and sets packet to nil after successful processing
170+
packetsHandler(packets[:packetsCount])
180171
}
181172

182173
if errors.Is(err, tun.ErrTooManySegments) {

0 commit comments

Comments
 (0)