feat: add TCP_NODELAY utilities and apply them to connections and listeners for improved latency
cognitive-glitch committed
Dec 9, 2025 at 10:06 UTC
572d31d18c8e4d5dd48ec0b9ef9f18430bc8d5ee
8 files changed
+173
-20
cmd/relay-server/bps_manager.go
+1
-1
@@ -12,7 +12,7 @@ import (
12
// BPSManager manages per-lease bytes-per-second rate limiting
13
type BPSManager struct {
14
mu sync.Mutex
15
- bpsLimits map[string]int64 // leaseID -> bytes-per-second (0 = unlimited)
15
+ bpsLimits map[string]int64 // leaseID -> bytes-per-second (0 = unlimited)
16
bpsBuckets map[string]*ratelimit.Bucket // leaseID -> rate limit bucket
17
defaultBPS int64 // default bytes-per-second for new leases
18
}
cmd/relay-server/view.go
+9
-3
@@ -176,13 +176,19 @@ func serveHTTP(addr string, serv *portal.RelayServer, bpsManager *BPSManager, no
176
})
177
178
srv := &http.Server{
179
- Addr: addr,
179
Handler: handler,
180
}
181
182
+ // Create TCP listener with TCP_NODELAY enabled for low-latency relay protocol
183
+ listener, err := net.Listen("tcp", addr)
184
+ if err != nil {
185
+ log.Fatal().Err(err).Msgf("[server] failed to listen on %s", addr)
186
+ }
187
+ noDelayListener := utils.NewTCPNoDelayListener(listener)
188
+
189
go func() {
190
log.Info().Msgf("[server] http: %s", addr)
185
- if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
191
+ if err := srv.Serve(noDelayListener); err != nil && err != http.ErrServerClosed {
192
log.Error().Err(err).Msg("[server] http error")
193
cancel()
194
}
@@ -558,7 +564,7 @@ func convertLeaseEntriesToRows(serv *portal.RelayServer) []leaseRow {
564
Metadata: lease.Metadata,
565
}
566
561
- if row.Hide != true {
567
+ if !row.Hide {
568
rows = append(rows, row)
569
}
570
}
portal/core/cryptoops/handshaker_test.go
+15
-5
@@ -11,14 +11,17 @@ import (
11
12
"golang.org/x/crypto/curve25519"
13
"gosuda.org/portal/portal/core/proto/rdsec"
14
+ "gosuda.org/portal/utils"
15
)
16
16
-// pipeConn creates a bidirectional pipe for testing using TCP loopback
17
+// pipeConn creates a bidirectional pipe for testing using TCP loopback.
18
+// TCP_NODELAY is enabled on both connections for low-latency testing.
19
func pipeConn() (net.Conn, net.Conn) {
18
- listener, err := net.Listen("tcp", "127.0.0.1:0")
20
+ rawListener, err := net.Listen("tcp", "127.0.0.1:0")
21
if err != nil {
22
panic(err)
23
}
24
+ listener := utils.NewTCPNoDelayListener(rawListener)
25
26
connCh := make(chan net.Conn, 1)
27
go func() {
@@ -34,6 +37,9 @@ func pipeConn() (net.Conn, net.Conn) {
37
if err != nil {
38
panic(err)
39
}
40
+ if err := utils.SetTCPNoDelay(clientConn); err != nil {
41
+ panic(err)
42
+ }
43
44
serverConn := <-connCh
45
return clientConn, serverConn
@@ -682,11 +688,12 @@ func TestRealNetworkConnection(t *testing.T) {
688
clientCred, _ := NewCredential()
689
serverCred, _ := NewCredential()
690
685
- // Start server
686
- listener, err := net.Listen("tcp", "127.0.0.1:0")
691
+ // Start server with TCP_NODELAY enabled
692
+ rawListener, err := net.Listen("tcp", "127.0.0.1:0")
693
if err != nil {
694
t.Fatalf("Failed to start listener: %v", err)
695
}
696
+ listener := utils.NewTCPNoDelayListener(rawListener)
697
defer listener.Close()
698
699
serverAddr := listener.Addr().String()
@@ -707,11 +714,14 @@ func TestRealNetworkConnection(t *testing.T) {
714
serverSecure, serverErr = serverHandshaker.ServerHandshake(conn, []string{"test-alpn"})
715
}()
716
710
- // Connect client
717
+ // Connect client with TCP_NODELAY enabled
718
clientConn, err := net.Dial("tcp", serverAddr)
719
if err != nil {
720
t.Fatalf("Failed to connect: %v", err)
721
}
722
+ if err := utils.SetTCPNoDelay(clientConn); err != nil {
723
+ t.Fatalf("Failed to set TCP_NODELAY: %v", err)
724
+ }
725
726
clientHandshaker := NewHandshaker(clientCred)
727
clientSecure, clientErr := clientHandshaker.ClientHandshake(clientConn, "test-alpn")
portal/integration_test.go
+6
-2
@@ -10,6 +10,7 @@ import (
10
"github.com/stretchr/testify/require"
11
"gosuda.org/portal/portal/core/cryptoops"
12
"gosuda.org/portal/portal/core/proto/rdverb"
13
+ "gosuda.org/portal/utils"
14
)
15
16
// generateTestCredential creates a new credential for testing
@@ -26,9 +27,10 @@ func TestIntegration_FullFlow(t *testing.T) {
27
server.Start()
28
defer server.Stop()
29
29
- // Create a listener for the server
30
- listener, err := net.Listen("tcp", "127.0.0.1:0")
30
+ // Create a listener for the server with TCP_NODELAY enabled
31
+ rawListener, err := net.Listen("tcp", "127.0.0.1:0")
32
require.NoError(t, err)
33
+ listener := utils.NewTCPNoDelayListener(rawListener)
34
defer listener.Close()
35
36
go func() {
@@ -47,6 +49,7 @@ func TestIntegration_FullFlow(t *testing.T) {
49
hostCred := generateTestCredential(t)
50
hostConn, err := net.Dial("tcp", serverAddr)
51
require.NoError(t, err)
52
+ require.NoError(t, utils.SetTCPNoDelay(hostConn))
53
54
hostClient := NewRelayClient(hostConn)
55
require.NotNil(t, hostClient)
@@ -75,6 +78,7 @@ func TestIntegration_FullFlow(t *testing.T) {
78
peerCred := generateTestCredential(t)
79
peerConn, err := net.Dial("tcp", serverAddr)
80
require.NoError(t, err)
81
+ require.NoError(t, utils.SetTCPNoDelay(peerConn))
82
83
peerClient := NewRelayClient(peerConn)
84
require.NotNil(t, peerClient)
sdk/sdk.go
+1
-1
@@ -302,7 +302,7 @@ func (g *Client) listenerWorker(server *connRelay) {
302
303
if !exists {
304
log.Warn().Str("lease_id", lease).Msg("[SDK] No listener found for lease, closing connection")
305
- incoming.SecureConnection.Close() // Close unused connection
305
+ incoming.Close() // Close unused connection
306
continue
307
}
308
sdk/sdk_e2e_test.go
+18
-6
@@ -59,14 +59,18 @@ func TestE2E_ClientToAppThroughRelay(t *testing.T) {
59
}
60
})
61
62
+ // Create TCP listener with TCP_NODELAY enabled for low-latency relay protocol
63
+ relayListener, err := net.Listen("tcp", relayAddr)
64
+ require.NoError(t, err, "Failed to create relay listener")
65
+ noDelayListener := utils.NewTCPNoDelayListener(relayListener)
66
+
67
relayHTTPServer := &http.Server{
63
- Addr: relayAddr,
68
Handler: relayMux,
69
}
70
71
go func() {
72
log.Info().Str("addr", relayAddr).Msg("[TEST] Relay HTTP server starting")
69
- if err := relayHTTPServer.ListenAndServe(); err != nil && err != http.ErrServerClosed {
73
+ if err := relayHTTPServer.Serve(noDelayListener); err != nil && err != http.ErrServerClosed {
74
log.Error().Err(err).Msg("[TEST] Relay HTTP server error")
75
}
76
}()
@@ -219,12 +223,16 @@ func TestE2E_MultipleConnections(t *testing.T) {
223
relayServer.HandleConnection(stream)
224
})
225
226
+ // Create TCP listener with TCP_NODELAY enabled
227
+ relayListener, err := net.Listen("tcp", relayAddr)
228
+ require.NoError(t, err, "Failed to create relay listener")
229
+ noDelayListener := utils.NewTCPNoDelayListener(relayListener)
230
+
231
relayHTTPServer := &http.Server{
223
- Addr: relayAddr,
232
Handler: relayMux,
233
}
234
227
- go relayHTTPServer.ListenAndServe()
235
+ go relayHTTPServer.Serve(noDelayListener)
236
defer func() {
237
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
238
defer cancel()
@@ -359,12 +367,16 @@ func TestE2E_ConnectionTimeout(t *testing.T) {
367
relayServer.HandleConnection(stream)
368
})
369
370
+ // Create TCP listener with TCP_NODELAY enabled
371
+ relayListener, err := net.Listen("tcp", relayAddr)
372
+ require.NoError(t, err, "Failed to create relay listener")
373
+ noDelayListener := utils.NewTCPNoDelayListener(relayListener)
374
+
375
relayHTTPServer := &http.Server{
363
- Addr: relayAddr,
376
Handler: relayMux,
377
}
378
367
- go relayHTTPServer.ListenAndServe()
379
+ go relayHTTPServer.Serve(noDelayListener)
380
defer func() {
381
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
382
defer cancel()
utils/utils.go
+51
-2
@@ -5,21 +5,70 @@ import (
5
"fmt"
6
"io"
7
"mime"
8
+ "net"
9
"net/http"
10
"net/url"
11
"regexp"
12
"strings"
13
14
"github.com/gorilla/websocket"
15
+ "github.com/rs/zerolog/log"
16
17
"gosuda.org/portal/portal/utils/wsstream"
18
)
19
20
+// SetTCPNoDelay enables TCP_NODELAY on a TCP connection to disable Nagle's algorithm.
21
+// Returns nil for non-TCP connections (e.g., Unix sockets, WebSocket over WASM).
22
+func SetTCPNoDelay(conn net.Conn) error {
23
+ if tcpConn, ok := conn.(*net.TCPConn); ok {
24
+ return tcpConn.SetNoDelay(true)
25
+ }
26
+ return nil
27
+}
28
+
29
+// TCPNoDelayListener wraps a net.Listener to enable TCP_NODELAY on accepted connections.
30
+type TCPNoDelayListener struct {
31
+ net.Listener
32
+}
33
+
34
+// Accept accepts a connection and enables TCP_NODELAY.
35
+func (l *TCPNoDelayListener) Accept() (net.Conn, error) {
36
+ conn, err := l.Listener.Accept()
37
+ if err != nil {
38
+ return nil, err
39
+ }
40
+ if err := SetTCPNoDelay(conn); err != nil {
41
+ log.Debug().Err(err).Msg("failed to set TCP_NODELAY on accepted connection")
42
+ }
43
+ return conn, nil
44
+}
45
+
46
+// NewTCPNoDelayListener wraps a listener to enable TCP_NODELAY on accepted connections.
47
+func NewTCPNoDelayListener(l net.Listener) *TCPNoDelayListener {
48
+ return &TCPNoDelayListener{Listener: l}
49
+}
50
+
51
// NewWebSocketDialer returns a dialer that establishes WebSocket connections
19
-// and wraps them as io.ReadWriteCloser.
52
+// and wraps them as io.ReadWriteCloser. TCP_NODELAY is enabled on the underlying
53
+// TCP connection to minimize latency for interactive relay protocols.
54
func NewWebSocketDialer() func(context.Context, string) (io.ReadWriteCloser, error) {
55
+ dialer := &websocket.Dialer{
56
+ NetDialContext: func(ctx context.Context, network, addr string) (net.Conn, error) {
57
+ d := &net.Dialer{}
58
+ conn, err := d.DialContext(ctx, network, addr)
59
+ if err != nil {
60
+ return nil, err
61
+ }
62
+ if err := SetTCPNoDelay(conn); err != nil {
63
+ log.Debug().Err(err).Msg("failed to set TCP_NODELAY on WebSocket connection")
64
+ }
65
+ return conn, nil
66
+ },
67
+ HandshakeTimeout: websocket.DefaultDialer.HandshakeTimeout,
68
+ }
69
+
70
return func(ctx context.Context, url string) (io.ReadWriteCloser, error) {
22
- wsConn, _, err := websocket.DefaultDialer.Dial(url, nil)
71
+ wsConn, _, err := dialer.DialContext(ctx, url, nil)
72
if err != nil {
73
return nil, err
74
}
utils/utils_test.go
+72
@@ -1,9 +1,11 @@
1
package utils
2
3
import (
4
+ "net"
5
"testing"
6
7
"github.com/stretchr/testify/assert"
8
+ "github.com/stretchr/testify/require"
9
)
10
11
func TestIsURLSafeName(t *testing.T) {
@@ -258,3 +260,73 @@ func TestIsSubdomain(t *testing.T) {
260
})
261
}
262
}
263
+
264
+func TestSetTCPNoDelay(t *testing.T) {
265
+ listener, err := net.Listen("tcp", "127.0.0.1:0")
266
+ require.NoError(t, err)
267
+ defer listener.Close()
268
+
269
+ go func() {
270
+ conn, _ := listener.Accept()
271
+ if conn != nil {
272
+ conn.Close()
273
+ }
274
+ }()
275
+
276
+ conn, err := net.Dial("tcp", listener.Addr().String())
277
+ require.NoError(t, err)
278
+ defer conn.Close()
279
+
280
+ // Should succeed on TCP connection
281
+ err = SetTCPNoDelay(conn)
282
+ require.NoError(t, err)
283
+
284
+ // Verify it's a TCP connection
285
+ _, ok := conn.(*net.TCPConn)
286
+ require.True(t, ok, "expected *net.TCPConn")
287
+}
288
+
289
+func TestSetTCPNoDelay_NonTCP(t *testing.T) {
290
+ // Test with a non-TCP connection (using a pipe)
291
+ server, client := net.Pipe()
292
+ defer server.Close()
293
+ defer client.Close()
294
+
295
+ // Should return nil for non-TCP connections (no-op)
296
+ err := SetTCPNoDelay(client)
297
+ require.NoError(t, err)
298
+}
299
+
300
+func TestTCPNoDelayListener(t *testing.T) {
301
+ rawListener, err := net.Listen("tcp", "127.0.0.1:0")
302
+ require.NoError(t, err)
303
+
304
+ listener := NewTCPNoDelayListener(rawListener)
305
+ defer listener.Close()
306
+
307
+ go func() {
308
+ conn, err := net.Dial("tcp", listener.Addr().String())
309
+ if err == nil {
310
+ conn.Close()
311
+ }
312
+ }()
313
+
314
+ conn, err := listener.Accept()
315
+ require.NoError(t, err)
316
+ defer conn.Close()
317
+
318
+ // Verify it's a TCP connection (TCP_NODELAY was set during Accept)
319
+ _, ok := conn.(*net.TCPConn)
320
+ require.True(t, ok, "expected *net.TCPConn from TCPNoDelayListener")
321
+}
322
+
323
+func TestTCPNoDelayListener_Addr(t *testing.T) {
324
+ rawListener, err := net.Listen("tcp", "127.0.0.1:0")
325
+ require.NoError(t, err)
326
+
327
+ listener := NewTCPNoDelayListener(rawListener)
328
+ defer listener.Close()
329
+
330
+ // Verify Addr() returns the correct address
331
+ require.Equal(t, rawListener.Addr(), listener.Addr())
332
+}