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 docs/DEPLOYMENT.md
Original file line number Diff line number Diff line change
Expand Up @@ -248,7 +248,10 @@ 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.
indefinitely. Ordinary HTTP forwarding also counts upload/download body
progress as activity, including small buffered responses. Stalled request or
response bodies and blocked writes remain bounded; header and keepalive
timeouts are unchanged.

Before leaving a client running, verify a real authenticated relay path with
the same connection flags:
Expand Down
13 changes: 13 additions & 0 deletions docs/WEB_COVER.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,19 @@ finish its response while continuing to receive an independent upload. Other
streams on the same connection remain usable. Applications that need to keep
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.

H2's handshake budget covers both TLS negotiation and the initial HTTP/2
preface/SETTINGS write. Caller cancellation or client shutdown also interrupts
that initialization, before the connection enters the reusable session pool.

There is no `autocar/2` ALPN or AutoCAR binary stream header on these paths.
The web ALPNs are `h2`, `h3`, and `http/1.1`. Native and web transports remain
separate modes and are not wire-compatible.
Expand Down
98 changes: 98 additions & 0 deletions integration/proxy_activity_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@ import (
"io"
"net"
"net/http"
"net/http/httptest"
"net/url"
"strings"
"testing"
"time"
Expand All @@ -17,6 +19,102 @@ import (
"github.com/cppla/autocar/internal/tunnel"
)

func TestHTTPForwardWebH2ActiveBodiesSurviveIdleTimeout(t *testing.T) {
const idle = 250 * time.Millisecond
const chunks = 16
const interval = 50 * time.Millisecond
dialer := startActivityWebH2Client(t)
for _, direction := range []string{"download", "upload"} {
t.Run(direction, func(t *testing.T) {
chunk := []byte("body-chunk\x00\xff")
payload := bytes.Repeat(chunk, chunks)
origin := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if direction == "upload" {
got, err := io.ReadAll(r.Body)
if err != nil || !bytes.Equal(got, payload) {
http.Error(w, "incomplete upload", http.StatusBadRequest)
return
}
w.WriteHeader(http.StatusNoContent)
return
}
w.Header().Set("Content-Length", fmt.Sprint(len(payload)))
for range chunks {
if _, err := w.Write(chunk); err != nil {
return
}
w.(http.Flusher).Flush()
select {
case <-r.Context().Done():
return
case <-time.After(interval):
}
}
}))
defer origin.Close()
address := startActivityHTTPProxy(t, dialer, idle)
client := &http.Client{
Timeout: operationTimeout,
Transport: &http.Transport{Proxy: http.ProxyURL(&url.URL{Scheme: "http", Host: address})},
}
defer client.CloseIdleConnections()
method := http.MethodGet
var body io.Reader
if direction == "upload" {
method = http.MethodPost
reader, writer := io.Pipe()
defer reader.Close()
defer writer.Close()
body = reader
done := make(chan struct{})
t.Cleanup(func() {
select {
case <-done:
case <-time.After(operationTimeout):
t.Error("upload writer did not stop")
}
})
go func() {
defer close(done)
defer writer.Close()
for range chunks {
if _, err := writer.Write(chunk); err != nil {
return
}
time.Sleep(interval)
}
}()
}
request, err := http.NewRequest(method, origin.URL, body)
if err != nil {
t.Fatal(err)
}
if direction == "upload" {
request.ContentLength = int64(len(payload))
}
response, err := client.Do(request)
if err != nil {
t.Fatal(err)
}
defer response.Body.Close()
got, err := io.ReadAll(response.Body)
if err != nil {
t.Fatalf("active %s over H2: %v (%d bytes)", direction, err, len(got))
}
wantStatus := http.StatusNoContent
if direction == "download" {
wantStatus = http.StatusOK
if !bytes.Equal(got, payload) {
t.Error("download bytes differ")
}
}
if response.StatusCode != wantStatus {
t.Fatalf("active %s over H2: status %d, want %d: %s", direction, response.StatusCode, wantStatus, got)
}
})
}
}

