From e9f62979ffb1543d3c8f02b94055c03d8f1d5771 Mon Sep 17 00:00:00 2001 From: wwqgtxx Date: Wed, 5 Aug 2026 23:32:26 +0800 Subject: [PATCH] chore: improve ZeroTier packet bridge throughput --- adapter/outbound/zerotier.go | 137 ++++++++++++++++++++++------------- go.mod | 2 +- go.sum | 4 +- 3 files changed, 89 insertions(+), 54 deletions(-) diff --git a/adapter/outbound/zerotier.go b/adapter/outbound/zerotier.go index ba5ec64b..a25d5a8c 100644 --- a/adapter/outbound/zerotier.go +++ b/adapter/outbound/zerotier.go @@ -32,7 +32,10 @@ const ( // zeroTierFrameQueueSize absorbs short multi-flow bursts so data frames do // not crowd handshake and window-update traffic out of the single ordered // bridge consumer. - zeroTierFrameQueueSize = 2048 + zeroTierFrameQueueSize = 2048 + // zeroTierFrameBatchSize amortizes runtime validation, lock acquisition, + // and IP-stack delivery while retaining callback order. + zeroTierFrameBatchSize = 64 zeroTierFrameDropLogInterval = 10 * time.Second ) @@ -1049,20 +1052,33 @@ func (z *ZeroTier) resolverForNetworkConfig(config ZT.NetworkConfigData) (resolv } func (z *ZeroTier) runStackPackets(runtime *zeroTierRuntime, device ipStack) { - buffer := make([]byte, 64*1024) - buffers := [][]byte{buffer} - sizes := []int{0} + mtu, err := device.MTU() + if err != nil || mtu < 1 { + mtu = 64 * 1024 + } + batchSize := device.BatchSize() + if batchSize < 1 { + batchSize = 1 + } + storage := make([]byte, mtu*batchSize) + buffers := make([][]byte, batchSize) + for index := range buffers { + start := index * mtu + buffers[index] = storage[start : start+mtu : start+mtu] + } + sizes := make([]int, batchSize) + writeErrors := make([]error, 0, batchSize) for z.ctx.Err() == nil { - if _, err := device.Read(buffers, sizes, 0); err != nil { + count, readErr := device.Read(buffers, sizes, 0) + if readErr != nil { if z.ctx.Err() == nil { - invalidated := z.invalidateDevice(device, fmt.Errorf("ZeroTier stack read failed: %w", err)) - if invalidated && !errors.Is(err, net.ErrClosed) && !errors.Is(err, os.ErrClosed) { - log.Errorln("[ZeroTier](%s) stack read: %v", z.Name(), err) + invalidated := z.invalidateDevice(device, fmt.Errorf("ZeroTier stack read failed: %w", readErr)) + if invalidated && !errors.Is(readErr, net.ErrClosed) && !errors.Is(readErr, os.ErrClosed) { + log.Errorln("[ZeroTier](%s) stack read: %v", z.Name(), readErr) } } return } - packet := buffer[:sizes[0]] z.operationMu.RLock() z.stateMu.RLock() current := z.runtime == runtime && z.tunDevice == device @@ -1071,68 +1087,87 @@ func (z *ZeroTier) runStackPackets(runtime *zeroTierRuntime, device ipStack) { z.operationMu.RUnlock() return } - err := runtime.ipLink.WritePacket(packet) + writeErrors = writeErrors[:0] + for index := 0; index < count; index++ { + if writeErr := runtime.ipLink.WritePacket(buffers[index][:sizes[index]]); writeErr != nil { + writeErrors = append(writeErrors, writeErr) + } + } z.operationMu.RUnlock() - if err != nil { - log.Debugln("[ZeroTier](%s) send IP packet: %v", z.Name(), err) + for _, writeErr := range writeErrors { + log.Debugln("[ZeroTier](%s) send IP packet: %v", z.Name(), writeErr) } } } func (z *ZeroTier) runInboundFrames() { + frames := make([]zeroTierInboundFrame, zeroTierFrameBatchSize) + packets := make([][]byte, 0, zeroTierFrameBatchSize) + frameErrors := make([]error, 0, zeroTierFrameBatchSize) for { + var first zeroTierInboundFrame select { - case inbound := <-z.frameCh: - z.stateMu.RLock() - current := z.runtime == inbound.runtime - z.stateMu.RUnlock() - if !current || inbound.runtime == nil { - continue - } - z.handleInboundFrame(inbound.runtime, inbound.frame) + case first = <-z.frameCh: case <-z.ctx.Done(): return } + frames[0] = first + count := 1 + drain: + for count < len(frames) { + select { + case frames[count] = <-z.frameCh: + count++ + default: + break drain + } + } + device, output, processingErrors := z.processInboundFrames(frames[:count], packets[:0], frameErrors[:0]) + for _, err := range processingErrors { + log.Debugln("[ZeroTier](%s) process inbound frame: %v", z.Name(), err) + } + if device != nil && len(output) != 0 { + if _, writeErr := device.Write(output, 0); writeErr != nil && z.ctx.Err() == nil { + invalidated := z.invalidateDevice(device, fmt.Errorf("ZeroTier stack write failed: %w", writeErr)) + if invalidated && !errors.Is(writeErr, net.ErrClosed) && !errors.Is(writeErr, os.ErrClosed) { + log.Debugln("[ZeroTier](%s) stack write: %v", z.Name(), writeErr) + } + } + } + for index := 0; index < count; index++ { + frames[index] = zeroTierInboundFrame{} + } + packets, frameErrors = output, processingErrors } } -func (z *ZeroTier) handleInboundFrame(runtime *zeroTierRuntime, frame ZT.Frame) { +// processInboundFrames converts one callback-order batch while preventing a +// concurrent configuration update from mutating its IP link. Stack delivery +// follows after the read lock is released. +func (z *ZeroTier) processInboundFrames(frames []zeroTierInboundFrame, packets [][]byte, frameErrors []error) (ipStack, [][]byte, []error) { z.operationMu.RLock() z.stateMu.RLock() - current := z.runtime == runtime - z.stateMu.RUnlock() - if !current { - z.operationMu.RUnlock() - return - } - packet, err := runtime.ipLink.HandleFrame(frame) - if len(packet) == 0 { - z.operationMu.RUnlock() - if err != nil { - log.Debugln("[ZeroTier](%s) process inbound frame: %v", z.Name(), err) - } - return - } - z.stateMu.RLock() - current = z.runtime == runtime + runtime := z.runtime device := z.tunDevice z.stateMu.RUnlock() - z.operationMu.RUnlock() - var writeErr error - var invalidated bool - if current && device != nil { - if _, writeErr = device.Write([][]byte{packet}, 0); writeErr != nil { - if z.ctx.Err() == nil { - invalidated = z.invalidateDevice(device, fmt.Errorf("ZeroTier stack write failed: %w", writeErr)) - } + if runtime == nil { + z.operationMu.RUnlock() + return nil, packets, frameErrors + } + for _, inbound := range frames { + if inbound.runtime != runtime { + continue + } + packet, err := runtime.ipLink.HandleFrame(inbound.frame) + if err != nil { + frameErrors = append(frameErrors, err) + } + if len(packet) != 0 { + packets = append(packets, packet) } } - if err != nil { - log.Debugln("[ZeroTier](%s) process inbound frame: %v", z.Name(), err) - } - if invalidated && !errors.Is(writeErr, net.ErrClosed) && !errors.Is(writeErr, os.ErrClosed) { - log.Debugln("[ZeroTier](%s) stack write: %v", z.Name(), writeErr) - } + z.operationMu.RUnlock() + return device, packets, frameErrors } func (z *ZeroTier) networkStackFor(destination netip.Addr) (*ZTIP.Link, ipStack, error) { diff --git a/go.mod b/go.mod index a6c2cebd..72f4f958 100644 --- a/go.mod +++ b/go.mod @@ -27,7 +27,7 @@ require ( github.com/metacubex/jls-tls v0.0.0-20260723084315-67adc0e2f796 github.com/metacubex/kcp-go v0.0.0-20260105040817-550693377604 github.com/metacubex/mhurl v0.1.0 - github.com/metacubex/mipstack v0.0.0-20260805095554-85edeb894267 + github.com/metacubex/mipstack v0.0.0-20260805152730-93b8f3669159 github.com/metacubex/mlkem v0.1.0 github.com/metacubex/quic-go v0.61.1-0.20260727080200-2548683b76f4 github.com/metacubex/randv2 v0.2.0 diff --git a/go.sum b/go.sum index fbc07e59..5d2bb6ef 100644 --- a/go.sum +++ b/go.sum @@ -134,8 +134,8 @@ github.com/metacubex/kcp-go v0.0.0-20260105040817-550693377604 h1:hJwCVlE3ojViC3 github.com/metacubex/kcp-go v0.0.0-20260105040817-550693377604/go.mod h1:lpmN3m269b3V5jFCWtffqBLS4U3QQoIid9ugtO+OhVc= github.com/metacubex/mhurl v0.1.0 h1:ZdW4Zxe3j3uJ89gNytOazHu6kbHn5owutN/VfXOI8GE= github.com/metacubex/mhurl v0.1.0/go.mod h1:2qpQImCbXoUs6GwJrjuEXKelPyoimsIXr07eNKZdS00= -github.com/metacubex/mipstack v0.0.0-20260805095554-85edeb894267 h1:3gJU5rCj53j5owaynkYUQ1UcDKVvpkKk/7y+5nrLT0A= -github.com/metacubex/mipstack v0.0.0-20260805095554-85edeb894267/go.mod h1:+bbwALZI0pbi2auSG5A3ptdpV2DZ2eLObziXo+P7oj0= +github.com/metacubex/mipstack v0.0.0-20260805152730-93b8f3669159 h1:1BEjP/A0LEm3WX12BLGsWk7jPTuytBrEYBTJmfkiYVU= +github.com/metacubex/mipstack v0.0.0-20260805152730-93b8f3669159/go.mod h1:+bbwALZI0pbi2auSG5A3ptdpV2DZ2eLObziXo+P7oj0= github.com/metacubex/mlkem v0.1.0 h1:wFClitonSFcmipzzQvax75beLQU+D7JuC+VK1RzSL8I= github.com/metacubex/mlkem v0.1.0/go.mod h1:amhaXZVeYNShuy9BILcR7P0gbeo/QLZsnqCdL8U2PDQ= github.com/metacubex/nftables v0.0.0-20260426003805-208c2c1ba2cb h1:wk6mHYPURSUvWcUv72gNP79oiylFsscBSDPJ6ieV6Iw=