mirror of
https://github.com/tladesignz/dnstt.git
synced 2026-10-01 19:38:02 +03:00
This enlarges a few buffers and windows, with the goal of improving download performance. kcp's SetWindowSize controls the number of unacknowledged packets that are allowed. smux's MaxStreamBuffer is another kind of "receive window" that is advertised to the peer of how much we are willing to receive at once. The default MaxStreamBuffer is 64 KB, but kcptun overrides the default to 2 MB. turbotunnel's QueueSize is the size of internal buffers in QueuePacketConn and RemoteMap; empirically I found that the server would sometimes fill its outgoing buffer if SetWindowSize and QueueSize were equal, so I set QueueSize to be twice SetWindowSize. https://lists.torproject.org/pipermail/anti-censorship-team/2021-July/000178.html https://gitlab.torproject.org/tpo/anti-censorship/pluggable-transports/snowflake/-/merge_requests/48 The changes have a large effect on a direct -udp connection without a recursive resolver—which, however, is a discouraged configuration. Through a recursive resolver, the improvements are more modest. If I really crank up the buffer sizes, I can get surprisingly fast downloads over a direct -udp connection (over 1 MB/s), but a connection through a resolver doesn't keep getting faster and may even get slower. I want to avoid a bufferbloat situation with oversized buffers, too. I manually explored a small neighborhood of parameter values and picked some settings that looked reasonable. The tables below show the test results. The test is downloading 10 MiB between two servers with 100 ms RTT between them. Server: dnstt-server -udp :53 -privkey-file server.key t.example.com 127.0.0.1:9321 ncat -l -k -v 9321 --send-only --sh-exec 'dd bs=1M count=10 if=/dev/urandom' Client: dnstt-client -pubkey-file server.pub t.example.com 127.0.0.1:7000 ncat --recv-only 127.0.0.1 7000 | pv -t -r -a -b -i 0.2 > /dev/null I did the download under every treatment twice and recorded the download rate in KiB/s. "Server drops" comes from hacking some log messages to turbotunnel.QueuePacketConn to track how often the "Drop the incoming packet" (QueueIncoming method) and "Drop the outgoing packet" (WriteTo) cases happen. resolver method QueueSize MaxStreamBuffer SetWindowSize KiB/s KiB/s -------- ------ --------- --------------- ------------- ----- ----- direct udp 64 64*1024 (32, 32) 169 173 (status before this commit) dns.google udp 64 64*1024 (32, 32) 63.8 64.3 (status before this commit) dns.google doh 64 64*1024 (32, 32) 125 122 (status before this commit) resolver method QueueSize MaxStreamBuffer SetWindowSize KiB/s KiB/s -------- ------ --------- --------------- ------------- ----- ----- direct udp 64 1*1024*1024 (32, 32) 172 174 dns.google udp 64 1*1024*1024 (32, 32) 57.3 58.4 server drops dns.google doh 64 1*1024*1024 (32, 32) 128 128 resolver method QueueSize MaxStreamBuffer SetWindowSize KiB/s KiB/s -------- ------ --------- --------------- ------------- ----- ----- direct udp 64 1*1024*1024 (64, 64) 322 305 dns.google udp 64 1*1024*1024 (64, 64) 72.5 70.9 server drops dns.google doh 64 1*1024*1024 (64, 64) 136 139 server drops resolver method QueueSize MaxStreamBuffer SetWindowSize KiB/s KiB/s -------- ------ --------- --------------- ------------- ----- ----- direct udp 128 1*1024*1024 (64, 64) 321 325 (this commit) dns.google udp 128 1*1024*1024 (64, 64) 82.5 78.5 (this commit) dns.google doh 128 1*1024*1024 (64, 64) 129 131 (this commit) resolver method QueueSize MaxStreamBuffer SetWindowSize KiB/s KiB/s -------- ------ --------- --------------- ------------- ----- ----- direct udp 2048 4*1024*1024 (1024, 1024) 1240 1060 server drops dns.google udp 2048 4*1024*1024 (1024, 1024) 73.5 81.4 dns.google doh 2048 4*1024*1024 (1024, 1024) 115 129
163 lines
5.8 KiB
Go
163 lines
5.8 KiB
Go
package turbotunnel
|
|
|
|
import (
|
|
"net"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
)
|
|
|
|
// taggedPacket is a combination of a []byte and a net.Addr, encapsulating the
|
|
// return type of PacketConn.ReadFrom.
|
|
type taggedPacket struct {
|
|
P []byte
|
|
Addr net.Addr
|
|
}
|
|
|
|
// QueuePacketConn implements net.PacketConn by storing queues of packets. There
|
|
// is one incoming queue (where packets are additionally tagged by the source
|
|
// address of the peer that sent them). There are many outgoing queues, one for
|
|
// each remote peer address that has been recently seen. The QueueIncoming
|
|
// method inserts a packet into the incoming queue, to eventually be returned by
|
|
// ReadFrom. WriteTo inserts a packet into an address-specific outgoing queue,
|
|
// which can later by accessed through the OutgoingQueue method.
|
|
//
|
|
// Besides the outgoing queues, there is also a one-element "stash" for each
|
|
// remote peer address. You can stash a packet using the Stash method, and get
|
|
// it back later by receiving from the channel returned by Unstash. The stash is
|
|
// meant as a convenient place to temporarily store a single packet, such as
|
|
// when you've read one too many packets from the send queue and need to store
|
|
// the extra packet to be processed first in the next pass. It's the caller's
|
|
// responsibility to Unstash what they have Stashed. Calling Stash does not put
|
|
// the packet at the head of the send queue; if there is the possibility that a
|
|
// packet has been stashed, it must be checked for by calling Unstash in
|
|
// addition to OutgoingQueue.
|
|
type QueuePacketConn struct {
|
|
remotes *RemoteMap
|
|
localAddr net.Addr
|
|
recvQueue chan taggedPacket
|
|
closeOnce sync.Once
|
|
closed chan struct{}
|
|
// What error to return when the QueuePacketConn is closed.
|
|
err atomic.Value
|
|
}
|
|
|
|
// NewQueuePacketConn makes a new QueuePacketConn, set to track recent peers
|
|
// for at least a duration of timeout.
|
|
func NewQueuePacketConn(localAddr net.Addr, timeout time.Duration) *QueuePacketConn {
|
|
return &QueuePacketConn{
|
|
remotes: NewRemoteMap(timeout),
|
|
localAddr: localAddr,
|
|
recvQueue: make(chan taggedPacket, QueueSize),
|
|
closed: make(chan struct{}),
|
|
}
|
|
}
|
|
|
|
// QueueIncoming queues and incoming packet and its source address, to be
|
|
// returned in a future call to ReadFrom.
|
|
func (c *QueuePacketConn) QueueIncoming(p []byte, addr net.Addr) {
|
|
select {
|
|
case <-c.closed:
|
|
// If we're closed, silently drop it.
|
|
return
|
|
default:
|
|
}
|
|
// Copy the slice so that the caller may reuse it.
|
|
buf := make([]byte, len(p))
|
|
copy(buf, p)
|
|
select {
|
|
case c.recvQueue <- taggedPacket{buf, addr}:
|
|
default:
|
|
// Drop the incoming packet if the receive queue is full.
|
|
}
|
|
}
|
|
|
|
// OutgoingQueue returns the queue of outgoing packets corresponding to addr,
|
|
// creating it if necessary. The contents of the queue will be packets that are
|
|
// written to the address in question using WriteTo.
|
|
func (c *QueuePacketConn) OutgoingQueue(addr net.Addr) <-chan []byte {
|
|
return c.remotes.SendQueue(addr)
|
|
}
|
|
|
|
// Stash places p in the stash for addr, if the stash is not already occupied.
|
|
// Returns true if the packet was placed in the stash, or false if the stash was
|
|
// already occupied. This method is similar to WriteTo, except that it puts the
|
|
// packet in the stash queue (accessible via Unstash), rather than the outgoing
|
|
// queue (accessible via OutgoingQueue).
|
|
func (c *QueuePacketConn) Stash(p []byte, addr net.Addr) bool {
|
|
return c.remotes.Stash(addr, p)
|
|
}
|
|
|
|
// Unstash returns the channel that represents the stash for addr.
|
|
func (c *QueuePacketConn) Unstash(addr net.Addr) <-chan []byte {
|
|
return c.remotes.Unstash(addr)
|
|
}
|
|
|
|
// ReadFrom returns a packet and address previously stored by QueueIncoming.
|
|
func (c *QueuePacketConn) ReadFrom(p []byte) (int, net.Addr, error) {
|
|
select {
|
|
case <-c.closed:
|
|
return 0, nil, &net.OpError{Op: "read", Net: c.LocalAddr().Network(), Addr: c.LocalAddr(), Err: c.err.Load().(error)}
|
|
default:
|
|
}
|
|
select {
|
|
case <-c.closed:
|
|
return 0, nil, &net.OpError{Op: "read", Net: c.LocalAddr().Network(), Addr: c.LocalAddr(), Err: c.err.Load().(error)}
|
|
case packet := <-c.recvQueue:
|
|
return copy(p, packet.P), packet.Addr, nil
|
|
}
|
|
}
|
|
|
|
// WriteTo queues an outgoing packet for the given address. The queue can later
|
|
// be retrieved using the OutgoingQueue method.
|
|
func (c *QueuePacketConn) WriteTo(p []byte, addr net.Addr) (int, error) {
|
|
select {
|
|
case <-c.closed:
|
|
return 0, &net.OpError{Op: "write", Net: c.LocalAddr().Network(), Addr: c.LocalAddr(), Err: c.err.Load().(error)}
|
|
default:
|
|
}
|
|
// Copy the slice so that the caller may reuse it.
|
|
buf := make([]byte, len(p))
|
|
copy(buf, p)
|
|
select {
|
|
case c.remotes.SendQueue(addr) <- buf:
|
|
return len(buf), nil
|
|
default:
|
|
// Drop the outgoing packet if the send queue is full.
|
|
return len(buf), nil
|
|
}
|
|
}
|
|
|
|
// closeWithError unblocks pending operations and makes future operations fail
|
|
// with the given error. If err is nil, it becomes errClosedPacketConn.
|
|
func (c *QueuePacketConn) closeWithError(err error) error {
|
|
var newlyClosed bool
|
|
c.closeOnce.Do(func() {
|
|
newlyClosed = true
|
|
// Store the error to be returned by future PacketConn
|
|
// operations.
|
|
if err == nil {
|
|
err = errClosedPacketConn
|
|
}
|
|
c.err.Store(err)
|
|
close(c.closed)
|
|
})
|
|
if !newlyClosed {
|
|
return &net.OpError{Op: "close", Net: c.LocalAddr().Network(), Addr: c.LocalAddr(), Err: c.err.Load().(error)}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Close unblocks pending operations and makes future operations fail with a
|
|
// "closed connection" error.
|
|
func (c *QueuePacketConn) Close() error {
|
|
return c.closeWithError(nil)
|
|
}
|
|
|
|
// LocalAddr returns the localAddr value that was passed to NewQueuePacketConn.
|
|
func (c *QueuePacketConn) LocalAddr() net.Addr { return c.localAddr }
|
|
|
|
func (c *QueuePacketConn) SetDeadline(t time.Time) error { return errNotImplemented }
|
|
func (c *QueuePacketConn) SetReadDeadline(t time.Time) error { return errNotImplemented }
|
|
func (c *QueuePacketConn) SetWriteDeadline(t time.Time) error { return errNotImplemented }
|