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
6 changes: 6 additions & 0 deletions docs/DEPLOYMENT.md
Original file line number Diff line number Diff line change
Expand Up @@ -244,6 +244,12 @@ behavior and limits are in [WEB_COVER.md](WEB_COVER.md).
Default local endpoints are loopback-only SOCKS5 `127.0.0.1:1080` and HTTP
`127.0.0.1:8080`. Use `socks5h://` when the relay should resolve names.

For SOCKS5 TCP and HTTP CONNECT tunnels, `--idle-timeout` (default `5m`)
measures inactivity across both directions: an active download or upload does
not need reverse-direction application traffic to stay open. A blocked write
still has its own timeout, so an unresponsive receiver cannot retain a tunnel
indefinitely. Ordinary forwarded HTTP bodies keep per-operation timeouts.

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

Expand Down
6 changes: 6 additions & 0 deletions docs/WEB_COVER.md
Original file line number Diff line number Diff line change
Expand Up @@ -165,6 +165,12 @@ both pacing fields as `not-applicable`, and `--pacing=fixed-rate` is rejected.
`h2` does not advertise UDP. The H3-to-H2 fallback carries only new TCP
streams: UDP never falls back to H2 and fails when H3 is unavailable.

Only a newly authenticated CONNECT-UDP response is fresh evidence that the H3
path has recovered. Sending on a cached UDP target merely queues a datagram
locally; it does not clear the TCP fallback cooldown or change the last
successful transport. A newly authenticated target rejection proves path
health without being counted as a successful target connection.

## CONNECT-UDP boundary

The H3 UDP path follows
Expand Down
211 changes: 211 additions & 0 deletions integration/proxy_activity_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,211 @@
package integration_test

import (
"bufio"
"bytes"
"context"
"fmt"
"io"
"net"
"net/http"
"strings"
"testing"
"time"

"github.com/cppla/autocar/internal/proxy"
"github.com/cppla/autocar/internal/transport"
"github.com/cppla/autocar/internal/tunnel"
)

// A completed upload must not leave a write-deadline timer armed on the H2
// stream, and its quiet read direction must not expire an active download.
// Unlike an echo or a half-close test, both directions stay open throughout.
func TestHTTPConnectWebH2ActiveDownloadSurvivesIdleTimeout(t *testing.T) {
const idleTimeout = 250 * time.Millisecond
const interval = 50 * time.Millisecond
request := []byte("go!")
chunk := []byte("download-chunk\x00\xff")
payload := bytes.Repeat(chunk, 16)
dialer := startActivityWebH2Client(t)

for _, pipelined := range []bool{false, true} {
name := "request_after_connect"
if pipelined {
name = "request_with_connect"
}
t.Run(name, func(t *testing.T) {
target, targetResult := startActivityDownloadTarget(t, request, chunk, 16, interval)
proxyAddress := startActivityHTTPProxy(t, dialer, idleTimeout)
conn := dialWithDeadline(t, proxyAddress)
defer conn.Close()
head := fmt.Sprintf("CONNECT %s HTTP/1.1\r\nHost: %s\r\n\r\n", target, target)
initial := []byte(head)
if pipelined {
initial = append(initial, request...)
}
started := time.Now()
writeFull(t, conn, initial)
reader := bufio.NewReader(conn)
status, err := reader.ReadString('\n')
if err != nil || !strings.Contains(status, " 200 ") {
t.Fatalf("CONNECT response = %q, err = %v", status, err)
}
for {
line, err := reader.ReadString('\n')
if err != nil {
t.Fatalf("read CONNECT headers: %v", err)
}
if line == "\r\n" {
break
}
}
if !pipelined {
writeFull(t, conn, request)
}

got := make([]byte, len(payload))
n, err := io.ReadFull(reader, got)
if err != nil {
t.Fatalf("active H2 download stopped after %d/%d bytes and %s: %v", n, len(got), time.Since(started), err)
}
if !bytes.Equal(got, payload) {
t.Fatal("download payload differs from target bytes")
}
if elapsed := time.Since(started); elapsed <= 2*idleTimeout {
t.Fatalf("download lasted %s, want more than two idle intervals", elapsed)
}
// The upload side must still work after the long download, without
// relying on EOF or CloseWrite to complete either direction.
writeFull(t, conn, []byte("!"))
select {
case err := <-targetResult:
if err != nil {
t.Fatalf("target exchange: %v", err)
}
case <-time.After(operationTimeout):
t.Fatal("target did not receive final acknowledgment")
}
})
}
}

