Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion cmd/autocar/preflight.go
Original file line number Diff line number Diff line change
Expand Up @@ -86,13 +86,16 @@ func checkServerAddresses(mode, udp, tcp string, disableFallback bool) error {
// These bounds mirror the transport constructors, which cannot be invoked
// for a server without binding sockets. Use them on both the check and normal
// startup paths, before any listener is prepared.
func validateServerLocalLimits(token string, handshakeTimeout, dialTimeout time.Duration, maxUpload, maxDownload uint64) error {
func validateServerLocalLimits(token string, handshakeTimeout, dialTimeout, destinationWriteTimeout time.Duration, maxUpload, maxDownload uint64) error {
if len(token) < protocol.MinTokenLength || len(token) > protocol.MaxTokenLength {
return fmt.Errorf("tunnel: token length must be between %d and %d bytes", protocol.MinTokenLength, protocol.MaxTokenLength)
}
if handshakeTimeout < 0 || dialTimeout < 0 {
return errors.New("--handshake-timeout and --dial-timeout must not be negative")
}
if destinationWriteTimeout < 0 {
return errors.New("--destination-write-timeout must not be negative; zero uses the 5m default")
}
if maxUpload > protocol.MaxRate || maxDownload > protocol.MaxRate {
return errors.New("tunnel: pacing rate exceeds protocol maximum")
}
Expand Down
1 change: 1 addition & 0 deletions cmd/autocar/preflight_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -237,6 +237,7 @@ func TestPreflightServerRejectsInvalidConfiguration(t *testing.T) {
{"CIDR policy", "invalid denied CIDR", []string{"--deny-cidrs", "not-a-CIDR"}},
{"dial timeout", "--dial-timeout", []string{"--dial-timeout", "-1s"}},
{"handshake timeout", "--handshake-timeout", []string{"--handshake-timeout", "-1s"}},
{"destination write timeout", "--destination-write-timeout", []string{"--destination-write-timeout", "-1s"}},
{"connections", "--max-connections", []string{"--max-connections", "0"}},
{"UDP sessions", "--max-client-udp-sessions", []string{"--max-udp-sessions", "1"}},
{"rate bound", "protocol maximum", []string{"--pacing", "fixed-rate", "--max-upload-mbps", "8000001", "--max-download-mbps", "1"}},
Expand Down
56 changes: 30 additions & 26 deletions cmd/autocar/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ func runServer(parent context.Context, args []string) error {
allowClientRates := fs.Bool("allow-client-rates", false, "native protocol: allow authenticated clients to request rates within server maxima")
dialTimeout := fs.Duration("dial-timeout", 4*time.Second, "remote destination dial timeout")
handshakeTimeout := fs.Duration("handshake-timeout", 10*time.Second, "authentication and initial stream-open timeout")
destinationWriteTimeout := fs.Duration("destination-write-timeout", 5*time.Minute, "completion timeout for each TCP destination write of at most 32 KiB; zero uses 5m (not an idle timeout)")
if err := parseFlagsWithConfig(fs, args); err != nil {
return err
}
Expand Down Expand Up @@ -128,7 +129,7 @@ func runServer(parent context.Context, args []string) error {
if err != nil {
return err
}
if err := validateServerLocalLimits(token, *handshakeTimeout, *dialTimeout, maxUpload, maxDownload); err != nil {
if err := validateServerLocalLimits(token, *handshakeTimeout, *dialTimeout, *destinationWriteTimeout, maxUpload, maxDownload); err != nil {
return err
}
certificate, err := security.LoadKeyPair(*certFile, *keyFile)
Expand Down Expand Up @@ -188,22 +189,23 @@ func runServer(parent context.Context, args []string) error {
}
if serverProtocol == "web" {
webServer, err := tunnel.ListenWeb(tunnel.WebServerConfig{
TCPAddress: *tcpListen,
UDPAddress: *listen,
Token: token,
TLSConfig: tlsConfig,
Cover: coverHandler,
Dialer: safeDialer,
HandshakeTimeout: *handshakeTimeout,
DialTimeout: *dialTimeout,
MaxConcurrentStreams: *maxStreams,
StreamAdmission: streamAdmission,
MaxConnections: *maxConnections,
MaxClientConnections: *maxClientConnections,
MaxUDPSessions: *maxUDPSessions,
MaxClientUDPSessions: *maxClientUDPSessions,
MaxUDPDestinations: *maxUDPDestinations,
UDPReceiveQueue: *udpReceiveQueue,
TCPAddress: *tcpListen,
UDPAddress: *listen,
Token: token,
TLSConfig: tlsConfig,
Cover: coverHandler,
Dialer: safeDialer,
HandshakeTimeout: *handshakeTimeout,
DialTimeout: *dialTimeout,
DestinationWriteTimeout: *destinationWriteTimeout,
MaxConcurrentStreams: *maxStreams,
StreamAdmission: streamAdmission,
MaxConnections: *maxConnections,
MaxClientConnections: *maxClientConnections,
MaxUDPSessions: *maxUDPSessions,
MaxClientUDPSessions: *maxClientUDPSessions,
MaxUDPDestinations: *maxUDPDestinations,
UDPReceiveQueue: *udpReceiveQueue,
})
if err != nil {
return err
Expand All @@ -229,6 +231,7 @@ func runServer(parent context.Context, args []string) error {
AllowClientRates: *allowClientRates,
HandshakeTimeout: *handshakeTimeout,
DialTimeout: *dialTimeout,
DestinationWriteTimeout: *destinationWriteTimeout,
MaxConcurrentStreams: *maxStreams,
StreamAdmission: streamAdmission,
MaxConnections: *maxConnections,
Expand All @@ -249,15 +252,16 @@ func runServer(parent context.Context, args []string) error {
var tlsServer *tunnel.TLSServer
if !*disableFallback {
tlsServer, err = tunnel.ListenTLS(tunnel.TLSServerConfig{
Address: *tcpListen,
Token: token,
TLSConfig: tlsConfig,
Dialer: safeDialer,
HandshakeTimeout: *handshakeTimeout,
DialTimeout: *dialTimeout,
MaxConcurrentStreams: *maxStreams,
StreamAdmission: streamAdmission,
MaxClientConnections: *maxClientFallbackConnections,
Address: *tcpListen,
Token: token,
TLSConfig: tlsConfig,
Dialer: safeDialer,
HandshakeTimeout: *handshakeTimeout,
DialTimeout: *dialTimeout,
DestinationWriteTimeout: *destinationWriteTimeout,
MaxConcurrentStreams: *maxStreams,
StreamAdmission: streamAdmission,
MaxClientConnections: *maxClientFallbackConnections,
})
if err != nil {
return err
Expand Down
104 changes: 104 additions & 0 deletions cmd/autocar/server_destination_timeout_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
package main

import (
"context"
"flag"
"io"
"strings"
"testing"
"time"
)

func TestServerDestinationWriteTimeoutCLIAndCheck(t *testing.T) {
clearPreflightEnvironment(t)
files := newPreflightFiles(t)
dnsCalls := denyPreflightDNS(t)
for _, mode := range []string{"native", "web"} {
for _, value := range []string{"", "0s", "37ms", "8m"} {
t.Run(mode+"/"+value, func(t *testing.T) {
args := append(files.serverArgs(), "--protocol", mode)
if mode == "web" {
args = append(args, "--cover-root", files.cover)
}
if value != "" {
args = append(args, "--destination-write-timeout", value)
}
if err := runServer(context.Background(), args); err != nil {
t.Fatal(err)
}
})
}
for _, check := range []string{"true", "false"} {
args := append(files.serverArgs(), "--protocol", mode, "--check="+check, "--destination-write-timeout=-1s")
if mode == "web" {
args = append(args, "--cover-root", files.cover)
}
if err := runServer(context.Background(), args); err == nil || !strings.Contains(err.Error(), "--destination-write-timeout") {
t.Fatalf("%s check=%s accepted negative timeout: %v", mode, check, err)
}
}
}
if got := dnsCalls.Load(); got != 0 {
t.Fatalf("local timeout checks attempted %d DNS operations", got)
}
}

func TestServerDestinationWriteTimeoutJSON(t *testing.T) {
clearPreflightEnvironment(t)
files := newPreflightFiles(t)
for _, mode := range []string{"native", "web"} {
for _, value := range []string{`"0s"`, `"37ms"`, `"-1s"`, `0`, `300`, `true`, `"bad-duration"`} {
t.Run(mode+"/"+value, func(t *testing.T) {
path := writeTestCommandConfig(t, `{"destination-write-timeout":`+value+`}`)
args := append(files.serverArgs(), "--protocol", mode, "--config", path)
if mode == "web" {
args = append(args, "--cover-root", files.cover)
}
err := runServer(context.Background(), args)
valid := value == `"0s"` || value == `"37ms"`
if (err == nil) != valid {
t.Fatalf("duration %s: error=%v valid=%t", value, err, valid)
}
})
}
}
}

func TestDestinationWriteTimeoutConfigPrecedence(t *testing.T) {
path := writeTestCommandConfig(t, `{"destination-write-timeout":"37s"}`)
for _, test := range []struct {
args []string
want time.Duration
}{
{nil, 5 * time.Minute},
{[]string{"--config", path}, 37 * time.Second},
{[]string{"--config", path, "--destination-write-timeout=9s"}, 9 * time.Second},
{[]string{"--config", path, "--destination-write-timeout=0s"}, 0},
} {
fs := flag.NewFlagSet("server", flag.ContinueOnError)
fs.SetOutput(io.Discard)
value := fs.Duration("destination-write-timeout", 5*time.Minute, "")
if err := parseFlagsWithConfig(fs, test.args); err != nil {
t.Fatal(err)
}
if *value != test.want {
t.Fatalf("timeout=%s want=%s", *value, test.want)
}
}
for _, value := range []string{`0`, `"invalid"`} {
path := writeTestCommandConfig(t, `{"destination-write-timeout":`+value+`}`)
fs := flag.NewFlagSet("server", flag.ContinueOnError)
fs.SetOutput(io.Discard)
fs.Duration("destination-write-timeout", 5*time.Minute, "")
if err := parseFlagsWithConfig(fs, []string{"--config", path, "--destination-write-timeout=9s"}); err == nil {
t.Fatal("invalid config type/format was hidden by CLI override")
}
}
}

func TestDestinationWriteTimeoutConfigIsServerOnly(t *testing.T) {
path := writeTestCommandConfig(t, `{"destination-write-timeout":"30s"}`)
if err := runClient(context.Background(), []string{"--config", path}); err == nil || !strings.Contains(err.Error(), "unknown or unsupported config option") {
t.Fatalf("client accepted the server-only destination timeout: %v", err)
}
}
19 changes: 19 additions & 0 deletions docs/DEPLOYMENT.md
Original file line number Diff line number Diff line change
Expand Up @@ -253,6 +253,25 @@ progress as activity, including small buffered responses. Stalled request or
response bodies and blocked writes remain bounded; header and keepalive
timeouts are unchanged.

The relay's separate `--destination-write-timeout` defaults to `5m`; `0`
also selects `5m`, rather than disabling the limit. It bounds the completion
of each TCP destination write, split into chunks of at most 32 KiB. Each
successfully completed chunk gets a fresh budget for the next write. This is
not a kernel-level no-progress timer: a write can time out after making partial
progress, and that error is retained rather than silently retried.

This server setting applies to native QUIC/TLS and web H2/H3 TCP tunnels, not
UDP or ordinary cover requests. The write deadline is cleared after every
write, including failures; when no write is pending it does not count idle
time, change read deadlines, or limit total upload duration. A stalled target
therefore releases its stream slot within the pending write's timeout, even
in the H3 response-FIN/reset edge case described in [WEB_COVER.md](WEB_COVER.md).
That is bounded cleanup, not a promise of immediate reset detection. Configure
it on the server as `--destination-write-timeout=30s` or JSON
`"destination-write-timeout": "30s"`; choose a longer value if individual
32 KiB writes can legitimately take longer. Negative durations are rejected by
normal startup and `--check`.

Before leaving a client running, verify a real authenticated relay path with
the same connection flags:

Expand Down
19 changes: 14 additions & 5 deletions docs/WEB_COVER.md
Original file line number Diff line number Diff line change
Expand Up @@ -45,11 +45,20 @@ uploading after receiving a destination EOF require the H3 stream transport.
Canceled H2 requests and H3 send-side stream errors close the destination directly,
including when both relay workers are blocked on destination I/O. Server or
physical-connection shutdown also releases H3 destinations after a response
FIN. One H3 edge case remains: after a clean response FIN, resetting only that
stream cannot interrupt an already blocked destination upload write through
the current QUIC API. The stream slot can remain occupied until the destination
unblocks or the physical connection closes; a sibling stream can remain usable.
This is not a guarantee that every canceled stream is immediately reclaimed.
FIN. After a clean response FIN, resetting only that H3 stream still cannot
directly interrupt an already blocked destination upload write through the
current QUIC API. The server's `--destination-write-timeout` now bounds that
pending write, releasing the stream slot without closing usable siblings.
Its default is `5m`; `0` selects that default and negative values are invalid.
This is bounded cleanup, not immediate reset detection.

The limit is the completion deadline of each TCP destination write, in chunks
of at most 32 KiB—not a kernel-level no-progress timer or a total upload limit.
After a chunk completes successfully, the next write gets a new budget.
Partial-write errors remain errors. Deadlines are cleared after writes, do not run while no write
is pending, and never change destination read deadlines. Native QUIC/TLS TCP
tunnels share this setting; UDP and ordinary cover traffic do not. See
[deployment settings](DEPLOYMENT.md) for CLI and JSON examples.

H2's handshake budget covers both TLS negotiation and the initial HTTP/2
preface/SETTINGS write. Caller cancellation or client shutdown also interrupts
Expand Down
3 changes: 2 additions & 1 deletion examples/server.json
Original file line number Diff line number Diff line change
Expand Up @@ -3,5 +3,6 @@
"listen": ":8443",
"cert": "server.crt",
"key": "server.key",
"token-file": "relay-token"
"token-file": "relay-token",
"destination-write-timeout": "5m"
}
42 changes: 25 additions & 17 deletions internal/tunnel/common.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,11 +21,12 @@ import (
)

const (
defaultHandshakeTimeout = 10 * time.Second
defaultDialTimeout = 10 * time.Second
defaultMaxStreams = 1024
defaultMaxConnections = 256
defaultMaxClientConnections = 32
defaultHandshakeTimeout = 10 * time.Second
defaultDialTimeout = 10 * time.Second
defaultDestinationWriteTimeout = 5 * time.Minute
defaultMaxStreams = 1024
defaultMaxConnections = 256
defaultMaxClientConnections = 32
)

// RemoteError is returned when the authenticated exit rejects a CONNECT
Expand All @@ -44,11 +45,12 @@ func (e *RemoteError) Error() string {
}

type serverCore struct {
tokenHash [sha256.Size]byte
dialer transport.Dialer
handshakeTimeout time.Duration
dialTimeout time.Duration
sem chan struct{}
tokenHash [sha256.Size]byte
dialer transport.Dialer
handshakeTimeout time.Duration
dialTimeout time.Duration
destinationWriteTimeout time.Duration
sem chan struct{}
}

// StreamAdmission is a concurrency budget for active relay streams. Pass the
Expand All @@ -70,14 +72,15 @@ func NewStreamAdmission(limit int) (*StreamAdmission, error) {
}

func newServerCore(token string, dialer transport.Dialer, handshakeTimeout, dialTimeout time.Duration, maxStreams int) (*serverCore, error) {
return newServerCoreWithAdmission(token, dialer, handshakeTimeout, dialTimeout, maxStreams, nil)
return newServerCoreWithAdmission(token, dialer, handshakeTimeout, dialTimeout, 0, maxStreams, nil)
}

func newServerCoreWithAdmission(
token string,
dialer transport.Dialer,
handshakeTimeout time.Duration,
dialTimeout time.Duration,
destinationWriteTimeout time.Duration,
maxStreams int,
admission *StreamAdmission,
) (*serverCore, error) {
Expand All @@ -88,7 +91,7 @@ func newServerCoreWithAdmission(
netDialer := &net.Dialer{Timeout: defaultDialTimeout, KeepAlive: 30 * time.Second}
dialer = netDialer
}
if handshakeTimeout < 0 || dialTimeout < 0 || maxStreams < 0 {
if handshakeTimeout < 0 || dialTimeout < 0 || destinationWriteTimeout < 0 || maxStreams < 0 {
return nil, errors.New("tunnel: timeout and concurrency limits cannot be negative")
}
if handshakeTimeout == 0 {
Expand All @@ -97,6 +100,9 @@ func newServerCoreWithAdmission(
if dialTimeout == 0 {
dialTimeout = defaultDialTimeout
}
if destinationWriteTimeout == 0 {
destinationWriteTimeout = defaultDestinationWriteTimeout
}
if maxStreams == 0 {
if admission == nil {
maxStreams = defaultMaxStreams
Expand All @@ -122,11 +128,12 @@ func newServerCoreWithAdmission(
}
}
return &serverCore{
tokenHash: sha256.Sum256([]byte(token)),
dialer: dialer,
handshakeTimeout: handshakeTimeout,
dialTimeout: dialTimeout,
sem: admission.sem,
tokenHash: sha256.Sum256([]byte(token)),
dialer: dialer,
handshakeTimeout: handshakeTimeout,
dialTimeout: dialTimeout,
destinationWriteTimeout: destinationWriteTimeout,
sem: admission.sem,
}, nil
}

Expand Down Expand Up @@ -217,6 +224,7 @@ func (s *serverCore) handleStream(
_ = protocol.WriteResponse(stream, protocol.Response{Status: protocol.StatusDialFailed, Message: "destination unavailable"})
return
}
upstream = s.boundDestinationWrites(upstream)
defer upstream.Close()
options.Response.Status = protocol.StatusOK
if err := protocol.WriteResponse(stream, options.Response); err != nil {
Expand Down
Loading
Loading