// 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.
Expand Down
4 changes: 2 additions & 2 deletions internal/proxy/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,8 +30,8 @@ type Config struct {
HandshakeTimeout time.Duration
DialTimeout 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.
// tunnel or an HTTP body transfer. Upload/download progress keeps reads
// alive; blocked writes retain an independent timeout.
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
Expand Down
32 changes: 31 additions & 1 deletion internal/proxy/http.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,8 @@ type HTTPServer struct {
transport *http.Transport
}

type httpConnContextKey struct{}

// NewHTTPServer validates cfg and creates an HTTP forward proxy.
func NewHTTPServer(cfg Config) (*HTTPServer, error) {
normalized, err := normalizeConfig(cfg)
Expand All @@ -56,6 +58,12 @@ func NewHTTPServer(cfg Config) (*HTTPServer, error) {
ReadHeaderTimeout: normalized.handshakeTimeout,
IdleTimeout: normalized.idleTimeout,
MaxHeaderBytes: 64 << 10,
ConnContext: func(ctx context.Context, conn net.Conn) context.Context {
if tracked, ok := conn.(*trackedConn); ok {
return context.WithValue(ctx, httpConnContextKey{}, tracked)
}
return ctx
},
ConnState: func(conn net.Conn, state http.ConnState) {
tracked, ok := conn.(*trackedConn)
if !ok {
Expand Down Expand Up @@ -272,7 +280,14 @@ func (s *HTTPServer) serveForward(w http.ResponseWriter, r *http.Request) {
}
}
w.WriteHeader(response.StatusCode)
_, _ = io.Copy(w, response.Body)
var body io.Reader = response.Body
if conn, ok := r.Context().Value(httpConnContextKey{}).(*trackedConn); ok {
// net/http waits for disconnects in a background read after the
// request body ends. Origin data is activity even while the response
// writer buffers small chunks and the client sends nothing more.
body = &httpResponseActivityReader{Reader: body, conn: conn}
}
_, _ = io.Copy(w, body)
for key, values := range response.Trailer {
if isHopByHopHeader(key) {
continue
Expand All @@ -281,6 +296,21 @@ func (s *HTTPServer) serveForward(w http.ResponseWriter, r *http.Request) {
}
}

type httpResponseActivityReader struct {
io.Reader
conn *trackedConn
}

func (r *httpResponseActivityReader) Read(p []byte) (int, error) {
n, err := r.Reader.Read(p)
if n > 0 {
if activityErr := r.conn.refreshReadActivity(); err == nil {
err = activityErr
}
}
return n, err
}

func validateAbsoluteTarget(target *url.URL) error {
host := target.Hostname()
if host == "" {
Expand Down
168 changes: 168 additions & 0 deletions internal/proxy/http_activity_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,168 @@
package proxy

import (
"bytes"
"errors"
"fmt"
"io"
"net"
"net/http"
"net/http/httptest"
"testing"
"time"
)

func TestHTTPForwardActiveBodiesSurviveIdleTimeout(t *testing.T) {
const idle = 200 * time.Millisecond
const chunks = 16
const interval = idle / 5
for _, direction := range []string{"download", "upload"} {
t.Run(direction, func(t *testing.T) {
chunk := bytes.Repeat([]byte("a"), 16)
origin := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if direction == "upload" {
got, err := io.ReadAll(r.Body)
if err != nil || !bytes.Equal(got, bytes.Repeat(chunk, chunks)) {
http.Error(w, "incomplete upload", http.StatusBadRequest)
return
}
w.WriteHeader(http.StatusNoContent)
return
}
w.Header().Set("Content-Length", fmt.Sprint(chunks*len(chunk)))
for range chunks {
if _, err := w.Write(chunk); err != nil {
return
}
w.(http.Flusher).Flush()
select {
case <-r.Context().Done():
return
case <-time.After(interval):
}
}
}))
defer origin.Close()
server, proxyURL, stop := startHTTPProxy(t, Config{Dialer: directDialer(), IdleTimeout: idle})
defer stop(server)
client := proxyHTTPClient(t, proxyURL, nil)
defer client.CloseIdleConnections()
method := http.MethodGet
var body io.Reader
if direction == "upload" {
method = http.MethodPost
reader, writer := io.Pipe()
defer reader.Close()
defer writer.Close()
body = reader
done := make(chan struct{})
t.Cleanup(func() {
select {
case <-done:
case <-time.After(time.Second):
t.Error("upload writer did not stop")
}
})
go func() {
defer close(done)
defer writer.Close()
for range chunks {
if _, err := writer.Write(chunk); err != nil {
return
}
time.Sleep(interval)
}
}()
}
request, err := http.NewRequest(method, origin.URL, body)
if err != nil {
t.Fatal(err)
}
if direction == "upload" {
request.ContentLength = int64(chunks * len(chunk))
}
response, err := client.Do(request)
if err != nil {
t.Fatal(err)
}
defer response.Body.Close()
got, err := io.ReadAll(response.Body)
if err != nil {
t.Fatalf("read active %s: %v (%d bytes)", direction, err, len(got))
}
wantStatus := http.StatusNoContent
if direction == "download" {
wantStatus = http.StatusOK
if !bytes.Equal(got, bytes.Repeat(chunk, chunks)) {
t.Errorf("active download = %d bytes, want %d", len(got), chunks*len(chunk))
}
}
if response.StatusCode != wantStatus {
t.Fatalf("active %s status = %d, want %d: %s", direction, response.StatusCode, wantStatus, got)
}
})
}
}

func TestHTTPActivityPreservesExternalReadDeadline(t *testing.T) {
local, peer := net.Pipe()
defer local.Close()
defer peer.Close()
conn := &trackedConn{Conn: local}
conn.setActivityTimeout(time.Hour)
// net/http uses this past deadline to stop its disconnect reader before
// hijacking a CONNECT. Neither forwarded body data nor writes may undo it.
if err := conn.SetReadDeadline(time.Now().Add(-time.Second)); err != nil {
t.Fatal(err)
}
if err := conn.refreshReadActivity(); err != nil {
t.Fatal(err)
}
done := make(chan error, 1)
go func() {
_, err := conn.Read(make([]byte, 1))
done <- err
}()
select {
case err := <-done:
var timeout net.Error
if !errors.As(err, &timeout) || !timeout.Timeout() {
t.Fatalf("read error = %v, want preserved deadline", err)
}
case <-time.After(time.Second):
t.Fatal("activity extended net/http's external read deadline")
}
}

func TestHTTPActivityDoesNotExtendBlockedWrite(t *testing.T) {
local, peer := net.Pipe()
defer local.Close()
defer peer.Close()
conn := &trackedConn{Conn: local}
conn.setActivityTimeout(100 * time.Millisecond)
done := make(chan error, 1)
go func() {
_, err := conn.Write([]byte("blocked"))
done <- err
}()
ticker := time.NewTicker(10 * time.Millisecond)
defer ticker.Stop()
limit := time.NewTimer(time.Second)
defer limit.Stop()
for {
select {
case <-ticker.C:
if err := conn.refreshReadActivity(); err != nil {
t.Fatal(err)
}
case err := <-done:
var timeout net.Error
if !errors.As(err, &timeout) || !timeout.Timeout() {
t.Fatalf("write error = %v, want stall timeout", err)
}
return
case <-limit.C:
t.Fatal("read activity kept a blocked client write alive")
}
}
}
Loading
Loading