func startActivityWebH2Client(t *testing.T) *tunnel.WebH2Client {
t.Helper()
material := newTLSMaterial(t)
server, err := tunnel.ListenWebH2(tunnel.WebH2ServerConfig{
Address: "127.0.0.1:0", Token: testToken, TLSConfig: material.server,
Cover: http.NotFoundHandler(), Dialer: &net.Dialer{Timeout: operationTimeout},
HandshakeTimeout: operationTimeout, DialTimeout: operationTimeout,
})
if err != nil {
t.Fatalf("listen activity-test H2 relay: %v", err)
}
ctx, cancel := context.WithCancel(context.Background())
done := make(chan error, 1)
go func() { done <- server.Serve(ctx) }()
t.Cleanup(func() {
cancel()
_ = server.Close()
waitForServe(t, "activity-test H2 relay", done)
})
client, err := tunnel.NewWebH2Client(tunnel.WebH2ClientConfig{
ServerAddress: server.Addr().String(), Token: testToken, TLSConfig: material.client,
HandshakeTimeout: operationTimeout, DialTimeout: operationTimeout,
})
if err != nil {
t.Fatalf("create activity-test H2 client: %v", err)
}
t.Cleanup(func() { _ = client.Close() })
return client
}

func startActivityHTTPProxy(t *testing.T, dialer transport.Dialer, idleTimeout time.Duration) string {
t.Helper()
server, err := proxy.NewHTTPServer(proxy.Config{
Dialer: dialer, HandshakeTimeout: operationTimeout, DialTimeout: operationTimeout,
IdleTimeout: idleTimeout, MaxConnections: 4,
})
if err != nil {
t.Fatalf("create activity-test HTTP proxy: %v", err)
}
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen activity-test HTTP proxy: %v", err)
}
done := make(chan error, 1)
go func() { done <- server.Serve(listener) }()
t.Cleanup(func() {
ctx, cancel := context.WithTimeout(context.Background(), operationTimeout)
defer cancel()
if err := server.Shutdown(ctx); err != nil && !isExpectedClose(err) {
t.Errorf("shut down activity-test HTTP proxy: %v", err)
}
waitForServe(t, "activity-test HTTP proxy", done)
})
return listener.Addr().String()
}

