From eb82a4dbc9382e5bba86a22f18d306d32f1bf866 Mon Sep 17 00:00:00 2001 From: Balaji Ganesan Date: Tue, 25 Aug 2026 20:16:48 -0700 Subject: [PATCH] feat(worker): allow the stateful tunnel to run over TCP instead of QUIC Edge proxy CPU is proportional to bytes moved, and today every byte is encrypted and decrypted twice: once on the worker to edge proxy leg, and again on the edge proxy to grpc-proxy leg. On a measured high-volume workload the edge proxy achieved roughly 33 MB/s per core, reproduced across three runs at two different scales. Horizontal scaling does not relieve it, because QUIC tunnels do not migrate between proxy pods once established. This adds an opt-in TCP path, selected per pod by NVCF_WORKER_TCP_TUNNEL so it can be enabled for a single function and reverted by removing the variable. Everything else is derived from the HTTP/3 connection config the proxy already sends, so no control plane change is required. The connection is made in two nested steps, and the nesting is the point: 1. TLS to the regional TCP entry point, then an authority-form CONNECT naming the target pod. The edge proxy matches this and turns the hop into an opaque byte pipe. Because it does not parse what flows inside, HTTP/1.1's rule that a response may not begin before the request body completes does not apply, which is what makes a long-lived bidirectional tunnel workable over TCP at all. Sending the tunnel as a plain HTTP/1.1 request instead was tried and does not work: the connection establishes and then no bytes move in either direction. 2. The ordinary POST /v1/proxy inside the pipe. grpc-proxy sees a plain HTTP/1.1 request on a real TCP connection, so http.Hijacker works and the server side needs no change. Endpoint derivation drops the pod label and selects the TCP service name: ..proxy. -> .tcp-proxy. Pod targeting travels in the CONNECT authority rather than the DNS name, so this needs neither a wildcard DNS record nor a wildcard certificate: the regional name already resolves and is already covered by the TCP load balancer's certificate. Failure falls through to the existing quicConnect path, so enabling the variable cannot take a worker offline. Experimental and off by default. Unit tests cover the endpoint derivation and the test override. The network path is not yet exercised end to end; that needs a matching edge proxy route and will be measured against the existing QUIC baseline before any wider use. Co-Authored-By: Balaji Ganesan --- src/libraries/go/worker/proxy/proxy.go | 158 ++++++++++++++++++ .../go/worker/proxy/tcptunnel_test.go | 87 ++++++++++ 2 files changed, 245 insertions(+) create mode 100644 src/libraries/go/worker/proxy/tcptunnel_test.go diff --git a/src/libraries/go/worker/proxy/proxy.go b/src/libraries/go/worker/proxy/proxy.go index 699fe7558..273cf5129 100644 --- a/src/libraries/go/worker/proxy/proxy.go +++ b/src/libraries/go/worker/proxy/proxy.go @@ -29,6 +29,8 @@ import ( "net/http" "net/http/httputil" "net/url" + "os" + "strings" "sync" "sync/atomic" "time" @@ -360,6 +362,17 @@ func getClientConnFromProxy(ctx context.Context, work *pb.WorkerInvokeFunctionRe zap.L().Warn("tcp connection attempt failed", zap.Error(err)) } case *pb.WorkerInvokeFunctionRequest_StatefulConfig_ConnectionConfig_Http3Config: + // Opt-in experiment: carry the tunnel over TCP instead of QUIC, + // deriving the endpoint from the HTTP/3 config the proxy already + // sent. Failure falls through to QUIC below, so enabling this + // cannot take a worker offline. + if tcpTunnelEnabled { + clientConn, err = tcpTunnelConnect(ctx, work.RequestId, config.Http3Config) + if err == nil { + break + } + zap.L().Warn("tcp tunnel attempt failed, falling back to quic", zap.Error(err)) + } clientConn, err = quicConnect(ctx, work.RequestId, config.Http3Config, h3) if err != nil { zap.L().Warn("quic connection attempt failed", zap.Error(err)) @@ -390,6 +403,151 @@ func traceError(span trace.Span, err error) error { return err } +// Opt-in experiment: carry the worker tunnel over TCP rather than QUIC. +// +// Proxy CPU is proportional to bytes moved, and today every byte is encrypted +// and decrypted twice: once on the worker to edge proxy leg, and again on the +// edge proxy to grpc-proxy leg. Moving the worker leg to TCP removes QUIC from +// both. +// +// Enabled per pod so it can be applied to a single function and reverted by +// removing the variable. Everything else is derived from the HTTP/3 connection +// config the proxy already sends, so no control plane change is required. +var ( + tcpTunnelEnabled = os.Getenv("NVCF_WORKER_TCP_TUNNEL") == "1" || + os.Getenv("NVCF_WORKER_TCP_TUNNEL") == "true" + // Overrides for testing. Empty means derive from the HTTP/3 proxy URI. + tcpTunnelHostOverride = os.Getenv("NVCF_WORKER_TCP_TUNNEL_HOST") + tcpTunnelPort = envOrDefault("NVCF_WORKER_TCP_TUNNEL_PORT", "10086") +) + +func envOrDefault(key, def string) string { + if v := os.Getenv(key); v != "" { + return v + } + return def +} + +// tcpTunnelDialHost turns the pod-targeted HTTP/3 host into the regional TCP +// entry point, by dropping the pod label and selecting the TCP service name: +// +// ..proxy. -> .tcp-proxy. +// +// The pod label is dropped because it does not resolve on the TCP side. Pod +// targeting travels in the CONNECT authority instead, which the edge proxy +// rewrites. That is deliberate: it avoids needing a wildcard DNS record and a +// wildcard certificate, since the regional name already resolves and is already +// covered by the TCP load balancer's certificate. +func tcpTunnelDialHost(h3Host string) (string, error) { + if tcpTunnelHostOverride != "" { + return tcpTunnelHostOverride, nil + } + labels := strings.Split(h3Host, ".") + if len(labels) < 4 || labels[2] != "proxy" { + return "", fmt.Errorf("cannot derive tcp tunnel host from %q", h3Host) + } + rest := append([]string{labels[1], "tcp-proxy"}, labels[3:]...) + return strings.Join(rest, "."), nil +} + +// tcpTunnelConnect establishes the stateful tunnel over TCP. +// +// Two nested steps, and the nesting is the point: +// +// 1. TLS to the regional TCP entry point, then an authority-form CONNECT +// naming the target pod. The edge proxy matches this and turns the hop into +// an opaque byte pipe to that pod. Because it does not parse what flows +// inside, HTTP/1.1's rule that a response may not begin before the request +// body completes does not apply, which is what makes a long-lived +// bidirectional tunnel possible over TCP at all. +// 2. The ordinary POST /v1/proxy inside the pipe. grpc-proxy sees a plain +// HTTP/1.1 request on a real TCP connection, so http.Hijacker works and the +// server side needs no change. +func tcpTunnelConnect(ctx context.Context, requestId string, connectionConfig *pb.WorkerInvokeFunctionRequest_StatefulConfig_ConnectionConfig_HTTP3ConnectionConfig) (net.Conn, error) { + tracer := otel.GetTracerProvider().Tracer("nvcf-worker-lib") + ctx, span := tracer.Start(ctx, "CONNECT /v1/proxy", + trace.WithSpanKind(trace.SpanKindClient), + trace.WithAttributes( + attribute.String("proxy-type", "tcp-tunnel"), + semconv.HTTPURL(connectionConfig.ProxyURI), + )) + defer span.End() + + proxyURL, err := url.Parse(connectionConfig.ProxyURI) + if err != nil { + return nil, traceError(span, err) + } + podHost := proxyURL.Hostname() + dialHost, err := tcpTunnelDialHost(podHost) + if err != nil { + return nil, traceError(span, err) + } + connectAuthority := net.JoinHostPort(podHost, tcpTunnelPort) + + dialer := &tls.Dialer{ + NetDialer: &net.Dialer{Timeout: 3 * time.Second}, + Config: &tls.Config{ServerName: dialHost, InsecureSkipVerify: quicInsecure}, + } + c, err := dialer.DialContext(ctx, "tcp", net.JoinHostPort(dialHost, "443")) + if err != nil { + return nil, traceError(span, fmt.Errorf("dialing tcp tunnel %q failed: %w", dialHost, err)) + } + + // Authority-form CONNECT. A request-target carrying a path is not matched + // as a CONNECT by the edge proxy, which is why the HTTP/3 path uses POST. + if _, err = fmt.Fprintf(c, "CONNECT %s HTTP/1.1\r\nHost: %s\r\n\r\n", connectAuthority, connectAuthority); err != nil { + _ = c.Close() + return nil, traceError(span, err) + } + br := bufio.NewReader(c) + tunnelResp, err := http.ReadResponse(br, &http.Request{Method: http.MethodConnect}) + if err != nil { + _ = c.Close() + return nil, traceError(span, err) + } + if tunnelResp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(io.LimitReader(tunnelResp.Body, 512)) + _ = c.Close() + return nil, traceError(span, fmt.Errorf("tcp tunnel CONNECT to %s returned %d: %s", + connectAuthority, tunnelResp.StatusCode, string(body))) + } + + // Inside the pipe, speak exactly what the HTTP/1 server on grpc-proxy + // expects. Path form here, because its mux routes on path. + inner, err := http.NewRequestWithContext(ctx, http.MethodPost, "http://"+connectAuthority+"/v1/proxy", http.NoBody) + if err != nil { + _ = c.Close() + return nil, traceError(span, err) + } + inner.ContentLength = -1 + inner.Header.Set("Authorization", "Bearer "+connectionConfig.ProxyAuthorizationToken) + inner.Header.Set("X-Request-ID", requestId) + otelhttptrace.Inject(ctx, inner) + if err = inner.Write(c); err != nil { + _ = c.Close() + return nil, traceError(span, err) + } + + resp, err := http.ReadResponse(br, nil) + if err != nil { + _ = c.Close() + return nil, traceError(span, err) + } + if resp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) + _ = resp.Body.Close() + _ = c.Close() + err := fmt.Errorf("unexpected status %d from tunnelled proxy request: %s", resp.StatusCode, string(body)) + if resp.StatusCode == http.StatusUnauthorized || resp.StatusCode == http.StatusForbidden { + err = errors.Join(err, ErrAuth) + } + return nil, traceError(span, err) + } + + // Deliberately not closing the body: the stream is used directly from here. + return buffconn.NewBufConn(c, br), nil +} + func quicConnect(ctx context.Context, requestId string, connectionConfig *pb.WorkerInvokeFunctionRequest_StatefulConfig_ConnectionConfig_HTTP3ConnectionConfig, h3 *h3ConnectionCache) (net.Conn, error) { tracer := otel.GetTracerProvider().Tracer("nvcf-worker-lib") ctx, span := tracer.Start(ctx, "CONNECT /v1/proxy", diff --git a/src/libraries/go/worker/proxy/tcptunnel_test.go b/src/libraries/go/worker/proxy/tcptunnel_test.go new file mode 100644 index 000000000..f3d9f6c00 --- /dev/null +++ b/src/libraries/go/worker/proxy/tcptunnel_test.go @@ -0,0 +1,87 @@ +/* +SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +SPDX-License-Identifier: Apache-2.0 + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package proxy + +import "testing" + +func TestTCPTunnelDialHost(t *testing.T) { + // The pod label is dropped on purpose: it does not resolve on the TCP side. + // Pod targeting travels in the CONNECT authority, which the edge proxy + // rewrites. + cases := []struct { + name string + in string + want string + wantErr bool + }{ + { + name: "regional host with pod label", + in: "10-0-0-1.region-1.proxy.example.com", + want: "region-1.tcp-proxy.example.com", + }, + { + name: "deeper domain", + in: "10-0-0-1.region-1.proxy.sub.example.com", + want: "region-1.tcp-proxy.sub.example.com", + }, + { + name: "service label is not proxy", + in: "10-0-0-1.region-1.something.example.com", + wantErr: true, + }, + { + name: "too few labels", + in: "proxy.local", + wantErr: true, + }, + } + + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + got, err := tcpTunnelDialHost(tc.in) + if tc.wantErr { + if err == nil { + t.Fatalf("expected an error for %q, got %q", tc.in, got) + } + return + } + if err != nil { + t.Fatalf("unexpected error for %q: %v", tc.in, err) + } + if got != tc.want { + t.Errorf("tcpTunnelDialHost(%q) = %q, want %q", tc.in, got, tc.want) + } + }) + } +} + +func TestTCPTunnelDialHostOverride(t *testing.T) { + // The override exists so a test can point at an arbitrary endpoint without + // depending on the naming convention holding. + old := tcpTunnelHostOverride + t.Cleanup(func() { tcpTunnelHostOverride = old }) + + tcpTunnelHostOverride = "explicit.example.com" + got, err := tcpTunnelDialHost("10-0-0-1.region-1.proxy.example.com") + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if got != "explicit.example.com" { + t.Errorf("override ignored: got %q", got) + } +}