From a106d1a8714e07fb172db747ee974569c9ace7a4 Mon Sep 17 00:00:00 2001 From: cppla Date: Sat, 3 Oct 2026 02:22:09 +0800 Subject: [PATCH 1/2] fix(tunnel): restore H2 TLS session resumption --- docs/WEB_COVER.md | 21 + internal/tunnel/web_fingerprint.go | 49 +- .../web_fingerprint_resumption_policy_test.go | 443 ++++++++++++++++++ internal/tunnel/web_h2_resumption_test.go | 430 +++++++++++++++++ 4 files changed, 942 insertions(+), 1 deletion(-) create mode 100644 internal/tunnel/web_fingerprint_resumption_policy_test.go create mode 100644 internal/tunnel/web_h2_resumption_test.go diff --git a/docs/WEB_COVER.md b/docs/WEB_COVER.md index e029590..5600c40 100644 --- a/docs/WEB_COVER.md +++ b/docs/WEB_COVER.md @@ -332,6 +332,27 @@ Chrome. The `native` rollback profile disables this client image. The H2 client separately uses the fixed `chrome-133` uTLS ClientHello profile. That describes only its TLS ClientHello; H2 settings, header order, flow control, connection reuse, payload sizes and timing retain their implementation behavior. +Source builds after v1.0.1 also support ordinary TLS 1.3 session resumption for +this profile when the caller enables a TLS session cache. An empty cache keeps +the fixed cold ClientHello shape; a valid cached ticket adds uTLS's native +`pre_shared_key` extension and binder as the last extension, as required by +[RFC 8446 section 4.2.11](https://www.rfc-editor.org/rfc/rfc8446.html#section-4.2.11). +The cache is private to one configured H2 client. A nil caller cache or +`SessionTicketsDisabled` keeps full handshakes. This does not enable 0-RTT: +every new physical connection completes TLS and starts fresh proxy authentication, +even when TLS resumes; connection-scoped proxy tickets are never inherited. + +Resumption retains the previously verified TLS session rather than repeating a +full certificate exchange or calling `VerifyPeerCertificate` again. Callers +requiring fresh per-connection certificate policy must disable session tickets; +the native profile additionally supports `VerifyConnection`. Cached tickets are +not an unconditional privacy improvement: reusing a ticket can let passive +observers correlate connections +([RFC 8446 appendix C.4](https://www.rfc-editor.org/rfc/rfc8446.html#appendix-C.4)). +The bounded loopback regression verifies actual client/server resumption, +cold/warm ClientHello policy, and fresh authentication. It is not a browser +comparison or evidence of lower classifier accuracy. + Authenticated H2 and H3 CONNECT requests explicitly suppress the Go HTTP libraries' default `User-Agent` and automatic `Accept-Encoding: gzip` values. AutoCAR does not invent browser headers until the complete request-header set diff --git a/internal/tunnel/web_fingerprint.go b/internal/tunnel/web_fingerprint.go index c852996..49dd0e3 100644 --- a/internal/tunnel/web_fingerprint.go +++ b/internal/tunnel/web_fingerprint.go @@ -6,6 +6,8 @@ import ( "errors" "fmt" "net" + "sync" + "sync/atomic" utls "github.com/refraction-networking/utls" ) @@ -87,7 +89,18 @@ func newWebH2TLSClientConn(raw net.Conn, config *tls.Config, profile Fingerprint if err != nil { return nil, err } - return &webH2UTLSConn{UConn: utls.UClient(raw, utlsConfig, utls.HelloChrome_133)}, nil + // The fixed cold profile has no pre_shared_key extension. Use a fresh + // custom copy when resumption is enabled so uTLS can populate its own + // PSK identity and binder, without changing the empty-cache wire shape. + resume := utlsConfig.ClientSessionCache != nil && !utlsConfig.SessionTicketsDisabled + hello := utls.HelloChrome_133 + if resume { + hello = utls.HelloCustom + } + return &webH2UTLSConn{ + UConn: utls.UClient(raw, utlsConfig, hello), + prepareChrome133: resume, + }, nil case FingerprintNative: return tls.Client(raw, config.Clone()), nil default: @@ -135,6 +148,7 @@ func chrome133UTLSConfig(input *tls.Config, sessionCache utls.ClientSessionCache MaxVersion: utls.VersionTLS13, SessionTicketsDisabled: input.SessionTicketsDisabled, ClientSessionCache: sessionCache, + OmitEmptyPsk: true, DynamicRecordSizingDisabled: input.DynamicRecordSizingDisabled, KeyLogWriter: input.KeyLogWriter, }, nil @@ -176,6 +190,10 @@ func cloneCertificateForUTLS(input tls.Certificate) utls.Certificate { type webH2UTLSConn struct { *utls.UConn + prepareChrome133 bool + prepareOnce sync.Once + prepareErr error + handshakeComplete atomic.Bool } // HandshakeContext enforces the application's TLS 1.3-only policy after the @@ -185,12 +203,41 @@ type webH2UTLSConn struct { // without ceasing to be that profile. Failing before any HTTP bytes are sent // keeps the policy fail-closed even if a future caller forgets a second check. func (c *webH2UTLSConn) HandshakeContext(ctx context.Context) error { + if c.prepareChrome133 { + if !c.handshakeComplete.Load() { + if err := ctx.Err(); err != nil { + return err + } + } + c.prepareOnce.Do(func() { + // ApplyPreset generates entropy and key shares: keep it inside the + // caller's handshake budget rather than in the connection factory. + // Each connection owns the extension pointers mutated by uTLS. + spec, err := utls.UTLSIdToSpec(utls.HelloChrome_133) + if err == nil { + // TLS 1.3 requires this extension to be last. OmitEmptyPsk + // suppresses it until a valid cached session is available. + spec.Extensions = append(spec.Extensions, &utls.UtlsPreSharedKeyExtension{}) + err = c.UConn.ApplyPreset(&spec) + } + c.prepareErr = err + }) + if c.prepareErr != nil { + return c.prepareErr + } + if !c.handshakeComplete.Load() { + if err := ctx.Err(); err != nil { + return err + } + } + } if err := c.UConn.HandshakeContext(ctx); err != nil { return err } if c.UConn.ConnectionState().Version != utls.VersionTLS13 { return errors.New("tunnel: web-cover HTTP/2 connection did not negotiate TLS 1.3") } + c.handshakeComplete.Store(true) return nil } diff --git a/internal/tunnel/web_fingerprint_resumption_policy_test.go b/internal/tunnel/web_fingerprint_resumption_policy_test.go new file mode 100644 index 0000000..28d2c65 --- /dev/null +++ b/internal/tunnel/web_fingerprint_resumption_policy_test.go @@ -0,0 +1,443 @@ +package tunnel + +import ( + "context" + "crypto/rand" + "crypto/tls" + "crypto/x509" + "errors" + "io" + "net" + "reflect" + "strings" + "sync/atomic" + "testing" + "time" + + utls "github.com/refraction-networking/utls" +) + +func TestChrome133ResumptionPolicyFactoryIsLazy(t *testing.T) { + for _, test := range []struct { + name string + disabled bool + cache bool + wantID utls.ClientHelloID + }{ + {name: "enabled-cache", cache: true, wantID: utls.HelloCustom}, + {name: "nil-cache", wantID: utls.HelloChrome_133}, + // Supplying a cache cannot override SessionTicketsDisabled. + {name: "tickets-disabled-with-cache", disabled: true, cache: true, wantID: utls.HelloChrome_133}, + } { + t.Run(test.name, func(t *testing.T) { + entropy := &webResumptionPolicyRand{} + raw := &webResumptionPolicyConn{} + config := &tls.Config{ + ServerName: "cover.example", + Rand: entropy, + ClientSessionCache: tls.NewLRUClientSessionCache(1), + SessionTicketsDisabled: test.disabled, + } + var cache utls.ClientSessionCache + if test.cache { + cache = utls.NewLRUClientSessionCache(1) + } + client, err := newWebH2TLSClientConn(raw, config, FingerprintChrome133, cache) + if err != nil { + t.Fatal(err) + } + chrome, ok := client.(*webH2UTLSConn) + if !ok { + t.Fatalf("Chrome client type = %T", client) + } + // Check uTLS's exported selected profile, not the wrapper's private + // preparation fields. No ClientHello has been built or sent yet. + if !reflect.DeepEqual(chrome.ClientHelloID, test.wantID) { + t.Fatalf("selected ClientHelloID = %v, want %v", chrome.ClientHelloID, test.wantID) + } + if entropy.reads.Load() != 0 { + t.Fatalf("connection factory consumed entropy: %d reads", entropy.reads.Load()) + } + raw.assertUntouched(t) + }) + } +} + +func TestChrome133ResumptionPolicyPreCancelledHandshakeRemainsReusable(t *testing.T) { + serverTLS, clientTLS := testTLSConfigs(t) + serverTLS.MinVersion = tls.VersionTLS13 + serverTLS.MaxVersion = tls.VersionTLS13 + serverTLS.NextProtos = []string{webH2ALPN} + // This fixture tests preparation/cancellation, not ticket issuance. + serverTLS.SessionTicketsDisabled = true + clientTLS.ClientSessionCache = tls.NewLRUClientSessionCache(1) + entropy := &webResumptionPolicyRand{} + clientTLS.Rand = entropy + clientPipe, serverPipe := net.Pipe() + defer clientPipe.Close() + defer serverPipe.Close() + raw := &webResumptionPolicyConn{Conn: clientPipe} + client, err := newWebH2TLSClientConn(raw, clientTLS, FingerprintChrome133, newWebH2UTLSSessionCache(clientTLS)) + if err != nil { + t.Fatal(err) + } + preCancelled, cancel := context.WithCancel(context.Background()) + cancel() + if err := client.HandshakeContext(preCancelled); !errors.Is(err, context.Canceled) { + t.Fatalf("pre-cancelled handshake error = %v, want context.Canceled", err) + } + if entropy.reads.Load() != 0 { + t.Fatalf("pre-cancelled handshake consumed entropy: %d reads", entropy.reads.Load()) + } + raw.assertUntouched(t) + + // Reuse exactly the same connection with a live context and an actual TLS + // peer. Success proves cancellation did not consume the preparation once. + ctx, stop := context.WithTimeout(context.Background(), 2*time.Second) + defer stop() + deadline, _ := ctx.Deadline() + _ = clientPipe.SetDeadline(deadline) + _ = serverPipe.SetDeadline(deadline) + serverDone := webResumptionPolicyRunHandshake(t, func() error { + return tls.Server(serverPipe, serverTLS).HandshakeContext(ctx) + }, clientPipe, serverPipe) + if err := client.HandshakeContext(ctx); err != nil { + t.Fatalf("live handshake after cancellation: %v", err) + } + select { + case err := <-serverDone: + if err != nil { + t.Fatalf("server handshake: %v", err) + } + case <-ctx.Done(): + t.Fatal("server handshake did not finish") + } + state := client.ConnectionState() + if !state.HandshakeComplete || state.Version != tls.VersionTLS13 || state.NegotiatedProtocol != webH2ALPN || len(state.VerifiedChains) == 0 { + t.Fatalf("live handshake state = %+v", state) + } + if entropy.reads.Load() == 0 || raw.reads.Load() == 0 || raw.writes.Load() == 0 { + t.Fatal("live handshake did not consume entropy and exchange TLS records") + } + reads, writes, randomReads := raw.reads.Load(), raw.writes.Load(), entropy.reads.Load() + if err := client.HandshakeContext(ctx); err != nil { + t.Fatalf("repeat completed handshake: %v", err) + } + if err := client.HandshakeContext(preCancelled); err != nil { + t.Fatalf("completed handshake with cancelled context: %v", err) + } + if raw.reads.Load() != reads || raw.writes.Load() != writes || entropy.reads.Load() != randomReads { + t.Fatal("repeat completed handshake performed preparation or socket I/O again") + } +} + +func TestChrome133ResumptionPolicyConcurrentPreCancellationDoesNotWaitForPeer(t *testing.T) { + entropy := &webResumptionPolicyRand{} + config := &tls.Config{ + ServerName: "cover.example", + Rand: entropy, + ClientSessionCache: tls.NewLRUClientSessionCache(1), + } + clientPipe, serverPipe := net.Pipe() + defer clientPipe.Close() + defer serverPipe.Close() + raw := &webResumptionPolicyConn{Conn: clientPipe, writeStarted: make(chan struct{}, 1)} + client, err := newWebH2TLSClientConn(raw, config, FingerprintChrome133, newWebH2UTLSSessionCache(config)) + if err != nil { + t.Fatal(err) + } + firstContext, cancelFirst := context.WithTimeout(context.Background(), 2*time.Second) + defer cancelFirst() + deadline, _ := firstContext.Deadline() + _ = clientPipe.SetDeadline(deadline) + firstDone := webResumptionPolicyRunHandshake(t, func() error { + return client.HandshakeContext(firstContext) + }, clientPipe, serverPipe) + select { + case <-raw.writeStarted: + // No peer reads this pipe, so the first TLS write remains blocked. + case <-firstContext.Done(): + t.Fatal("first handshake did not reach its TLS write") + } + reads, writes, randomReads := raw.reads.Load(), raw.writes.Load(), entropy.reads.Load() + preCancelled, cancel := context.WithCancel(context.Background()) + cancel() + secondDone := webResumptionPolicyRunHandshake(t, func() error { + return client.HandshakeContext(preCancelled) + }, clientPipe, serverPipe) + secondReturned := false + select { + case err := <-secondDone: + secondReturned = true + if !errors.Is(err, context.Canceled) { + t.Errorf("second pre-cancelled handshake error = %v", err) + } + case <-time.After(500 * time.Millisecond): + t.Error("pre-cancelled handshake waited for the blocked first handshake") + } + if raw.reads.Load() != reads || raw.writes.Load() != writes || entropy.reads.Load() != randomReads || raw.closes.Load() != 0 { + t.Error("second pre-cancelled handshake touched entropy or the raw connection") + } + select { + case err := <-firstDone: + t.Errorf("first handshake unexpectedly finished without a peer: %v", err) + firstDone = nil + default: + } + // Unblock and join both workers even if the cancellation assertion fails. + cancelFirst() + _ = clientPipe.Close() + _ = serverPipe.Close() + if firstDone != nil { + select { + case <-firstDone: + case <-time.After(time.Second): + t.Error("first handshake did not join after cancellation") + } + } + if !secondReturned { + select { + case <-secondDone: + case <-time.After(time.Second): + t.Error("second handshake did not join after cancellation") + } + } +} + +func TestChrome133ResumptionPolicyPreparationErrorIsStableAndDoesNotTouchSocket(t *testing.T) { + fault := errors.New("resumption-policy entropy fault") + entropy := &webResumptionPolicyRand{fault: fault} + raw := &webResumptionPolicyConn{} + config := &tls.Config{ + ServerName: "cover.example", + Rand: entropy, + ClientSessionCache: tls.NewLRUClientSessionCache(1), + } + client, err := newWebH2TLSClientConn(raw, config, FingerprintChrome133, newWebH2UTLSSessionCache(config)) + if err != nil { + t.Fatalf("factory eagerly prepared the handshake: %v", err) + } + if entropy.reads.Load() != 0 { + t.Fatal("factory read failing entropy source") + } + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + first := client.HandshakeContext(ctx) + if first == nil || !strings.Contains(first.Error(), fault.Error()) { + t.Fatalf("preparation error = %v, want entropy fault", first) + } + reads := entropy.reads.Load() + if reads == 0 { + t.Fatal("preparation did not reach the failing entropy source") + } + for attempt := 0; attempt < 3; attempt++ { + if err := client.HandshakeContext(ctx); err == nil || err.Error() != first.Error() { + t.Fatalf("repeat preparation error = %v, want stable %v", err, first) + } + if entropy.reads.Load() != reads { + t.Fatal("repeat failed preparation consumed entropy again") + } + raw.assertUntouched(t) + } + if client.ConnectionState().HandshakeComplete { + t.Fatal("failed preparation marked the handshake complete") + } +} + +func TestChrome133ResumptionPolicyStillVerifiesPeer(t *testing.T) { + for _, name := range []string{"untrusted-roots", "wrong-server-name"} { + t.Run(name, func(t *testing.T) { + serverTLS, clientTLS := testTLSConfigs(t) + serverTLS.MinVersion = tls.VersionTLS13 + serverTLS.MaxVersion = tls.VersionTLS13 + serverTLS.NextProtos = []string{webH2ALPN} + clientTLS.ClientSessionCache = tls.NewLRUClientSessionCache(1) + if name == "untrusted-roots" { + clientTLS.RootCAs = x509.NewCertPool() + } else { + clientTLS.ServerName = "wrong.example" + } + // TCP buffering permits the verification alert and server flight + // to cross without net.Pipe's synchronous-write deadlock. + clientPipe, serverPipe := webResumptionPolicyTCPPair(t) + defer clientPipe.Close() + defer serverPipe.Close() + raw := &webResumptionPolicyConn{Conn: clientPipe} + client, err := newWebH2TLSClientConn(raw, clientTLS, FingerprintChrome133, newWebH2UTLSSessionCache(clientTLS)) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + deadline, _ := ctx.Deadline() + _ = clientPipe.SetDeadline(deadline) + _ = serverPipe.SetDeadline(deadline) + serverDone := webResumptionPolicyRunHandshake(t, func() error { + return tls.Server(serverPipe, serverTLS).HandshakeContext(ctx) + }, clientPipe, serverPipe) + err = client.HandshakeContext(ctx) + if name == "untrusted-roots" { + var untrusted x509.UnknownAuthorityError + if !errors.As(err, &untrusted) { + t.Fatalf("untrusted-root handshake error = %v", err) + } + } else { + var wrongName x509.HostnameError + if !errors.As(err, &wrongName) { + t.Fatalf("wrong-name handshake error = %v", err) + } + } + if client.ConnectionState().HandshakeComplete { + t.Fatal("certificate verification failure completed the handshake") + } + if raw.reads.Load() == 0 || raw.writes.Load() == 0 { + t.Fatal("verification test did not exchange TLS records with its peer") + } + _ = clientPipe.Close() + _ = serverPipe.Close() + select { + case <-serverDone: + case <-time.After(time.Second): + t.Fatal("rejected TLS peer did not join") + } + }) + } +} + +func webResumptionPolicyTCPPair(t *testing.T) (net.Conn, net.Conn) { + t.Helper() + listener, err := net.ListenTCP("tcp", &net.TCPAddr{IP: net.IPv4(127, 0, 0, 1)}) + if err != nil { + t.Fatal(err) + } + defer listener.Close() + _ = listener.SetDeadline(time.Now().Add(time.Second)) + client, err := net.DialTimeout("tcp", listener.Addr().String(), time.Second) + if err != nil { + t.Fatal(err) + } + server, err := listener.Accept() + if err != nil { + _ = client.Close() + t.Fatal(err) + } + return client, server +} + +func webResumptionPolicyRunHandshake(t *testing.T, run func() error, clientPipe, serverPipe net.Conn) <-chan error { + t.Helper() + done := make(chan error, 1) + joined := make(chan struct{}) + // A buffered result can arrive before the worker's deferred cleanup. + // Wait for a separate completion signal on success and Fatal paths. + t.Cleanup(func() { + _ = clientPipe.Close() + _ = serverPipe.Close() + select { + case <-joined: + case <-time.After(2 * time.Second): + t.Error("TLS handshake worker did not join after closing both pipe ends") + } + }) + go func() { + defer close(joined) + defer close(done) + done <- run() + }() + return done +} + +type webResumptionPolicyRand struct { + reads atomic.Int64 + fault error +} + +func (r *webResumptionPolicyRand) Read(p []byte) (int, error) { + r.reads.Add(1) + if r.fault != nil { + return 0, r.fault + } + return rand.Read(p) +} + +// A nil underlying connection fails fast on unexpected I/O, without touching +// a socket. A real pipe lets the same counters observe a bounded TLS exchange. +type webResumptionPolicyConn struct { + net.Conn + reads atomic.Int64 + writes atomic.Int64 + closes atomic.Int64 + writeStarted chan struct{} +} + +func (c *webResumptionPolicyConn) Read(p []byte) (int, error) { + c.reads.Add(1) + if c.Conn == nil { + return 0, io.ErrClosedPipe + } + return c.Conn.Read(p) +} + +func (c *webResumptionPolicyConn) Write(p []byte) (int, error) { + c.writes.Add(1) + if c.writeStarted != nil { + select { + case c.writeStarted <- struct{}{}: + default: + } + } + if c.Conn == nil { + return 0, io.ErrClosedPipe + } + return c.Conn.Write(p) +} + +func (c *webResumptionPolicyConn) Close() error { + c.closes.Add(1) + if c.Conn == nil { + return nil + } + return c.Conn.Close() +} + +func (c *webResumptionPolicyConn) LocalAddr() net.Addr { + if c.Conn == nil { + return &net.TCPAddr{IP: net.IPv4(127, 0, 0, 1), Port: 1} + } + return c.Conn.LocalAddr() +} + +func (c *webResumptionPolicyConn) RemoteAddr() net.Addr { + if c.Conn == nil { + return &net.TCPAddr{IP: net.IPv4(127, 0, 0, 1), Port: 443} + } + return c.Conn.RemoteAddr() +} + +func (c *webResumptionPolicyConn) SetDeadline(deadline time.Time) error { + if c.Conn == nil { + return nil + } + return c.Conn.SetDeadline(deadline) +} + +func (c *webResumptionPolicyConn) SetReadDeadline(deadline time.Time) error { + if c.Conn == nil { + return nil + } + return c.Conn.SetReadDeadline(deadline) +} + +func (c *webResumptionPolicyConn) SetWriteDeadline(deadline time.Time) error { + if c.Conn == nil { + return nil + } + return c.Conn.SetWriteDeadline(deadline) +} + +func (c *webResumptionPolicyConn) assertUntouched(t *testing.T) { + t.Helper() + if c.reads.Load() != 0 || c.writes.Load() != 0 || c.closes.Load() != 0 { + t.Fatalf("raw connection touched: reads=%d writes=%d closes=%d", c.reads.Load(), c.writes.Load(), c.closes.Load()) + } +} diff --git a/internal/tunnel/web_h2_resumption_test.go b/internal/tunnel/web_h2_resumption_test.go new file mode 100644 index 0000000..75ee800 --- /dev/null +++ b/internal/tunnel/web_h2_resumption_test.go @@ -0,0 +1,430 @@ +package tunnel + +import ( + "context" + "crypto/tls" + "crypto/x509" + "encoding/binary" + "errors" + "fmt" + "net" + "net/http" + "reflect" + "slices" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/cppla/autocar/internal/security" + "github.com/cppla/autocar/internal/transport" + utls "github.com/refraction-networking/utls" + "golang.org/x/net/http2" +) + +type webH2ResumptionTLSCache struct { + inner tls.ClientSessionCache + hits atomic.Int64 + puts atomic.Int64 +} + +func (c *webH2ResumptionTLSCache) Get(key string) (*tls.ClientSessionState, bool) { + state, ok := c.inner.Get(key) + if ok { + c.hits.Add(1) + } + return state, ok +} + +func (c *webH2ResumptionTLSCache) Put(key string, state *tls.ClientSessionState) { + if state != nil { + c.puts.Add(1) + } + c.inner.Put(key, state) +} + +type webH2ResumptionUTLSCache struct { + inner utls.ClientSessionCache + hits atomic.Int64 + puts atomic.Int64 +} + +func (c *webH2ResumptionUTLSCache) Get(key string) (*utls.ClientSessionState, bool) { + state, ok := c.inner.Get(key) + if ok { + c.hits.Add(1) + } + return state, ok +} + +func (c *webH2ResumptionUTLSCache) Put(key string, state *utls.ClientSessionState) { + if state != nil { + c.puts.Add(1) + } + c.inner.Put(key, state) +} + +// Capture raw TLS records only in memory, to parse the actual first ClientHello +// received by the trusted loopback TLS server. Never log ticket or auth bytes. +type webH2ResumptionWireConn struct { + net.Conn + mu sync.Mutex + data []byte +} + +func (c *webH2ResumptionWireConn) Read(p []byte) (int, error) { + n, err := c.Conn.Read(p) + if n > 0 { + c.mu.Lock() + if len(c.data)+n <= 64<<10 { + c.data = append(c.data, p[:n]...) + } + c.mu.Unlock() + } + return n, err +} + +func (c *webH2ResumptionWireConn) clientHello() (parsedClientHello, error) { + c.mu.Lock() + defer c.mu.Unlock() + if len(c.data) < 5 { + return parsedClientHello{}, errors.New("missing real TLS record") + } + length := 5 + int(binary.BigEndian.Uint16(c.data[3:5])) + if length > len(c.data) { + return parsedClientHello{}, errors.New("incomplete real ClientHello record") + } + return parseTLSClientHello(c.data[:length]) +} + +type webH2ResumptionPhysicalResult struct { + sequence int + state tls.ConnectionState + hello parsedClientHello + err error +} + +type webH2ResumptionRequest struct { + physical int + bearerBytes int +} + +func webH2ResumptionTLSConfigs(t *testing.T) (*tls.Config, *tls.Config) { + t.Helper() + // A DNS SAN exercises SNI as well as trusted hostname validation, while + // the actual socket remains loopback and does not require DNS lookup. + const host = "resumption-cover.example" + certPEM, keyPEM, err := security.GenerateSelfSignedCertificate(security.CertificateOptions{Hosts: []string{host}, ValidFor: time.Hour}) + if err != nil { + t.Fatal(err) + } + certificate, err := tls.X509KeyPair(certPEM, keyPEM) + if err != nil { + t.Fatal(err) + } + roots := x509.NewCertPool() + if !roots.AppendCertsFromPEM(certPEM) { + t.Fatal("could not trust test DNS certificate") + } + return &tls.Config{Certificates: []tls.Certificate{certificate}}, &tls.Config{RootCAs: roots, ServerName: host} +} + +// The public client talks to the production handler over real TCP/TLS/H2. +// ServeConn lets this fixture observe the handshake without substituting auth +// or a fake TLS state; its context carries the same exact-connection bindings +// as ListenWebH2. Each accepted socket and worker is independently owned. +func webH2ResumptionOrigin(t *testing.T, serverTLS *tls.Config) (string, <-chan webH2ResumptionPhysicalResult, <-chan webH2ResumptionRequest, <-chan struct{}, *atomic.Int64) { + t.Helper() + listener, err := net.ListenTCP("tcp", &net.TCPAddr{IP: net.ParseIP("127.0.0.1")}) + if err != nil { + t.Fatal(err) + } + var dialCalls atomic.Int64 + core, err := newServerCoreWithAdmission(webTestToken, transport.DialFunc(func(context.Context, string, string) (net.Conn, error) { + dialCalls.Add(1) + return nil, errors.New("private resumption test destination failure") + }), 0, 0, 0, 0, nil) + if err != nil { + _ = listener.Close() + t.Fatal(err) + } + verifier, err := newWebAuthVerifier(mustWebAuthKey(t, webTestToken), nil, 0) + if err != nil { + _ = listener.Close() + t.Fatal(err) + } + handler := &webTunnelHandler{core: core, auth: verifier, cover: http.NotFoundHandler()} + config := serverTLS.Clone() + config.MinVersion, config.MaxVersion = tls.VersionTLS13, tls.VersionTLS13 + config.NextProtos = []string{webH2ALPN} + results := make(chan webH2ResumptionPhysicalResult, 2) + requests := make(chan webH2ResumptionRequest, 4) + handlerDone := make(chan struct{}, 4) + done := make(chan struct{}) + var mu sync.Mutex + closed := false + owned := make(map[net.Conn]struct{}) + go func() { + defer close(done) + var wg sync.WaitGroup + defer wg.Wait() + for sequence := range 2 { + _ = listener.SetDeadline(time.Now().Add(3 * time.Second)) + raw, acceptErr := listener.AcceptTCP() + if acceptErr != nil { + results <- webH2ResumptionPhysicalResult{sequence: sequence, err: acceptErr} + return + } + mu.Lock() + if closed { + mu.Unlock() + _ = raw.Close() + return + } + owned[raw] = struct{}{} + mu.Unlock() + wg.Add(1) + go func() { + defer wg.Done() + defer raw.Close() + defer func() { mu.Lock(); delete(owned, raw); mu.Unlock() }() + _ = raw.SetDeadline(time.Now().Add(4 * time.Second)) + wire := &webH2ResumptionWireConn{Conn: raw} + tlsConn := tls.Server(wire, config) + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + handshakeErr := tlsConn.HandshakeContext(ctx) + cancel() + if handshakeErr != nil { + results <- webH2ResumptionPhysicalResult{sequence: sequence, err: handshakeErr} + return + } + hello, helloErr := wire.clientHello() + if helloErr != nil { + results <- webH2ResumptionPhysicalResult{sequence: sequence, err: helloErr} + return + } + ctx = context.WithValue(context.Background(), webTLSConnectionContextKey{}, tlsConn) + ctx = context.WithValue(ctx, webServerConnectionAuthContextKey{}, newWebServerConnectionAuth(raw.Close)) + observed := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + // Four actual requests fit independently of the test consumer. + // Unexpected canceled requests must not strand a handler here. + defer func() { + select { + case handlerDone <- struct{}{}: + default: + } + }() + select { + case requests <- webH2ResumptionRequest{physical: sequence, bearerBytes: len(r.Header.Get("Proxy-Authorization"))}: + case <-r.Context().Done(): + return + } + handler.ServeHTTP(w, r) + }) + (&http2.Server{}).ServeConn(tlsConn, &http2.ServeConnOpts{Context: ctx, Handler: observed}) + results <- webH2ResumptionPhysicalResult{sequence: sequence, state: tlsConn.ConnectionState(), hello: hello} + }() + } + }() + t.Cleanup(func() { + _ = listener.Close() + mu.Lock() + closed = true + for conn := range owned { + _ = conn.Close() + } + mu.Unlock() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Error("real resumption Accept/TLS/H2 workers did not join") + } + if got := len(core.sem); got != 0 { + t.Errorf("destination admission retained %d slots", got) + } + }) + return listener.Addr().String(), results, requests, handlerDone, &dialCalls +} + +func webH2ResumptionCheckChromeHello(t *testing.T, hello parsedClientHello, warm bool) { + t.Helper() + wantCiphers := []uint16{0x1301, 0x1302, 0x1303, 0xc02b, 0xc02f, 0xc02c, 0xc030, 0xcca9, 0xcca8, 0xc013, 0xc014, 0x009c, 0x009d, 0x002f, 0x0035} + if got := withoutGREASE(hello.cipherSuites); !reflect.DeepEqual(got, wantCiphers) { + t.Errorf("Chrome cold/warm cipher order changed: %#v", got) + } + // This exact non-GREASE membership is the original pinned Chrome-133 + // preset. Order is deliberately not frozen: the preset shuffles it. + wantExtensions := []uint16{0, 5, 10, 11, 13, 16, 18, 23, 27, 35, 43, 45, 51, 17613, 0xfe0d, 0xff01} + if warm { + wantExtensions = append(wantExtensions, 41) + } + gotExtensions := withoutGREASE(hello.extensions) + slices.Sort(gotExtensions) + slices.Sort(wantExtensions) + if !reflect.DeepEqual(gotExtensions, wantExtensions) { + t.Errorf("Chrome extension membership = %#v, want %#v", gotExtensions, wantExtensions) + } + if !containsGREASE(hello.cipherSuites) || !containsGREASE(hello.extensions) { + t.Error("Chrome ClientHello lost GREASE") + } + if !reflect.DeepEqual(hello.alpn, []string{webH2ALPN, webHTTP11ALPN}) { + t.Errorf("Chrome offered ALPN = %q", hello.alpn) + } +} + +func TestWebH2RealTLSResumptionAndFreshAuthentication(t *testing.T) { + for _, test := range []struct { + name string + profile FingerprintProfile + cacheEnabled bool + ticketsDisabled bool + wantResume bool + }{ + {name: "chrome_enabled", profile: FingerprintChrome133, cacheEnabled: true, wantResume: true}, + {name: "chrome_nil_cache", profile: FingerprintChrome133}, + {name: "chrome_tickets_disabled", profile: FingerprintChrome133, cacheEnabled: true, ticketsDisabled: true}, + {name: "native_positive", profile: FingerprintNative, cacheEnabled: true, wantResume: true}, + } { + t.Run(test.name, func(t *testing.T) { + serverTLS, clientTLS := webH2ResumptionTLSConfigs(t) + address, results, requests, handlerDone, dialCalls := webH2ResumptionOrigin(t, serverTLS) + config := clientTLS.Clone() + config.ClientSessionCache = nil + config.SessionTicketsDisabled = test.ticketsDisabled + standardCache := &webH2ResumptionTLSCache{inner: tls.NewLRUClientSessionCache(4)} + if test.cacheEnabled { + config.ClientSessionCache = standardCache + } + client, err := NewWebH2Client(WebH2ClientConfig{ServerAddress: address, Token: webTestToken, TLSConfig: config, FingerprintProfile: test.profile}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = client.Close() }) + var chromeCache *webH2ResumptionUTLSCache + if test.profile == FingerprintChrome133 { + if test.wantResume { + if client.utlsSessionCache == nil { + t.Fatal("enabled cache policy did not create a uTLS cache") + } + chromeCache = &webH2ResumptionUTLSCache{inner: client.utlsSessionCache} + client.utlsSessionCache = chromeCache + } else if client.utlsSessionCache != nil { + t.Fatal("nil-cache/tickets-disabled policy created a uTLS cache") + } + } + var physical [2]*webH2ClientSession + var states [2]tls.ConnectionState + var authKeys [2]webSessionKey + for request := range 4 { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + conn, dialErr := client.DialContext(ctx, "tcp", fmt.Sprintf("resumption-target%d.example:443", request)) + cancel() + if conn != nil { + _ = conn.Close() + t.Error("authenticated destination rejection returned a connection") + } + var rejection *WebConnectError + if !errors.As(dialErr, &rejection) || rejection.StatusCode != http.StatusBadGateway { + t.Fatalf("real authenticated HTTP/2 refusal: %v", dialErr) + } + select { + case <-handlerDone: + case <-time.After(2 * time.Second): + t.Fatal("authenticated handler did not return") + } + index := request / 2 + client.mu.Lock() + current := client.current + ready := current != nil && current.authState == webH2ClientAuthReady && current.auth != nil + if ready { + authKeys[index] = current.auth.key + } + client.mu.Unlock() + if !ready { + t.Fatal("verified response did not retain ready connection authentication") + } + if request%2 == 0 { + physical[index] = current + states[index] = current.conn.ConnectionState() + } else if current != physical[index] { + t.Fatal("continuation unexpectedly changed physical connection") + } + select { + case observation := <-requests: + if observation.physical != index { + t.Errorf("request %d used physical %d, want %d", request, observation.physical, index) + } + if request%2 == 0 && observation.bearerBytes <= 41 { + t.Errorf("new physical %d inherited a short auth ticket", index) + } + if request%2 == 1 && observation.bearerBytes != 41 { + t.Errorf("physical %d continuation bearer bytes=%d, want 41", index, observation.bearerBytes) + } + case <-time.After(2 * time.Second): + t.Fatal("real request observation did not arrive") + } + if request == 1 { + current.h2.SetDoNotReuse() + } + } + if physical[0] == physical[1] || authKeys[0] == authKeys[1] { + t.Error("replacement inherited physical connection or authentication key") + } + if got := dialCalls.Load(); got != 4 { + t.Errorf("authenticated destination calls=%d, want 4", got) + } + for sequence, state := range states { + wantResume := sequence == 1 && test.wantResume + if state.DidResume != wantResume { + t.Errorf("client physical %d DidResume=%t, want %t", sequence, state.DidResume, wantResume) + } + if state.Version != tls.VersionTLS13 || state.NegotiatedProtocol != webH2ALPN || len(state.VerifiedChains) == 0 || state.ServerName != config.ServerName { + t.Errorf("client physical %d lost verified TLS 1.3/h2 state", sequence) + } + } + _ = client.Close() + for range 2 { + select { + case result := <-results: + if result.err != nil { + t.Fatal(result.err) + } + warm := result.sequence == 1 && test.wantResume + if result.state.DidResume != warm { + t.Errorf("server physical %d DidResume=%t, want %t", result.sequence, result.state.DidResume, warm) + } + if result.state.Version != tls.VersionTLS13 || result.state.NegotiatedProtocol != webH2ALPN || result.state.ServerName != config.ServerName { + t.Errorf("server physical %d lost TLS 1.3/h2 state", result.sequence) + } + if containsUint16(result.hello.extensions, 42) { + t.Error("ClientHello unexpectedly offered TLS early_data") + } + if got := containsUint16(result.hello.extensions, 41); got != warm { + t.Errorf("physical %d actual wire PSK present=%t, want %t", result.sequence, got, warm) + } + if warm && result.hello.extensions[len(result.hello.extensions)-1] != 41 { + t.Error("warm ClientHello PSK extension is not last") + } + if test.profile == FingerprintChrome133 { + webH2ResumptionCheckChromeHello(t, result.hello, warm) + } + case <-time.After(2 * time.Second): + t.Fatal("real physical HTTP/2 worker did not return") + } + } + if chromeCache != nil { + if chromeCache.hits.Load() < 1 || chromeCache.puts.Load() < 1 { + t.Error("Chrome resumed without a real stored/loaded ticket control") + } + t.Logf("real Chrome TLS tickets puts=%d hits=%d", chromeCache.puts.Load(), chromeCache.hits.Load()) + } else if test.profile == FingerprintNative { + if standardCache.hits.Load() < 1 || standardCache.puts.Load() < 1 { + t.Error("native resumption cache positive control failed") + } + t.Logf("real native TLS tickets puts=%d hits=%d", standardCache.puts.Load(), standardCache.hits.Load()) + } else if standardCache.hits.Load() != 0 || standardCache.puts.Load() != 0 { + t.Error("Chrome caller cache policy leaked into standard cache use") + } + }) + } +} From 97dfb7c81fe7e68df686dd39b9f2c7d335127df5 Mon Sep 17 00:00:00 2001 From: cppla Date: Sat, 3 Oct 2026 02:47:34 +0800 Subject: [PATCH 2/2] fix(tunnel): bound PSK HelloRetryRequest compatibility fallback --- docs/WEB_COVER.md | 7 + internal/tunnel/web_fingerprint.go | 35 ++ internal/tunnel/web_h2_client.go | 26 +- .../tunnel/web_h2_resumption_fallback_test.go | 488 ++++++++++++++++ internal/tunnel/web_h2_resumption_hrr_test.go | 523 ++++++++++++++++++ 5 files changed, 1078 insertions(+), 1 deletion(-) create mode 100644 internal/tunnel/web_h2_resumption_fallback_test.go create mode 100644 internal/tunnel/web_h2_resumption_hrr_test.go diff --git a/docs/WEB_COVER.md b/docs/WEB_COVER.md index 5600c40..defb3f5 100644 --- a/docs/WEB_COVER.md +++ b/docs/WEB_COVER.md @@ -341,6 +341,13 @@ The cache is private to one configured H2 client. A nil caller cache or `SessionTicketsDisabled` keeps full handshakes. This does not enable 0-RTT: every new physical connection completes TLS and starts fresh proxy authentication, even when TLS resumes; connection-scoped proxy tickets are never inherited. +The pinned uTLS implementation cannot rebuild a populated PSK after a TLS 1.3 +HelloRetryRequest. For this narrowly recognized library limitation, the H2 client +closes the failed socket and retries once on a fresh connection without a ticket, +within the same remaining initialization timeout. Certificate/hostname checks, +TLS 1.3 and h2 are still mandatory; unrelated TLS failures are not retried. +The library may invalidate the failed cached ticket. This compatibility fallback +is a full handshake, not successful HRR resumption or a browser-equivalence claim. Resumption retains the previously verified TLS session rather than repeating a full certificate exchange or calling `VerifyPeerCertificate` again. Callers diff --git a/internal/tunnel/web_fingerprint.go b/internal/tunnel/web_fingerprint.go index 49dd0e3..a75048b 100644 --- a/internal/tunnel/web_fingerprint.go +++ b/internal/tunnel/web_fingerprint.go @@ -1,6 +1,7 @@ package tunnel import ( + "bytes" "context" "crypto/tls" "errors" @@ -76,6 +77,37 @@ type webH2TLSClientConn interface { ConnectionState() tls.ConnectionState } +var errWebH2PSKHelloRetryRequest = errors.New("tunnel: cached H2 TLS session requires an unsupported HelloRetryRequest") + +// This diagnostic is private to the pinned uTLS v1.8.2 implementation. It has +// no exported error type; recognize its exact text AND completed HRR state, +// never arbitrary certificate, entropy or I/O errors with similar text. +const webH2UnsupportedPSKHRR = "uTLS does not support reprocessing of PSK key triggered by HelloRetryRequest" + +func webH2UnsupportedPSKHelloRetryRequest(err error, state utls.PubClientHandshakeState) bool { + if err == nil || err.Error() != webH2UnsupportedPSKHRR || state.Hello == nil || state.ServerHello == nil { + return false + } + hello, server := state.Hello, state.ServerHello + hrrRandom := [...]byte{ // RFC 8446 section 4.1.3. + 0xcf, 0x21, 0xad, 0x74, 0xe5, 0x9a, 0x61, 0x11, + 0xbe, 0x1d, 0x8c, 0x02, 0x1e, 0x65, 0xb8, 0x91, + 0xc2, 0xa2, 0x11, 0x16, 0x7a, 0xbb, 0x8c, 0x5e, + 0x07, 0x9e, 0x09, 0xe2, 0xc8, 0xa8, 0x33, 0x9c, + } + if server.SupportedVersion != utls.VersionTLS13 || !bytes.Equal(server.Random, hrrRandom[:]) || + len(hello.PskIdentities) == 0 || len(hello.PskBinders) != len(hello.PskIdentities) { + return false + } + if server.SelectedGroup != 0 { + // uTLS reaches the unsupported-PSK diagnostic only after generating + // and installing the requested share. This excludes an entropy error + // with the same text during HRR key generation. + return len(hello.KeyShares) == 1 && hello.KeyShares[0].Group == server.SelectedGroup + } + return len(server.Cookie) > 0 +} + func newWebH2TLSClientConn(raw net.Conn, config *tls.Config, profile FingerprintProfile, sessionCache utls.ClientSessionCache) (webH2TLSClientConn, error) { if raw == nil { return nil, errors.New("tunnel: nil web-cover HTTP/2 TCP connection") @@ -232,6 +264,9 @@ func (c *webH2UTLSConn) HandshakeContext(ctx context.Context) error { } } if err := c.UConn.HandshakeContext(ctx); err != nil { + if c.prepareChrome133 && err.Error() == webH2UnsupportedPSKHRR && webH2UnsupportedPSKHelloRetryRequest(err, c.UConn.HandshakeState) { + return fmt.Errorf("%w: %w", errWebH2PSKHelloRetryRequest, err) + } return err } if c.UConn.ConnectionState().Version != utls.VersionTLS13 { diff --git a/internal/tunnel/web_h2_client.go b/internal/tunnel/web_h2_client.go index 6af17a1..dcb2378 100644 --- a/internal/tunnel/web_h2_client.go +++ b/internal/tunnel/web_h2_client.go @@ -497,6 +497,30 @@ func (c *WebH2Client) openSession(ctx context.Context) (*webH2ClientSession, err // contexts through initialization, not just through the TLS handshake. initializationCtx, initializationCancel := context.WithTimeout(dialCtx, c.handshakeTimeout) defer initializationCancel() + session, err := c.initializeSession(initializationCtx, raw, c.utlsSessionCache) + if !errors.Is(err, errWebH2PSKHelloRetryRequest) { + return session, err + } + // Pinned uTLS cannot rebuild a populated PSK after HRR. The failed + // attempt has already closed its raw socket and joined its watcher. + // Retry once on a fresh socket without tickets, with the SAME remaining + // initialization budget. No HTTP or proxy authentication was sent yet. + if cause := context.Cause(initializationCtx); cause != nil { + return nil, fmt.Errorf("tunnel: web-cover HTTP/2 TLS handshake: %w", cause) + } + raw, err = c.dialer.DialContext(initializationCtx, "tcp", c.address) + if err != nil { + if cause := context.Cause(initializationCtx); cause != nil { + err = cause + } + return nil, fmt.Errorf("tunnel: redial web-cover HTTP/2 server after TLS retry: %w", err) + } + return c.initializeSession(initializationCtx, raw, nil) +} + +// initializeSession owns one physical attempt. Its immutable raw parameter +// prevents a retired watcher's closure from affecting a replacement socket. +func (c *WebH2Client) initializeSession(initializationCtx context.Context, raw net.Conn, sessionCache utls.ClientSessionCache) (*webH2ClientSession, error) { rawClosed := make(chan struct{}) stopRawClose := context.AfterFunc(initializationCtx, func() { _ = raw.Close() @@ -508,7 +532,7 @@ func (c *WebH2Client) openSession(ctx context.Context) (*webH2ClientSession, err <-rawClosed } }() - tlsConn, err := newWebH2TLSClientConn(raw, c.tlsConfig, c.fingerprint, c.utlsSessionCache) + tlsConn, err := newWebH2TLSClientConn(raw, c.tlsConfig, c.fingerprint, sessionCache) if err != nil { _ = raw.Close() return nil, err diff --git a/internal/tunnel/web_h2_resumption_fallback_test.go b/internal/tunnel/web_h2_resumption_fallback_test.go new file mode 100644 index 0000000..9354c69 --- /dev/null +++ b/internal/tunnel/web_h2_resumption_fallback_test.go @@ -0,0 +1,488 @@ +package tunnel + +import ( + "context" + "crypto/tls" + "crypto/x509" + "errors" + "io" + "net" + "net/http" + "strings" + "sync" + "sync/atomic" + "syscall" + "testing" + "time" + + "github.com/cppla/autocar/internal/transport" +) + +func TestWebH2ResumptionFallbackInitializationSharesBudgetAndCancellation(t *testing.T) { + for _, stage := range []string{"redial", "h2_preface"} { + t.Run(stage, func(t *testing.T) { + for _, mode := range []string{"handshake_timeout", "caller_cancel", "client_close"} { + t.Run(mode, func(t *testing.T) { + fixture := webH2FallbackPrime(t) + const budget = 1500 * time.Millisecond + fixture.client.handshakeTimeout = budget + warm := webH2FallbackInstallWarm(t, fixture) + atFallback := make(chan context.Context, 1) + var fresh *webH2InitializationWire + var serverDone <-chan error + var verified atomic.Pointer[webH2InitializationWire] + fixture.client.tlsConfig.VerifyPeerCertificate = func(_ [][]byte, _ [][]*x509.Certificate) error { + if wire := verified.Load(); wire != nil { + wire.verified.Store(true) + } + return nil + } + warm.fresh = func(ctx context.Context) (net.Conn, error) { + if stage == "redial" { + atFallback <- ctx + <-ctx.Done() + return nil, ctx.Err() + } + raw, peer := net.Pipe() + freshWire := &webH2InitializationWire{Conn: raw, initializing: make(chan struct{}), closed: make(chan struct{})} + verified.Store(freshWire) + serverTLS := fixture.serverTLS.Clone() + serverTLS.MinVersion, serverTLS.MaxVersion = tls.VersionTLS13, tls.VersionTLS13 + serverTLS.NextProtos = []string{webH2ALPN} + serverTLS.SessionTicketsDisabled = true + serverCtx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + deadline, _ := serverCtx.Deadline() + _ = peer.SetDeadline(deadline) + serverDone = webH2FallbackWorker(t, func() error { + return tls.Server(peer, serverTLS).HandshakeContext(serverCtx) + }, func() { cancel(); _ = raw.Close(); _ = peer.Close() }) + atFallback <- ctx + return freshWire, nil + } + caller, cancel := context.WithCancelCause(context.Background()) + defer cancel(context.Canceled) + openDone := webH2FallbackOpenWorker(t, fixture.client, caller, warm) + firstRead := webH2FallbackReleaseWarm(t, warm, budget/3) + var fallbackCtx context.Context + select { + case fallbackCtx = <-atFallback: + case err := <-openDone: + t.Fatalf("initializer did not enter real HRR fallback: %v", err) + case <-time.After(2 * time.Second): + t.Fatal("fallback redial was not reached") + } + deadline, ok := fallbackCtx.Deadline() + // These causal timestamps bracket creation of the ORIGINAL + // initialization context, independent of scheduler delays. + // Waiting budget/3 before releasing the real HRR makes a + // newly restarted full budget fall outside this window. + if !ok || deadline.Before(warm.returned.Add(budget)) || deadline.After(firstRead.Add(budget)) { + t.Fatalf("fallback deadline %v not in original budget window [%v,%v]", deadline, warm.returned.Add(budget), firstRead.Add(budget)) + } + if stage == "h2_preface" { + // The fresh pointer/worker are published by the dial callback + // before its TLS exchange; wait for its actual H2 write. + select { + case <-verified.Load().initializing: + case err := <-openDone: + t.Fatalf("fallback did not reach H2 preface: %v", err) + case <-time.After(2 * time.Second): + t.Fatal("fresh TLS connection did not reach H2 preface") + } + fresh = verified.Load() + select { + case err := <-serverDone: + if err != nil { + t.Fatalf("fresh fallback TLS handshake: %v", err) + } + case <-time.After(2 * time.Second): + t.Fatal("fresh fallback TLS worker did not report") + } + select { + case <-fresh.closed: + t.Fatal("old attempt's cleanup prematurely closed the fresh connection") + default: + } + } + want := error(context.DeadlineExceeded) + switch mode { + case "caller_cancel": + want = errors.New("caller cancelled fallback initialization") + cancel(want) + case "client_close": + want = context.Canceled + if err := fixture.client.Close(); err != nil { + t.Fatal(err) + } + } + select { + case err := <-openDone: + if !errors.Is(err, want) { + t.Fatalf("fallback initialization error = %v, want %v", err, want) + } + case <-time.After(2 * time.Second): + t.Fatal("fallback initialization ignored its cancellation/budget") + } + if stage == "h2_preface" { + select { + case <-fresh.closed: + default: + t.Fatal("failed fallback initializer retained its fresh raw connection") + } + } + webH2FallbackCheckWarmHRR(t, fixture) + fixture.assertEmpty(t, 3) + }) + } + }) + } +} + +func TestWebH2ResumptionFallbackSecondFailureDoesNotLoop(t *testing.T) { + for _, mode := range []string{"redial_error", "fresh_tls_certificate_error"} { + t.Run(mode, func(t *testing.T) { + fixture := webH2FallbackPrime(t) + warm := webH2FallbackInstallWarm(t, fixture) + fault := errors.New("fresh fallback dial failed") + untrustedTLS, _ := webH2ResumptionTLSConfigs(t) + warm.fresh = func(ctx context.Context) (net.Conn, error) { + if mode == "redial_error" { + return nil, fault + } + return webH2FallbackOtherTLSPeer(t, ctx, untrustedTLS, nil) + } + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + defer cancel() + openDone := webH2FallbackOpenWorker(t, fixture.client, ctx, warm) + webH2FallbackReleaseWarm(t, warm, 0) + select { + case err := <-openDone: + if mode == "redial_error" { + if !errors.Is(err, fault) { + t.Fatalf("redial failure = %v", err) + } + } else { + var untrusted x509.UnknownAuthorityError + if !errors.As(err, &untrusted) { + t.Fatalf("fresh certificate failure = %v", err) + } + } + case <-ctx.Done(): + t.Fatal("second initialization failure did not return") + } + webH2FallbackCheckWarmHRR(t, fixture) + fixture.assertEmpty(t, 3) + }) + } +} + +func TestWebH2ResumptionFallbackDoesNotRetryOtherTLSErrors(t *testing.T) { + for _, mode := range []string{"certificate", "alpn", "matching_text_without_hrr"} { + t.Run(mode, func(t *testing.T) { + fixture := webH2FallbackPrime(t) + serverTLS := fixture.serverTLS.Clone() + serverTLS.SessionTicketsDisabled = true + serverTLS.CurvePreferences = []tls.CurveID{tls.X25519} + if mode == "certificate" { + serverTLS, _ = webH2ResumptionTLSConfigs(t) + serverTLS.SessionTicketsDisabled = true + } + if mode == "alpn" { + serverTLS.NextProtos = []string{webHTTP11ALPN} + } + // Keep the pinned dependency error text independent of production's + // classifier, so an old-production overlay can exercise this test. + fault := errors.New("uTLS does not support reprocessing of PSK key triggered by HelloRetryRequest") + if mode == "matching_text_without_hrr" { + fixture.client.tlsConfig.VerifyPeerCertificate = func(_ [][]byte, _ [][]*x509.Certificate) error { return fault } + } + observed := make(chan *webH2HRRWireConn, 1) + var raw *webH2FallbackRaw + fixture.client.dialer = transport.DialFunc(func(ctx context.Context, _, _ string) (net.Conn, error) { + if fixture.calls.Add(1) != 2 { + return nil, errors.New("unexpected retry for ordinary TLS failure") + } + conn, err := webH2FallbackOtherTLSPeer(t, ctx, serverTLS, observed) + if err != nil { + return nil, err + } + raw = &webH2FallbackRaw{Conn: conn, closed: make(chan struct{})} + return raw, nil + }) + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + defer cancel() + _, err := fixture.client.openSession(ctx) + switch mode { + case "certificate": + var untrusted x509.UnknownAuthorityError + if !errors.As(err, &untrusted) { + t.Fatalf("ordinary certificate failure = %v", err) + } + case "alpn": + if err == nil || !strings.Contains(err.Error(), "did not negotiate h2") { + t.Fatalf("ordinary ALPN failure = %v", err) + } + case "matching_text_without_hrr": + if !errors.Is(err, fault) { + t.Fatalf("certificate callback failure = %v", err) + } + } + select { + case <-raw.closed: + default: + t.Fatal("ordinary TLS failure retained its raw connection") + } + var wire *webH2HRRWireConn + select { + case wire = <-observed: + case <-time.After(2 * time.Second): + t.Fatal("ordinary TLS failure peer was not observed") + } + hello, hrr, _, parseErr := wire.snapshot() + if parseErr != nil || hrr || !containsUint16(hello.extensions, 41) { + t.Fatalf("ordinary failed warm attempt: parse=%v HRR=%t PSK=%t", parseErr, hrr, containsUint16(hello.extensions, 41)) + } + fixture.assertEmpty(t, 2) + }) + } +} + +type webH2FallbackFixture struct { + client *WebH2Client + serverTLS *tls.Config + attempts <-chan webH2HRRAttempt + calls atomic.Int64 + destination *atomic.Int64 +} + +// GET primes only TLS tickets and H2 initialization, deliberately bypassing +// CONNECT/app authentication. Healthy authenticated HRR tests live separately. +func webH2FallbackPrime(t *testing.T) *webH2FallbackFixture { + t.Helper() + serverTLS, config := webH2ResumptionTLSConfigs(t) + address, attempts, requests, handlerDone, destination := webH2HRRServer(t, serverTLS, "unused.invalid:443", 2) + config.ClientSessionCache = tls.NewLRUClientSessionCache(2) + client, err := NewWebH2Client(WebH2ClientConfig{ServerAddress: address, Token: webTestToken, TLSConfig: config, HandshakeTimeout: 3 * time.Second}) + if err != nil { + t.Fatal(err) + } + fixture := &webH2FallbackFixture{client: client, serverTLS: serverTLS, attempts: attempts, destination: destination} + client.dialer = transport.DialFunc(func(ctx context.Context, network, address string) (net.Conn, error) { + fixture.calls.Add(1) + return (&net.Dialer{}).DialContext(ctx, network, address) + }) + t.Cleanup(func() { _ = client.Close() }) + cache := &webH2ResumptionUTLSCache{inner: client.utlsSessionCache} + client.utlsSessionCache = cache + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + defer cancel() + session, err := client.openSession(ctx) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = session.raw.Close(); _ = session.h2.Close() }) + request, err := http.NewRequestWithContext(ctx, http.MethodGet, "https://"+address+"/prime", nil) + if err != nil { + t.Fatal(err) + } + response, err := session.h2.RoundTrip(request) + if err != nil { + t.Fatal(err) + } + _, readErr := io.ReadAll(response.Body) + _ = response.Body.Close() + if readErr != nil || response.StatusCode != http.StatusNotFound || cache.puts.Load() == 0 { + t.Fatalf("ordinary GET prime: status=%d read=%v ticket puts=%d", response.StatusCode, readErr, cache.puts.Load()) + } + select { + case request := <-requests: + if request.physical != 0 || request.bearerBytes != 0 { + t.Fatal("ticket prime unexpectedly used application authentication") + } + case <-ctx.Done(): + t.Fatal("ordinary GET prime was not observed") + } + select { + case <-handlerDone: + case <-ctx.Done(): + t.Fatal("ordinary prime handler did not return") + } + select { + case attempt := <-attempts: + if attempt.sequence != 0 || !attempt.hrr || attempt.err != nil || containsUint16(attempt.hello.extensions, 41) { + t.Fatalf("prime was not real cold HRR: %+v", attempt) + } + case <-ctx.Done(): + t.Fatal("cold HRR prime was not observed") + } + _ = session.raw.Close() + _ = session.h2.Close() + return fixture +} + +func (f *webH2FallbackFixture) assertEmpty(t *testing.T, wantCalls int64) { + t.Helper() + f.client.mu.Lock() + registered, current, selected := len(f.client.sessions), f.client.current, f.client.selected + f.client.mu.Unlock() + if f.calls.Load() != wantCalls || registered != 0 || current != nil || selected || f.destination.Load() != 0 { + t.Fatalf("calls/sessions/current/selected/destination=%d/%d/%p/%t/%d", f.calls.Load(), registered, current, selected, f.destination.Load()) + } +} + +type webH2FallbackRaw struct { + net.Conn + closed chan struct{} + once sync.Once +} + +func (c *webH2FallbackRaw) Close() error { + err := c.Conn.Close() + c.once.Do(func() { close(c.closed) }) + return err +} + +type webH2FallbackWarm struct { + *webH2FallbackRaw + started chan time.Time + release chan struct{} + readOnce sync.Once + releaseOnce sync.Once + returned time.Time + fresh func(context.Context) (net.Conn, error) +} + +func (c *webH2FallbackWarm) Read(p []byte) (int, error) { + c.readOnce.Do(func() { + c.started <- time.Now() + select { + case <-c.release: + case <-c.closed: + } + }) + return c.Conn.Read(p) +} + +func webH2FallbackInstallWarm(t *testing.T, fixture *webH2FallbackFixture) *webH2FallbackWarm { + t.Helper() + warm := &webH2FallbackWarm{started: make(chan time.Time, 1), release: make(chan struct{})} + fixture.client.dialer = transport.DialFunc(func(ctx context.Context, network, address string) (net.Conn, error) { + switch fixture.calls.Add(1) { + case 2: + raw, err := (&net.Dialer{}).DialContext(ctx, network, address) + if err != nil { + return nil, err + } + warm.webH2FallbackRaw = &webH2FallbackRaw{Conn: raw, closed: make(chan struct{})} + warm.returned = time.Now() + return warm, nil + case 3: + select { + case <-warm.closed: + default: + return nil, errors.New("fallback redial began before old raw Close completed") + } + return warm.fresh(ctx) + default: + return nil, errors.New("initializer exceeded its single fresh retry") + } + }) + return warm +} + +func webH2FallbackReleaseWarm(t *testing.T, warm *webH2FallbackWarm, delay time.Duration) time.Time { + t.Helper() + var started time.Time + select { + case started = <-warm.started: + case <-time.After(2 * time.Second): + t.Fatal("warm HRR TLS read was not reached") + } + if delay > 0 { + timer := time.NewTimer(delay) + defer timer.Stop() + <-timer.C + } + warm.releaseOnce.Do(func() { close(warm.release) }) + return started +} + +func webH2FallbackCheckWarmHRR(t *testing.T, fixture *webH2FallbackFixture) { + t.Helper() + // Prime evidence was consumed before the warm initializer began, so only + // this attempt can be present; no completion-order assumption is made. + select { + case attempt := <-fixture.attempts: + peerCloseError := errors.Is(attempt.err, io.EOF) || errors.Is(attempt.err, syscall.ECONNRESET) + if attempt.sequence != 1 || !attempt.hrr || !containsUint16(attempt.hello.extensions, 41) || !attempt.peerClosed || !peerCloseError { + t.Fatalf("warm fallback was not a real client-closed HRR+PSK failure: sequence=%d HRR=%t PSK=%t peerClosed=%t err=%T/%v", attempt.sequence, attempt.hrr, containsUint16(attempt.hello.extensions, 41), attempt.peerClosed, attempt.err, attempt.err) + } + t.Log("real P256 HRR + populated PSK confirmed; old physical peer closure before test cleanup") + case <-time.After(2 * time.Second): + t.Fatal("warm HRR attempt was not observed") + } +} + +func webH2FallbackOpenWorker(t *testing.T, client *WebH2Client, ctx context.Context, warm *webH2FallbackWarm) <-chan error { + t.Helper() + return webH2FallbackWorker(t, func() error { + session, err := client.openSession(ctx) + if session != nil { + _ = session.raw.Close() + _ = session.h2.Close() + return errors.New("failed initialization unexpectedly returned a session") + } + return err + }, func() { _ = client.Close(); warm.releaseOnce.Do(func() { close(warm.release) }) }) +} + +func webH2FallbackWorker(t *testing.T, run func() error, stop func()) <-chan error { + t.Helper() + done := make(chan error, 1) + joined := make(chan struct{}) + t.Cleanup(func() { + stop() + select { + case <-joined: + case <-time.After(2 * time.Second): + t.Error("fallback test worker did not join") + } + }) + go func() { defer close(joined); defer close(done); done <- run() }() + return done +} + +func webH2FallbackOtherTLSPeer(t *testing.T, ctx context.Context, serverTLS *tls.Config, observed chan<- *webH2HRRWireConn) (net.Conn, error) { + t.Helper() + listener, err := net.ListenTCP("tcp", &net.TCPAddr{IP: net.IPv4(127, 0, 0, 1)}) + if err != nil { + return nil, err + } + defer listener.Close() + _ = listener.SetDeadline(time.Now().Add(time.Second)) + raw, err := (&net.Dialer{Timeout: time.Second}).DialContext(ctx, "tcp", listener.Addr().String()) + if err != nil { + return nil, err + } + peer, err := listener.Accept() + if err != nil { + _ = raw.Close() + return nil, err + } + config := serverTLS.Clone() + config.MinVersion, config.MaxVersion = tls.VersionTLS13, tls.VersionTLS13 + if len(config.NextProtos) == 0 { + config.NextProtos = []string{webH2ALPN} + } + config.CurvePreferences = []tls.CurveID{tls.X25519} + config.SessionTicketsDisabled = true + wire := &webH2HRRWireConn{Conn: peer} + if observed != nil { + observed <- wire + } + serverCtx, cancel := context.WithTimeout(context.Background(), 3*time.Second) + deadline, _ := serverCtx.Deadline() + _ = peer.SetDeadline(deadline) + webH2FallbackWorker(t, func() error { return tls.Server(wire, config).HandshakeContext(serverCtx) }, func() { cancel(); _ = raw.Close(); _ = peer.Close() }) + return raw, nil +} diff --git a/internal/tunnel/web_h2_resumption_hrr_test.go b/internal/tunnel/web_h2_resumption_hrr_test.go new file mode 100644 index 0000000..c70e6d0 --- /dev/null +++ b/internal/tunnel/web_h2_resumption_hrr_test.go @@ -0,0 +1,523 @@ +package tunnel + +import ( + "bytes" + "context" + "crypto/tls" + "encoding/binary" + "errors" + "fmt" + "io" + "net" + "net/http" + "sync" + "sync/atomic" + "syscall" + "testing" + "time" + + "github.com/cppla/autocar/internal/transport" + "golang.org/x/net/http2" +) + +type webH2HRRWireConn struct { + net.Conn + mu sync.Mutex + in, out []byte + clientClosed bool +} + +func (c *webH2HRRWireConn) Read(p []byte) (int, error) { + n, err := c.Conn.Read(p) + c.mu.Lock() + if n > 0 && len(c.in)+n <= 64<<10 { + c.in = append(c.in, p[:n]...) + } + if errors.Is(err, io.EOF) || errors.Is(err, syscall.ECONNRESET) { + c.clientClosed = true + } + c.mu.Unlock() + return n, err +} + +func (c *webH2HRRWireConn) Write(p []byte) (int, error) { + n, err := c.Conn.Write(p) + if n > 0 { + c.mu.Lock() + if len(c.out)+n <= 64<<10 { + c.out = append(c.out, p[:n]...) + } + c.mu.Unlock() + } + return n, err +} + +func (c *webH2HRRWireConn) snapshot() (parsedClientHello, bool, bool, error) { + c.mu.Lock() + defer c.mu.Unlock() + if len(c.in) < 5 { + return parsedClientHello{}, false, c.clientClosed, errors.New("missing real ClientHello") + } + length := 5 + int(binary.BigEndian.Uint16(c.in[3:5])) + if length > len(c.in) { + return parsedClientHello{}, false, c.clientClosed, errors.New("incomplete real ClientHello") + } + hello, err := parseTLSClientHello(c.in[:length]) + // TLS 1.3 HRR is a real ServerHello with this fixed RFC 8446 random. + hrrRandom := []byte{0xcf, 0x21, 0xad, 0x74, 0xe5, 0x9a, 0x61, 0x11, 0xbe, 0x1d, 0x8c, 0x02, 0x1e, 0x65, 0xb8, 0x91, 0xc2, 0xa2, 0x11, 0x16, 0x7a, 0xbb, 0x8c, 0x5e, 0x07, 0x9e, 0x09, 0xe2, 0xc8, 0xa8, 0x33, 0x9c} + hrr := len(c.out) >= 43 && c.out[0] == 22 && c.out[5] == 2 && bytes.Equal(c.out[11:43], hrrRandom) + return hello, hrr, c.clientClosed, err +} + +type webH2HRRAttempt struct { + sequence int + state tls.ConnectionState + hello parsedClientHello + hrr bool + peerClosed bool + err error +} + +// Failed HRR and replacement handshakes finish on independent workers. +// Sequence is accept identity, not completion order. This collector belongs +// only to one subtest; it never shares observations with another fixture. +type webH2HRRAttemptCollector struct { + events <-chan webH2HRRAttempt + pending map[int]webH2HRRAttempt + seen map[int]struct{} +} + +func newWebH2HRRAttemptCollector(events <-chan webH2HRRAttempt) *webH2HRRAttemptCollector { + return &webH2HRRAttemptCollector{events: events, pending: make(map[int]webH2HRRAttempt), seen: make(map[int]struct{})} +} + +func (c *webH2HRRAttemptCollector) take(t *testing.T, sequence int) webH2HRRAttempt { + t.Helper() + if attempt, ok := c.pending[sequence]; ok { + delete(c.pending, sequence) + return attempt + } + timer := time.NewTimer(2 * time.Second) + defer timer.Stop() + for { + select { + case attempt := <-c.events: + if _, duplicate := c.seen[attempt.sequence]; duplicate { + t.Fatalf("duplicate real HRR attempt sequence %d", attempt.sequence) + } + c.seen[attempt.sequence] = struct{}{} + if attempt.sequence == sequence { + return attempt + } + c.pending[attempt.sequence] = attempt + case <-timer.C: + t.Fatalf("real HRR handshake worker %d did not report", sequence) + } + } +} + +type webH2HRRClientRaw struct { + net.Conn + closes atomic.Int64 +} + +func (c *webH2HRRClientRaw) Close() error { + err := c.Conn.Close() + c.closes.Add(1) + return err +} + +func webH2HRREchoOrigin(t *testing.T) (string, <-chan error) { + t.Helper() + listener, err := net.ListenTCP("tcp", &net.TCPAddr{IP: net.ParseIP("127.0.0.1")}) + if err != nil { + t.Fatal(err) + } + results := make(chan error, 4) + done := make(chan struct{}) + var mu sync.Mutex + closed := false + owned := make(map[net.Conn]struct{}) + go func() { + defer close(done) + var workers sync.WaitGroup + defer workers.Wait() + for range 4 { + _ = listener.SetDeadline(time.Now().Add(4 * time.Second)) + conn, err := listener.AcceptTCP() + if err != nil { + results <- err + return + } + mu.Lock() + if closed { + mu.Unlock() + _ = conn.Close() + return + } + owned[conn] = struct{}{} + mu.Unlock() + workers.Add(1) + go func() { + defer workers.Done() + defer conn.Close() + defer func() { mu.Lock(); delete(owned, conn); mu.Unlock() }() + _ = conn.SetDeadline(time.Now().Add(4 * time.Second)) + payload, readErr := io.ReadAll(io.LimitReader(conn, 8<<10)) + if readErr != nil { + results <- readErr + return + } + reply := append([]byte("hrr-reply:"), payload...) + n, writeErr := conn.Write(reply) + if writeErr == nil && n != len(reply) { + writeErr = io.ErrShortWrite + } + results <- writeErr + }() + } + }() + t.Cleanup(func() { + _ = listener.Close() + mu.Lock() + closed = true + for conn := range owned { + _ = conn.Close() + } + mu.Unlock() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Error("real HRR echo Accept/socket workers did not join") + } + }) + return listener.Addr().String(), results +} + +func webH2HRRServer(t *testing.T, serverTLS *tls.Config, destination string, expectedAttempts int) (string, <-chan webH2HRRAttempt, <-chan webH2ResumptionRequest, <-chan struct{}, *atomic.Int64) { + t.Helper() + listener, err := net.ListenTCP("tcp", &net.TCPAddr{IP: net.ParseIP("127.0.0.1")}) + if err != nil { + t.Fatal(err) + } + var destinationCalls atomic.Int64 + core, err := newServerCoreWithAdmission(webTestToken, transport.DialFunc(func(ctx context.Context, network, address string) (net.Conn, error) { + destinationCalls.Add(1) + if network != "tcp" || address != destination { + return nil, errors.New("unexpected HRR test destination") + } + return (&net.Dialer{}).DialContext(ctx, network, address) + }), 0, 0, 0, 0, nil) + if err != nil { + _ = listener.Close() + t.Fatal(err) + } + verifier, err := newWebAuthVerifier(mustWebAuthKey(t, webTestToken), nil, 0) + if err != nil { + _ = listener.Close() + t.Fatal(err) + } + handler := &webTunnelHandler{core: core, auth: verifier, cover: http.NotFoundHandler()} + config := serverTLS.Clone() + config.MinVersion, config.MaxVersion = tls.VersionTLS13, tls.VersionTLS13 + config.NextProtos = []string{webH2ALPN} + config.CurvePreferences = []tls.CurveID{tls.CurveP256} + attempts := make(chan webH2HRRAttempt, expectedAttempts) + requests := make(chan webH2ResumptionRequest, 4) + handlerDone := make(chan struct{}, 4) + done := make(chan struct{}) + var mu sync.Mutex + closed := false + owned := make(map[net.Conn]struct{}) + go func() { + defer close(done) + var workers sync.WaitGroup + defer workers.Wait() + for sequence := range expectedAttempts { + _ = listener.SetDeadline(time.Now().Add(4 * time.Second)) + raw, err := listener.AcceptTCP() + if err != nil { + attempts <- webH2HRRAttempt{sequence: sequence, err: err} + return + } + mu.Lock() + if closed { + mu.Unlock() + _ = raw.Close() + return + } + owned[raw] = struct{}{} + mu.Unlock() + workers.Add(1) + go func() { + defer workers.Done() + defer raw.Close() + defer func() { mu.Lock(); delete(owned, raw); mu.Unlock() }() + _ = raw.SetDeadline(time.Now().Add(4 * time.Second)) + wire := &webH2HRRWireConn{Conn: raw} + tlsConn := tls.Server(wire, config) + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + handshakeErr := tlsConn.HandshakeContext(ctx) + cancel() + hello, hrr, peerClosed, parseErr := wire.snapshot() + if handshakeErr == nil { + handshakeErr = parseErr + } + attempts <- webH2HRRAttempt{sequence: sequence, state: tlsConn.ConnectionState(), hello: hello, hrr: hrr, peerClosed: peerClosed, err: handshakeErr} + if handshakeErr != nil { + return + } + ctx = context.WithValue(context.Background(), webTLSConnectionContextKey{}, tlsConn) + ctx = context.WithValue(ctx, webServerConnectionAuthContextKey{}, newWebServerConnectionAuth(raw.Close)) + observed := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + defer func() { + select { + case handlerDone <- struct{}{}: + default: + } + }() + select { + case requests <- webH2ResumptionRequest{physical: sequence, bearerBytes: len(r.Header.Get("Proxy-Authorization"))}: + case <-r.Context().Done(): + return + } + handler.ServeHTTP(w, r) + }) + (&http2.Server{}).ServeConn(tlsConn, &http2.ServeConnOpts{Context: ctx, Handler: observed}) + }() + } + }() + t.Cleanup(func() { + _ = listener.Close() + mu.Lock() + closed = true + for conn := range owned { + _ = conn.Close() + } + mu.Unlock() + select { + case <-done: + case <-time.After(2 * time.Second): + t.Error("real HRR Accept/TLS/ServeConn workers did not join") + } + if got := len(core.sem); got != 0 { + t.Errorf("HRR destination admission retained %d slots", got) + } + }) + return listener.Addr().String(), attempts, requests, handlerDone, &destinationCalls +} + +func webH2HRRReadAttempt(t *testing.T, collector *webH2HRRAttemptCollector, sequence int, wantFailure, wantPSK, wantResume bool) { + t.Helper() + attempt := collector.take(t, sequence) + if attempt.sequence != sequence || !attempt.hrr { + t.Errorf("attempt sequence/real HRR=%d/%t, want %d/true", attempt.sequence, attempt.hrr, sequence) + } + if got := containsUint16(attempt.hello.extensions, 41); got != wantPSK { + t.Errorf("physical %d actual offered PSK=%t, want %t", sequence, got, wantPSK) + } + if containsUint16(attempt.hello.extensions, 42) { + t.Error("HRR attempt offered TLS early_data") + } + if wantFailure { + if !attempt.peerClosed || (!errors.Is(attempt.err, io.EOF) && !errors.Is(attempt.err, syscall.ECONNRESET)) { + t.Errorf("failed warm physical %d was not client-closed before cleanup: peerClosed=%t err=%T: %v", sequence, attempt.peerClosed, attempt.err, attempt.err) + } + } else { + if attempt.err != nil { + t.Errorf("physical %d handshake error: %v", sequence, attempt.err) + } + if attempt.state.Version != tls.VersionTLS13 || attempt.state.NegotiatedProtocol != webH2ALPN || attempt.state.DidResume != wantResume { + t.Errorf("physical %d TLS/h2/DidResume=%#x/%q/%t, want TLS13/h2/%t", sequence, attempt.state.Version, attempt.state.NegotiatedProtocol, attempt.state.DidResume, wantResume) + } + } + t.Logf("physical=%d realHRR=%t offeredPSK=%t serverDidResume=%t peerClosed=%t peerEOF=%t peerReset=%t failed=%t", sequence, attempt.hrr, wantPSK, attempt.state.DidResume, attempt.peerClosed, errors.Is(attempt.err, io.EOF), errors.Is(attempt.err, syscall.ECONNRESET), attempt.err != nil) +} + +func TestWebH2ResumptionHRRFallbackKeepsHealthyTunnel(t *testing.T) { + for _, test := range []struct { + name string + profile FingerprintProfile + cacheEnabled bool + ticketsDisabled bool + fallback bool + resumed bool + }{ + {name: "chrome_enabled", profile: FingerprintChrome133, cacheEnabled: true, fallback: true}, + {name: "chrome_nil_cache", profile: FingerprintChrome133}, + {name: "chrome_tickets_disabled", profile: FingerprintChrome133, cacheEnabled: true, ticketsDisabled: true}, + {name: "native_positive", profile: FingerprintNative, cacheEnabled: true, resumed: true}, + } { + t.Run(test.name, func(t *testing.T) { + target, targetResults := webH2HRREchoOrigin(t) + serverTLS, config := webH2ResumptionTLSConfigs(t) + config.SessionTicketsDisabled = test.ticketsDisabled + standardCache := &webH2ResumptionTLSCache{inner: tls.NewLRUClientSessionCache(4)} + if test.cacheEnabled { + config.ClientSessionCache = standardCache + } + physicalCount := 2 + if test.fallback { + physicalCount = 3 + } + address, attempts, requests, handlerDone, destinationCalls := webH2HRRServer(t, serverTLS, target, physicalCount) + collector := newWebH2HRRAttemptCollector(attempts) + client, err := NewWebH2Client(WebH2ClientConfig{ServerAddress: address, Token: webTestToken, TLSConfig: config, FingerprintProfile: test.profile}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = client.Close() }) + var chromeCache *webH2ResumptionUTLSCache + if test.fallback { + chromeCache = &webH2ResumptionUTLSCache{inner: client.utlsSessionCache} + client.utlsSessionCache = chromeCache + } + var rawMu sync.Mutex + var rawConnections []*webH2HRRClientRaw + client.dialer = transport.DialFunc(func(ctx context.Context, network, address string) (net.Conn, error) { + rawMu.Lock() + if test.fallback && len(rawConnections) == 2 && rawConnections[1].closes.Load() == 0 { + rawMu.Unlock() + return nil, errors.New("fresh HRR dial began before failed raw TCP Close completed") + } + rawMu.Unlock() + raw, err := (&net.Dialer{}).DialContext(ctx, network, address) + if err != nil { + return nil, err + } + observed := &webH2HRRClientRaw{Conn: raw} + rawMu.Lock() + rawConnections = append(rawConnections, observed) + rawMu.Unlock() + return observed, nil + }) + var successful [2]*webH2ClientSession + var authKeys [2]webSessionKey + for request := range 4 { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + conn, dialErr := client.DialContext(ctx, "tcp", target) + cancel() + if dialErr != nil { + if conn != nil { + _ = conn.Close() + } + // Preserve direct warm-PSK/HRR/peer-close evidence even on the old + // candidate, before fixture/client cleanup can close sockets. + if request == 2 && test.fallback { + webH2HRRReadAttempt(t, collector, 1, true, true, false) + } + t.Fatalf("healthy authenticated request %d after real HRR: %v", request, dialErr) + } + if conn == nil { + t.Fatal("healthy CONNECT returned nil connection") + } + if request%2 == 0 { + if request == 0 { + webH2HRRReadAttempt(t, collector, 0, false, false, false) + } else if test.fallback { + webH2HRRReadAttempt(t, collector, 1, true, true, false) + webH2HRRReadAttempt(t, collector, 2, false, false, false) + rawMu.Lock() + failedWasClosed := len(rawConnections) == 3 && rawConnections[1].closes.Load() > 0 + rawMu.Unlock() + if !failedWasClosed { + t.Error("warm failed raw TCP was not closed by client before fixture cleanup") + } + } else { + webH2HRRReadAttempt(t, collector, 1, false, test.resumed, test.resumed) + } + } + _ = conn.SetDeadline(time.Now().Add(2 * time.Second)) + payload := fmt.Sprintf("hrr-authenticated-payload-%d", request) + n, writeErr := io.WriteString(conn, payload) + if writeErr != nil || n != len(payload) { + _ = conn.Close() + t.Fatalf("healthy payload write: n=%d err=%v", n, writeErr) + } + writer, ok := conn.(interface{ CloseWrite() error }) + if !ok { + _ = conn.Close() + t.Fatal("H2 tunnel has no upload FIN") + } + if err := writer.CloseWrite(); err != nil { + _ = conn.Close() + t.Fatal(err) + } + reply, readErr := io.ReadAll(conn) + _ = conn.Close() + if readErr != nil || string(reply) != "hrr-reply:"+payload { + t.Fatalf("healthy full TCP response: %q err=%v", reply, readErr) + } + select { + case err := <-targetResults: + if err != nil { + t.Fatal(err) + } + case <-time.After(2 * time.Second): + t.Fatal("real TCP echo worker did not report") + } + select { + case <-handlerDone: + case <-time.After(2 * time.Second): + t.Fatal("healthy CONNECT handler did not return") + } + index := request / 2 + client.mu.Lock() + current := client.current + ready := current != nil && current.authState == webH2ClientAuthReady && current.auth != nil + if ready { + authKeys[index] = current.auth.key + } + client.mu.Unlock() + if !ready { + t.Fatal("healthy HRR tunnel did not retain authenticated session") + } + if request%2 == 0 { + successful[index] = current + state := current.conn.ConnectionState() + if state.DidResume != (index == 1 && test.resumed) || state.Version != tls.VersionTLS13 || state.NegotiatedProtocol != webH2ALPN || len(state.VerifiedChains) == 0 { + t.Error("public client lost verified TLS13/h2 or expected resumption policy") + } + } else if current != successful[index] { + t.Error("continuation changed its healthy physical connection") + } + select { + case observed := <-requests: + physical := index + if index == 1 && test.fallback { + physical = 2 + } + if observed.physical != physical || (request%2 == 0 && observed.bearerBytes <= 41) || (request%2 == 1 && observed.bearerBytes != 41) { + t.Errorf("request %d physical/bearer=%d/%d violated fresh-bootstrap/continuation policy", request, observed.physical, observed.bearerBytes) + } + case <-time.After(2 * time.Second): + t.Fatal("healthy request auth observation did not arrive") + } + if request == 1 { + current.h2.SetDoNotReuse() + } + } + if successful[0] == successful[1] || authKeys[0] == authKeys[1] { + t.Error("HRR replacement inherited a physical connection or proxy key") + } + if destinationCalls.Load() != 4 { + t.Errorf("healthy destination calls=%d, want 4", destinationCalls.Load()) + } + rawMu.Lock() + actualAttempts := len(rawConnections) + rawMu.Unlock() + if actualAttempts != physicalCount { + t.Errorf("actual raw TCP attempts=%d, want exactly %d", actualAttempts, physicalCount) + } + if chromeCache != nil { + if chromeCache.puts.Load() < 1 || chromeCache.hits.Load() != 1 { + t.Errorf("real HRR Chrome cache puts/hits=%d/%d, want stored ticket and exactly one loaded attempt", chromeCache.puts.Load(), chromeCache.hits.Load()) + } + t.Logf("Chrome HRR fallback actual TCP attempts=%d ticket puts=%d hits=%d", actualAttempts, chromeCache.puts.Load(), chromeCache.hits.Load()) + } else if test.resumed && (standardCache.puts.Load() < 1 || standardCache.hits.Load() < 1) { + t.Error("native resumed HRR cache positive control failed") + } + _ = client.Close() + }) + } +}