func startActivityDownloadTarget(t *testing.T, request, chunk []byte, count int, interval time.Duration) (string, <-chan error) {
t.Helper()
listener, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen activity-test target: %v", err)
}
ctx, cancel := context.WithCancel(context.Background())
result := make(chan error, 1)
stopped := make(chan struct{})
go func() {
defer close(stopped)
result <- func() error {
conn, err := listener.Accept()
if err != nil {
return err
}
defer conn.Close()
stop := context.AfterFunc(ctx, func() { _ = conn.Close() })
defer stop()
if err := conn.SetDeadline(time.Now().Add(operationTimeout)); err != nil {
return err
}
got := make([]byte, len(request))
if _, err := io.ReadFull(conn, got); err != nil {
return fmt.Errorf("read initial request: %w", err)
}
if !bytes.Equal(got, request) {
return fmt.Errorf("initial request = %q, want %q", got, request)
}
ticker := time.NewTicker(interval)
defer ticker.Stop()
for range count {
select {
case <-ticker.C:
case <-ctx.Done():
return ctx.Err()
}
if n, err := conn.Write(chunk); err != nil {
return fmt.Errorf("write download chunk: %w", err)
} else if n != len(chunk) {
return io.ErrShortWrite
}
}
ack := make([]byte, 1)
if _, err := io.ReadFull(conn, ack); err != nil {
return fmt.Errorf("read final acknowledgment: %w", err)
}
if ack[0] != '!' {
return fmt.Errorf("final acknowledgment = %q", ack)
}
return nil
}()
}()
t.Cleanup(func() {
cancel()
_ = listener.Close()
select {
case <-stopped:
case <-time.After(operationTimeout):
t.Error("activity-test target did not stop")
}
})
return listener.Addr().String(), result
}
5 changes: 4 additions & 1 deletion internal/proxy/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,10 @@ type Config struct {
Authenticator Authenticator
HandshakeTimeout time.Duration
DialTimeout time.Duration
IdleTimeout time.Duration
// IdleTimeout bounds inactivity across either direction of a CONNECT
// tunnel, and separately bounds a blocked write. Ordinary HTTP request
// and response bodies retain their per-operation inactivity bound.
IdleTimeout time.Duration
// MaxConnections is the number of accepted client TCP connections that may
// be active at once. Zero uses a conservative default; a negative value is
// rejected.
Expand Down
5 changes: 1 addition & 4 deletions internal/proxy/http.go
Original file line number Diff line number Diff line change
Expand Up @@ -329,10 +329,7 @@ func writeGatewayError(w http.ResponseWriter, err error) {

func writeAll(conn net.Conn, payload []byte, timeout time.Duration) error {
for len(payload) > 0 {
if timeout > 0 {
_ = conn.SetWriteDeadline(time.Now().Add(timeout))
}
written, err := conn.Write(payload)
written, err := writeWithStallDeadline(conn, payload, timeout)
if written < 0 || written > len(payload) {
return errors.New("proxy: invalid write count")
}
Expand Down
72 changes: 57 additions & 15 deletions internal/proxy/relay.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,12 +31,25 @@ func (c *activityConn) Read(p []byte) (int, error) {
}

func (c *activityConn) Write(p []byte) (int, error) {
if c.timeout > 0 {
if err := c.Conn.SetWriteDeadline(time.Now().Add(c.timeout)); err != nil {
return 0, err
}
return writeWithStallDeadline(c.Conn, p, c.timeout)
}

// Bound only the pending write. In particular, H2 streams implement deadlines
// by aborting the stream: leaving a completed write's timer armed would later
// terminate an otherwise active download with no more uploads.
func writeWithStallDeadline(conn net.Conn, p []byte, timeout time.Duration) (int, error) {
if timeout <= 0 {
return conn.Write(p)
}
return c.Conn.Write(p)
if err := conn.SetWriteDeadline(time.Now().Add(timeout)); err != nil {
return 0, err
}
n, err := conn.Write(p)
clearErr := conn.SetWriteDeadline(time.Time{})
if err == nil {
err = clearErr
}
return n, err
}

func (c *activityConn) CloseWrite() error {
Expand All @@ -61,11 +74,17 @@ func relay(left, right net.Conn, idleTimeout time.Duration) error {
err error
}
results := make(chan result, 2)
activity := &relayActivity{left: left, right: right, timeout: idleTimeout}
if err := activity.refresh(); err != nil {
_ = left.Close()
_ = right.Close()
return err
}
go func() {
results <- result{destination: right, err: copyHalf(right, left, idleTimeout)}
results <- result{destination: right, err: copyHalf(right, left, activity)}
}()
go func() {
results <- result{destination: left, err: copyHalf(left, right, idleTimeout)}
results <- result{destination: left, err: copyHalf(left, right, activity)}
}()

var relayErrors []error
Expand Down Expand Up @@ -101,32 +120,55 @@ func relay(left, right net.Conn, idleTimeout time.Duration) error {
return errors.Join(relayErrors...)
}

func copyHalf(dst, src net.Conn, idleTimeout time.Duration) error {
// relayActivity makes read inactivity a property of the whole tunnel, not
// either direction independently. Downloads, uploads and server-push streams
// may legitimately have no reverse-direction application data for minutes.
// Serializing the update prevents an older activity event from overwriting a
// newer deadline. Writes retain their own stall deadline for backpressure.
type relayActivity struct {
mu sync.Mutex
left, right net.Conn
timeout time.Duration
}

func (a *relayActivity) refresh() error {
if a.timeout <= 0 {
return nil
}
a.mu.Lock()
defer a.mu.Unlock()
deadline := time.Now().Add(a.timeout)
return errors.Join(a.left.SetReadDeadline(deadline), a.right.SetReadDeadline(deadline))
}

func copyHalf(dst, src net.Conn, activity *relayActivity) error {
bufferPtr := relayBufferPool.Get().(*[]byte)
defer relayBufferPool.Put(bufferPtr)
buffer := *bufferPtr
for {
if idleTimeout > 0 {
_ = src.SetReadDeadline(time.Now().Add(idleTimeout))
}
n, readErr := src.Read(buffer)
if n < 0 || n > len(buffer) {
return errors.New("proxy: invalid read count")
}
if n > 0 {
if err := activity.refresh(); err != nil {
return err
}
written := 0
for written < n {
if idleTimeout > 0 {
_ = dst.SetWriteDeadline(time.Now().Add(idleTimeout))
}
m, writeErr := dst.Write(buffer[written:n])
m, writeErr := writeWithStallDeadline(dst, buffer[written:n], activity.timeout)
if m < 0 || m > n-written {
return errors.New("proxy: invalid write count")
}
written += m
if writeErr != nil {
return writeErr
}
if m > 0 {
if err := activity.refresh(); err != nil {
return err
}
}
if m == 0 {
return io.ErrShortWrite
}
Expand Down
Loading
Loading