From 9db98021278ce511e75435bf23610f52880605ae Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Pawe=C5=82=20Gronowski?= Date: Thu, 8 Oct 2026 20:19:46 +0200 Subject: [PATCH 1/2] Update to go1.26.9 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit This release includes 15 security fixes following the security policy: - net/http: HTTP/2 server crash due to HPACK encoder race HTTP/2 servers could end up crashing due to inadvertently modifying its HPACK encoder concurrently. This happens because the server modifies the HPACK encoder from two goroutines without synchronization: one uses the encoder to encode a HEADERS frame as part of a response sent to a client and the other modifies the encoder's table size when handling a SETTINGS frame containing SETTINGS_HEADER_TABLE_SIZE that a client sends. A malicious client can repeatedly send a request while changing the header table size to crash the server. Fix this issue by not applying SETTINGS_HEADER_TABLE_SIZE immediately. Instead, buffer any SETTINGS_HEADER_TABLE_SIZE received, and only apply the new value prior to the next time the server writes a frame. Thanks to RyotaK (https://ryotak.net) of GMO Flatt Security Inc. for reporting this issue. This is CVE-2026-97032 and Go issue https://go.dev/issue/81867. - net/http: HTTP/2 server memory exhaustion due to Trailer headers When "Trailer" headers are sent by a client, the HTTP server internally uses the header values to populate the Request.Trailer map passed to the server handler. Because Request.Trailer is a map, each entry incurs memory overhead. For HTTP/2 servers, a malicious client can exploit this by sending a "Trailer" header that declares a large number of fields, causing the server to allocate a disproportionate amount of memory while bypassing Server.MaxHeaderValueCount and Server.MaxHeaderBytes limits. This exploit is not applicable for HTTP/1 servers, which do not support multiplexing a large number of requests over one TCP connection, and whose Server.MaxHeaderBytes are calculated differently. Server.MaxHeaderValueCount and Server.MaxHeaderBytes limits are now applied towards the trailer fields declared in "Trailer" headers. Thanks to RyotaK (https://ryotak.net) of GMO Flatt Security Inc. for reporting this issue. This is CVE-2026-78659 and Go issue https://go.dev/issue/81857. - crypto/tls: reject malformed ECH outer extension references Multiple ECH outer extension references are not permitted under RFC 9849; previously, a client could send a well-crafted packet that could trigger memory exhaustion in the server process by specifying multiple references. We now reject these as malformed and curb the memory amplification vector as a result. This is CVE-2026-97031 and Go issue https://go.dev/issue/81855. - cmd/go: checksum bypass for golang.org/fips140 Previously, a user operating inside of a malicious Go project that defines a bogus golang.org/fips140 and operates a malicious GOMODPROXY the user chooses to connect to can serve an arbitrary module in its place. We now unpack the trusted ziphash for the bundled golang.org/fips140 module and construct its entry in the GOMODCACHE such that it can be verified by the toolchain. This is CVE-2026-94444 and Go issue https://go.dev/issue/81833. - cmd/go: checksum database bypass for golang.org/toolchain Previously, a user operating inside of a malicious Go project that defines a bogus golang.org/toolchain go.sum entry and operates a malicious GOMODPROXY the user chooses to use can bypass the intended checksum. We now ensure that golang.org/toolchain always goes to the network for the canonical checksum. This is CVE-2026-94447 and Go issue https://go.dev/issue/81834. - html/template: reset context tracking on consecutive template expressions When a JavaScript template literal contains consecutive expressions, the context tracking state was not properly reset upon entering a new expression. We now ensure that template-literal expression entries correctly reset context variables so all subsequent regular expression literals are accurately recognized and escaped. This is CVE-2026-94448 and Go issue https://go.dev/issue/81821. - html/template: recognize yield as regexp preceder keyword A trusted template author may have previously written a valid template wherein the use of the yield keyword would not be correctly escaped. We now ensure that valid keyword uses are escaped and non-keyword uses are not escaped. This is CVE-2026-97030 and Go issue https://go.dev/issue/81823. - net/textproto, mime/multipart: memory limit bypass when parsing MIME headers Parsing a multipart form could bypass memory limits and read an arbitrarily long line into memory when the remaining limit at the start of a part was less than 400 bytes. Multipart form memory limits are now properly enforced in this situation. Thanks to Jakub Ciolek (https://ciolek.dev) for reporting this issue. This is CVE-2026-94440 and Go issue https://go.dev/issue/81741. - net/http: HTTP/1 client connection desynchronization after CONNECT rejection When http.Transport sends an HTTP/1 CONNECT request with a non-empty Request.Body, it writes the body directly to the connection without framing after the request headers. If the server rejects the CONNECT request with a non-2xx keep-alive response, Transport returns the connection to the idle pool. Because CONNECT requests do not have a request body, the server may interpret the trailing body bytes as a subsequent pipelined HTTP/1.1 request on the connection, leaving the pooled connection desynchronized and causing the next caller that reuses it to read the response to the injected request. In reverse proxies (including httputil.ReverseProxy) that forward CONNECT requests through a shared Transport, this can lead to cross-user response poisoning. The HTTP/1 transport now closes a connection after sending a CONNECT request, regardless of the response status. In addition, ReverseProxy now rejects incoming CONNECT requests with a 405 Method Not Allowed response. ReverseProxy has never handled CONNECT requests in a useful fashion (it does not convert the connection into a bidirectional tunnel), so we do not expect this change to negatively affect any current users. Thanks to Xclow3n (Rajat Raghav) for reporting this issue. This is CVE-2026-56866 and Go issue https://go.dev/issue/81740. - net/http: HTTP/1 server connection desynchronization after 2xx CONNECT response When an HTTP server handler sent a 2xx response to an HTTP/1 CONNECT request and returned without hijacking the connection, the server improperly continued to read and serve requests from the connection. Since a 2xx response to an HTTP/1 CONNECT converts the connection into a tunnel, the server should not treat the connection as continuing to contain HTTP. The impact of this misbehavior is mostly limited to potential request smuggling, where an intermediate proxy considers the data on the connection to be tunneled and the server considers it to be HTTP. The HTTP/1 server now always closes a connection after responding to a CONNECT request, regardless of the response status. Thanks to Jakub Ciolek (https://ciolek.dev) for reporting this issue. This is CVE-2026-94439 and Go issue https://go.dev/issue/81744. - net/http: excessive CPU consumption from repeated initial window changes A malicious HTTP/2 peer could cause excessive CPU consumption in the client or server by opening a large number of streams and then sending many small SETTINGS frames containing SETTINGS_INITIAL_WINDOW_SIZE values. The HTTP/2 client and server now efficiently handle changes to the initial window size (O(1) rather than O(number of streams)). Thanks to Jakub Ciolek (https://ciolek.dev) for reporting this issue. This is CVE-2026-78669 and Go issue https://go.dev/issue/81742. - net/http: HTTP/2 transport accepts malformed framing-related headers Historically, we have been rather lax about malformed framing-related headers in our HTTP/2 implementation, as they cannot interfere with HTTP/2 framing. However, this makes it possible for our HTTP/2 implementation to forward responses containing such headers to an HTTP/1 client when acting as a reverse proxy. If the HTTP/1 client also does not behave strictly enough, this can result in response smuggling. We now delete malformed framing-related headers when received by our HTTP/2 transport, so they will not be forwarded to a potentially vulnerable HTTP/1 client. Thanks to TJ Barton for reporting this issue. This is CVE-2026-78660 and Go issue https://go.dev/issue/81115. - os: Root.Mkdir(All) can follow junctions out of the root on Windows On Windows, when the target of Root.Mkdir or Root.MkdirAll was a junction pointing to an empty location, the operation would create a directory at the junction target even when that target was located outside the root. This only applies to operations where the last path component is a junction (path/to/junction, but not path/junction/target). Root.Mkdir and Root.MkdirAll now correctly avoid resolving junctions. Thanks to Daniele Ballarini for reporting this issue. This is CVE-2026-56857 and Go issue https://go.dev/issue/81739. - net/http: lack of limit on size of parsed Range headers When parsing a Range header containing a large number of small ranges, FileServer(FS), ServeContent, and ServeFile(FS) could consume an excessive amount of CPU. These functions now ignore Range headers containing more than 200 ranges. The limit is controlled by the new httpservecontentmaxranges= GODEBUG setting. Setting GODEBUG=httpservecontentmaxranges=0 disables the limit. Thanks to Jakub Ciolek (https://ciolek.dev) for reporting this issue. This is CVE-2026-78667 and Go issue https://go.dev/issue/81858. - net/http: double flow control refund on HTTP/2 server streams The HTTP/2 server could refund connection-level flow control twice for the same data: Once when a client resets a stream (refunding data for any sent-but-unread portion of the stream), and again when a request handler reads the buffered data. A malicious client could exploit this to bypass the configured connection-level flow control limit (MaxReceiveBufferPerConnection). Total buffered data is still limited by the concurrent stream limit and stream-level flow control. The HTTP/2 server now waits to refund connection-level flow control for reset streams until after the request handler is complete. Thanks to Ali Sherif (https://www.linkedin.com/in/ali-sherif-13812b276/) for reporting this issue. This is CVE-2026-78663 and Go issue https://go.dev/issue/81743. release notes: https://go.dev/doc/devel/release#go1.26.9 Signed-off-by: Paweł Gronowski --- .github/workflows/codeql.yml | 2 +- .github/workflows/test.yml | 2 +- .github/workflows/validate.yml | 2 +- .golangci.yml | 2 +- Dockerfile | 2 +- dockerfiles/Dockerfile.dev | 2 +- dockerfiles/Dockerfile.lint | 2 +- dockerfiles/Dockerfile.vendor | 2 +- 8 files changed, 8 insertions(+), 8 deletions(-) diff --git a/.github/workflows/codeql.yml b/.github/workflows/codeql.yml index 48e32d6ef1cb..ead4a4f0f270 100644 --- a/.github/workflows/codeql.yml +++ b/.github/workflows/codeql.yml @@ -64,7 +64,7 @@ jobs: name: Update Go uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0 with: - go-version: "1.26.8" + go-version: "1.26.9" cache: false - name: Initialize CodeQL diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index 459abb31639f..6b6a5ea56447 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -68,7 +68,7 @@ jobs: name: Set up Go uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0 with: - go-version: "1.26.8" + go-version: "1.26.9" cache: false - name: Test diff --git a/.github/workflows/validate.yml b/.github/workflows/validate.yml index 25507e1aabab..1396c8d7b0aa 100644 --- a/.github/workflows/validate.yml +++ b/.github/workflows/validate.yml @@ -101,7 +101,7 @@ jobs: name: Set up Go uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0 with: - go-version: "1.26.8" + go-version: "1.26.9" cache: false - name: Run gocompat check diff --git a/.golangci.yml b/.golangci.yml index ad5a06f1d831..ac804566511b 100644 --- a/.golangci.yml +++ b/.golangci.yml @@ -5,7 +5,7 @@ run: # which causes it to fallback to go1.17 semantics. # # TODO(thaJeztah): update "usetesting" settings to enable go1.24 features once our minimum version is go1.24 - go: "1.26.8" + go: "1.26.9" timeout: 5m diff --git a/Dockerfile b/Dockerfile index 02d1fc7d7dca..59f519dcbed1 100644 --- a/Dockerfile +++ b/Dockerfile @@ -8,7 +8,7 @@ ARG BASE_VARIANT=alpine ARG ALPINE_VERSION=3.23 ARG BASE_DEBIAN_DISTRO=bookworm -ARG GO_VERSION=1.26.8 +ARG GO_VERSION=1.26.9 # XX_VERSION specifies the version of the xx utility to use. # It must be a valid tag in the docker.io/tonistiigi/xx image repository. diff --git a/dockerfiles/Dockerfile.dev b/dockerfiles/Dockerfile.dev index 7f02956252e3..c1a7d04c27cc 100644 --- a/dockerfiles/Dockerfile.dev +++ b/dockerfiles/Dockerfile.dev @@ -1,6 +1,6 @@ # syntax=docker/dockerfile:1 -ARG GO_VERSION=1.26.8 +ARG GO_VERSION=1.26.9 # ALPINE_VERSION sets the version of the alpine base image to use, including for the golang image. # It must be a supported tag in the docker.io/library/alpine image repository diff --git a/dockerfiles/Dockerfile.lint b/dockerfiles/Dockerfile.lint index b0404ab8b801..e6acf6f71863 100644 --- a/dockerfiles/Dockerfile.lint +++ b/dockerfiles/Dockerfile.lint @@ -1,6 +1,6 @@ # syntax=docker/dockerfile:1 -ARG GO_VERSION=1.26.8 +ARG GO_VERSION=1.26.9 # ALPINE_VERSION sets the version of the alpine base image to use, including for the golang image. # It must be a supported tag in the docker.io/library/alpine image repository diff --git a/dockerfiles/Dockerfile.vendor b/dockerfiles/Dockerfile.vendor index d781ef983053..442e88286e1f 100644 --- a/dockerfiles/Dockerfile.vendor +++ b/dockerfiles/Dockerfile.vendor @@ -1,6 +1,6 @@ # syntax=docker/dockerfile:1 -ARG GO_VERSION=1.26.8 +ARG GO_VERSION=1.26.9 # ALPINE_VERSION sets the version of the alpine base image to use, including for the golang image. # It must be a supported tag in the docker.io/library/alpine image repository From 949bf401dd8204107463c1fa24aeaeb49349f333 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Pawe=C5=82=20Gronowski?= Date: Thu, 8 Oct 2026 20:25:01 +0200 Subject: [PATCH 2/2] vendor: golang.org/x/net v0.60.0 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit full diff: https://golang.org/x/net/compare/v0.59.0...v0.60.0 Signed-off-by: Paweł Gronowski --- vendor.mod | 2 +- vendor.sum | 4 +- vendor/golang.org/x/net/http2/flow.go | 103 +++++++-- vendor/golang.org/x/net/http2/frame.go | 34 ++- vendor/golang.org/x/net/http2/server.go | 110 +++++---- vendor/golang.org/x/net/http2/transport.go | 142 ++++++++---- .../golang.org/x/net/http2/transport_wrap.go | 214 ++++++++++++++++-- vendor/golang.org/x/net/http2/writesched.go | 7 +- .../x/net/internal/httpcommon/gzip.go | 132 +++++++++++ vendor/modules.txt | 2 +- 10 files changed, 598 insertions(+), 152 deletions(-) create mode 100644 vendor/golang.org/x/net/internal/httpcommon/gzip.go diff --git a/vendor.mod b/vendor.mod index 4a7f5853a80a..2245115c6cf0 100644 --- a/vendor.mod +++ b/vendor.mod @@ -102,7 +102,7 @@ require ( go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.46.0 // indirect go.opentelemetry.io/proto/otlp v1.11.0 // indirect golang.org/x/mod v0.41.0 // indirect - golang.org/x/net v0.59.0 // indirect + golang.org/x/net v0.60.0 // indirect golang.org/x/time v0.16.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260825221802-da73d73af1c5 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260825221802-da73d73af1c5 // indirect diff --git a/vendor.sum b/vendor.sum index bb416723a4cb..6754ef6f6e2e 100644 --- a/vendor.sum +++ b/vendor.sum @@ -193,8 +193,8 @@ golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20201021035429-f5854403a974/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU= -golang.org/x/net v0.59.0 h1:5zfYln+w5XCxwrnMMJPufRgNoXEaGxl0wo5GqPXyues= -golang.org/x/net v0.59.0/go.mod h1:2DA/G1UfVbCpQPeWTmMPGY7Cs2PkBkwu743bVX5PIVg= +golang.org/x/net v0.60.0 h1:79p50tfZlm0J9YfoDsSi639qSXNGVwEzOPLCxM2FsYU= +golang.org/x/net v0.60.0/go.mod h1:2DA/G1UfVbCpQPeWTmMPGY7Cs2PkBkwu743bVX5PIVg= golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= diff --git a/vendor/golang.org/x/net/http2/flow.go b/vendor/golang.org/x/net/http2/flow.go index b7dbd186957e..5066a39912b3 100644 --- a/vendor/golang.org/x/net/http2/flow.go +++ b/vendor/golang.org/x/net/http2/flow.go @@ -10,6 +10,9 @@ package http2 // flow control window update. const inflowMinRefresh = 4 << 10 +// maxFlowWindow is the maximum size of a flow control window. +const maxFlowWindow = (1 << 31) - 1 + // inflow accounts for an inbound flow control window. // It tracks both the latest window sent to the peer (used for enforcement) // and the accumulated unsent window. @@ -74,47 +77,99 @@ func takeInflows(f1, f2 *inflow, n uint32) bool { return true } -// outflow is the outbound flow control window's size. +// connOutflow is connection-level outbound flow control. +type connOutflow struct { + initial int32 // SETTINGS_INITIAL_WINDOW_SIZE, changes with settings updates + n int32 // connection-level flow control window + flowErr bool // set when a flow control error is encountered +} + +func (f *connOutflow) init() { + f.initial = initialWindowSize // initial stream window size + f.n = initialWindowSize // current connection window size +} + +func (f *connOutflow) changeInitialWindowSize(size int64) bool { + if size > maxFlowWindow { + f.flowErr = true + return false + } + f.initial = int32(size) + return true +} + +func (f *connOutflow) add(n int32) bool { + sum := int64(f.n) + int64(n) + if sum > maxFlowWindow { + f.flowErr = true + return false + } + f.n += n + return true +} + +// outflow is the stream-level outbound flow control window's size. type outflow struct { _ incomparable - // n is the number of DATA bytes we're allowed to send. - // An outflow is kept both on a conn and a per-stream. - n int32 + // delta is the difference between the stream's flow control window and + // the connection's initial window size (conn.initial). + // + // Another view is that delta is the number of flow control bytes provided to this + // stream in WINDOW_UPDATE frames, less the number of bytes sent on the stream. + delta int32 // conn points to the shared connection-level outflow that is - // shared by all streams on that conn. It is nil for the outflow - // that's on the conn directly. - conn *outflow + // shared by all streams on that conn. + conn *connOutflow } -func (f *outflow) setConnFlow(cf *outflow) { f.conn = cf } - -func (f *outflow) available() int32 { - n := f.n - if f.conn != nil && f.conn.n < n { - n = f.conn.n +func (f *outflow) available() (int32, bool) { + if f.conn == nil { + return maxFlowWindow, true // only happens in tests } - return n + if f.conn.flowErr { + // Block all sending once any stream observes a flow control error. + return 0, false + } + n := int64(f.conn.initial) + int64(f.delta) + if n > maxFlowWindow { + f.conn.flowErr = true + return 0, false + } + return min(int32(n), f.conn.n), true } func (f *outflow) take(n int32) { - if n > f.available() { - panic("internal error: took too much") + if f.conn == nil { + return // only happens in tests } - f.n -= n - if f.conn != nil { - f.conn.n -= n + avail, _ := f.available() + if n > avail { + panic("internal error: took too much") } + f.delta -= n + f.conn.n -= n } // add adds n bytes (positive or negative) to the flow control window. // It returns false if the sum would exceed 2^31-1. func (f *outflow) add(n int32) bool { - sum := f.n + n - if (sum > n) == (f.n > 0) { - f.n = sum - return true + if f.conn == nil { + return true // only happens in tests + } + avail := int64(f.conn.initial) + int64(f.delta) + if avail > maxFlowWindow { + // An earlier change to the initial window pushed this stream over the limit. + // This is a connection-level flow control error. + f.conn.flowErr = true + return false } - return false + if avail+int64(n) > maxFlowWindow { + // This update would push the stream over the limit. + // This is a stream-level flow control error. + return false + } + f.delta += n + return true } diff --git a/vendor/golang.org/x/net/http2/frame.go b/vendor/golang.org/x/net/http2/frame.go index be75badcc0ef..3363bbafc858 100644 --- a/vendor/golang.org/x/net/http2/frame.go +++ b/vendor/golang.org/x/net/http2/frame.go @@ -1711,7 +1711,8 @@ func (fr *Framer) readMetaFrame(hf *HeadersFrame) (Frame, error) { mh := &MetaHeadersFrame{ HeadersFrame: hf, } - var remainSize = fr.maxHeaderListSize() + var headersRemainSize = fr.maxHeaderListSize() + var trailersRemainSize = fr.maxHeaderListSize() var sawRegular bool var invalid error // pseudo header field errors @@ -1743,14 +1744,29 @@ func (fr *Framer) readMetaFrame(hf *HeadersFrame) (Frame, error) { return } - size := hf.Size() - if size > remainSize { + var remainSize *uint32 + var size uint32 + if hf.Name == "trailer" { + remainSize = &trailersRemainSize + fieldCount := strings.Count(hf.Value, ",") + 1 + // Rather than actually constructing hpack.HeaderField for each + // trailer field and cumulatively adding its Size, just do the math + // manually to avoid unnecessary work. This does make it so + // whitespaces after comma are counted against the budget, but that + // should be innocuous. + size = uint32(len(hf.Value)-fieldCount+1) + uint32(fieldCount)*hpack.HeaderField{}.Size() + } else { + remainSize = &headersRemainSize + size = hf.Size() + } + if size > *remainSize { hdec.SetEmitEnabled(false) mh.Truncated = true - remainSize = 0 + headersRemainSize = 0 + trailersRemainSize = 0 return } - remainSize -= size + *remainSize -= size mh.Fields = append(mh.Fields, hf) }) @@ -1766,10 +1782,10 @@ func (fr *Framer) readMetaFrame(hf *HeadersFrame) (Frame, error) { // skip parsing the fragment and close the connection. // // "Too much" is either any CONTINUATION frame after we've already - // exceeded the max header list size (in which case remainSize is 0), - // or a frame whose encoded size is more than twice the remaining - // header list bytes we're willing to accept. - if int64(len(frag)) > int64(2*remainSize) { + // exceeded the max header list size (if so, both budgets are 0), or a + // frame whose encoded size is more than twice the remaining header + // list bytes we're willing to accept. + if int64(len(frag)) > 2*int64(headersRemainSize+trailersRemainSize) { if VerboseLogs { log.Printf("http2: header list too large") } diff --git a/vendor/golang.org/x/net/http2/server.go b/vendor/golang.org/x/net/http2/server.go index a7d2053b6c79..a6cd66226d06 100644 --- a/vendor/golang.org/x/net/http2/server.go +++ b/vendor/golang.org/x/net/http2/server.go @@ -280,7 +280,6 @@ func (s *Server) serveConn(c net.Conn, opts *ServeConnOpts, newf func(*serverCon doneServing: make(chan struct{}), clientMaxStreams: math.MaxUint32, // Section 6.5.2: "Initially, there is no limit to this value" advMaxStreams: conf.MaxConcurrentStreams, - initialStreamSendWindowSize: initialWindowSize, initialStreamRecvWindowSize: conf.MaxUploadBufferPerStream, maxFrameSize: initialMaxFrameSize, pingTimeout: conf.PingTimeout, @@ -320,7 +319,7 @@ func (s *Server) serveConn(c net.Conn, opts *ServeConnOpts, newf func(*serverCon // These start at the RFC-specified defaults. If there is a higher // configured value for inflow, that will be updated when we send a // WINDOW_UPDATE shortly after sending SETTINGS. - sc.flow.add(initialWindowSize) + sc.flow.init() sc.inflow.init(initialWindowSize) sc.hpackEncoder = hpack.NewEncoder(&sc.headerWriteBuf) sc.hpackEncoder.SetMaxDynamicTableSizeLimit(conf.MaxEncoderHeaderTableSize) @@ -427,7 +426,7 @@ type serverConn struct { wroteFrameCh chan frameWriteResult // from writeFrameAsync -> serve, tickles more frame writes bodyReadCh chan bodyReadMsg // from handlers -> serve serveMsgCh chan interface{} // misc messages & code to send to / run on the serve loop - flow outflow // conn-wide (not stream-specific) outbound flow control + flow connOutflow // conn-wide (not stream-specific) outbound flow control inflow inflow // conn-wide inbound flow control tlsState *tls.ConnectionState // shared by all handlers, like net/http remoteAddrStr string @@ -441,6 +440,9 @@ type serverConn struct { sawFirstSettings bool // got the initial SETTINGS frame after the preface needToSendSettingsAck bool unackedSettings int // how many SETTINGS have we sent without ACKs? + pendingEncoderTableSize bool // peer changed SETTINGS_HEADER_TABLE_SIZE; apply to hpackEncoder before the next frame write + encoderTableSizeMin uint32 // smallest SETTINGS_HEADER_TABLE_SIZE since the last apply + encoderTableSize uint32 // latest SETTINGS_HEADER_TABLE_SIZE queuedControlFrames int // control frames in the writeSched queue clientMaxStreams uint32 // SETTINGS_MAX_CONCURRENT_STREAMS from client (our PUSH_PROMISE limit) advMaxStreams uint32 // our SETTINGS_MAX_CONCURRENT_STREAMS advertised the client @@ -451,7 +453,6 @@ type serverConn struct { maxPushPromiseID uint32 // ID of the last push promise (even), or 0 if there have been no pushes streams map[uint32]*stream unstartedHandlers []unstartedHandler - initialStreamSendWindowSize int32 initialStreamRecvWindowSize int32 maxFrameSize int32 peerMaxHeaderListSize uint32 // zero means unknown (default) @@ -521,7 +522,8 @@ type stream struct { // immutable: sc *serverConn id uint32 - body *pipe // non-nil if expecting DATA frames + body *pipe // non-nil if expecting DATA frames + reqBody *requestBody cw closeWaiter // closed wait stream transitions to closed state ctx context.Context cancelCtx func() @@ -1172,6 +1174,16 @@ func (sc *serverConn) startFrameWrite(wr FrameWriteRequest) { sc.writingFrame = true sc.needsFrameFlush = true + if sc.pendingEncoderTableSize { + // hpackEncoder may be in use by writeFrameAsync, so SETTINGS + // changes to it are deferred until no frame is being written. + // Replaying the smallest size before the latest one keeps the + // encoder's view identical to having applied every change + // (RFC 7541, Section 4.2). + sc.pendingEncoderTableSize = false + sc.hpackEncoder.SetMaxDynamicTableSize(sc.encoderTableSizeMin) + sc.hpackEncoder.SetMaxDynamicTableSize(sc.encoderTableSize) + } if wr.write.staysWithinBuffer(sc.bw.Available()) { sc.writingFrameAsync = false err := wr.write.writeFrame(sc) @@ -1271,6 +1283,11 @@ func (sc *serverConn) scheduleFrameWrite() { } sc.inFrameScheduleLoop = true for !sc.writingFrameAsync { + if sc.flow.flowErr && (!sc.inGoAway || sc.goAwayCode == ErrCodeNo) { + sc.inGoAway = true + sc.needToSendGoAway = true + sc.goAwayCode = ErrCodeFlowControl + } if sc.needToSendGoAway { sc.needToSendGoAway = false sc.startFrameWrite(FrameWriteRequest{ @@ -1294,6 +1311,9 @@ func (sc *serverConn) scheduleFrameWrite() { sc.startFrameWrite(wr) continue } + if sc.flow.flowErr { + continue + } } if sc.needsFrameFlush { sc.startFrameWrite(FrameWriteRequest{write: flushFrameWriter{}}) @@ -1526,6 +1546,10 @@ func (sc *serverConn) processWindowUpdate(f *WindowUpdateFrame) error { return nil } if !st.flow.add(int32(f.Increment)) { + if st.flow.conn.flowErr { + // This is a lazily-detected connection-level flow control error. + return sc.countError("bad_flow", ConnectionError(ErrCodeFlowControl)) + } return sc.countError("bad_flow", streamError(f.StreamID, ErrCodeFlowControl)) } default: // connection-level flow control @@ -1584,10 +1608,6 @@ func (sc *serverConn) closeStream(st *stream, err error) { } } if p := st.body; p != nil { - // Return any buffered unread bytes worth of conn-level flow control. - // See golang.org/issue/16481 - sc.sendWindowUpdate(nil, p.Len()) - p.CloseWithError(err) } if e, ok := err.(StreamError); ok { @@ -1641,7 +1661,12 @@ func (sc *serverConn) processSetting(s Setting) error { } switch s.ID { case SettingHeaderTableSize: - sc.hpackEncoder.SetMaxDynamicTableSize(s.Val) + // Applied by startFrameWrite; see comment there. + if !sc.pendingEncoderTableSize || s.Val < sc.encoderTableSizeMin { + sc.encoderTableSizeMin = s.Val + } + sc.encoderTableSize = s.Val + sc.pendingEncoderTableSize = true case SettingEnablePush: sc.pushEnabled = s.Val != 0 case SettingMaxConcurrentStreams: @@ -1672,28 +1697,14 @@ func (sc *serverConn) processSetting(s Setting) error { func (sc *serverConn) processSettingInitialWindowSize(val uint32) error { sc.serveG.check() - // Note: val already validated to be within range by - // processSetting's Valid call. - - // "A SETTINGS frame can alter the initial flow control window - // size for all current streams. When the value of - // SETTINGS_INITIAL_WINDOW_SIZE changes, a receiver MUST - // adjust the size of all stream flow control windows that it - // maintains by the difference between the new value and the - // old value." - old := sc.initialStreamSendWindowSize - sc.initialStreamSendWindowSize = int32(val) - growth := int32(val) - old // may be negative - for _, st := range sc.streams { - if !st.flow.add(growth) { - // 6.9.2 Initial Flow Control Window Size - // "An endpoint MUST treat a change to - // SETTINGS_INITIAL_WINDOW_SIZE that causes any flow - // control window to exceed the maximum size as a - // connection error (Section 5.4.1) of type - // FLOW_CONTROL_ERROR." - return sc.countError("setting_win_size", ConnectionError(ErrCodeFlowControl)) - } + if !sc.flow.changeInitialWindowSize(int64(val)) { + // 6.9.2 Initial Flow Control Window Size + // "An endpoint MUST treat a change to + // SETTINGS_INITIAL_WINDOW_SIZE that causes any flow + // control window to exceed the maximum size as a + // connection error (Section 5.4.1) of type + // FLOW_CONTROL_ERROR." + return sc.countError("setting_win_size", ConnectionError(ErrCodeFlowControl)) } return nil } @@ -1966,7 +1977,7 @@ func (sc *serverConn) processHeaders(f *MetaHeadersFrame) error { if st.reqTrailer != nil { st.trailer = make(http.Header) } - st.body = req.Body.(*requestBody).pipe // may be nil + st.body = st.reqBody.pipe // may be nil st.declBodyBytes = req.ContentLength handler := sc.handler.ServeHTTP @@ -1989,7 +2000,7 @@ func (sc *serverConn) processHeaders(f *MetaHeadersFrame) error { st.readDeadline = time.AfterFunc(sc.hs.ReadTimeout, st.onReadTimeout) } - return sc.scheduleHandler(id, rw, req, handler) + return sc.scheduleHandler(st, rw, req, handler) } func (sc *serverConn) upgradeRequest(req *http.Request) { @@ -2102,7 +2113,6 @@ func (sc *serverConn) newStream(id, pusherID uint32, state streamState, priority } st.cw.Init() st.flow.conn = &sc.flow // link to conn-level counter - st.flow.add(sc.initialStreamSendWindowSize) st.inflow.init(sc.initialStreamRecvWindowSize) if sc.hs.WriteTimeout > 0 { st.writeDeadline = time.AfterFunc(sc.hs.WriteTimeout, st.onWriteTimeout) @@ -2184,7 +2194,7 @@ func (sc *serverConn) newWriterAndRequest(st *stream, f *MetaHeadersFrame) (*res } else { req.ContentLength = -1 } - req.Body.(*requestBody).pipe = &pipe{ + st.reqBody.pipe = &pipe{ b: &dataBuffer{expected: req.ContentLength}, } } @@ -2204,7 +2214,7 @@ func (sc *serverConn) newWriterAndRequestNoBody(st *stream, rp httpcommon.Server return nil, nil, sc.countError(res.InvalidReason, streamError(st.id, ErrCodeProtocol)) } - body := &requestBody{ + st.reqBody = &requestBody{ conn: sc, stream: st, needsContinue: res.NeedsContinue, @@ -2220,7 +2230,7 @@ func (sc *serverConn) newWriterAndRequestNoBody(st *stream, rp httpcommon.Server ProtoMinor: 0, TLS: tlsState, Host: rp.Authority, - Body: body, + Body: st.reqBody, Trailer: res.Trailer, }).WithContext(st.ctx) rw := sc.newResponseWriter(st, req) @@ -2244,11 +2254,12 @@ type unstartedHandler struct { rw *responseWriter req *http.Request handler func(http.ResponseWriter, *http.Request) + body *pipe } // scheduleHandler starts a handler goroutine, // or schedules one to start as soon as an existing handler finishes. -func (sc *serverConn) scheduleHandler(streamID uint32, rw *responseWriter, req *http.Request, handler func(http.ResponseWriter, *http.Request)) error { +func (sc *serverConn) scheduleHandler(st *stream, rw *responseWriter, req *http.Request, handler func(http.ResponseWriter, *http.Request)) error { sc.serveG.check() maxHandlers := sc.advMaxStreams if sc.curHandlers < maxHandlers { @@ -2260,10 +2271,11 @@ func (sc *serverConn) scheduleHandler(streamID uint32, rw *responseWriter, req * return sc.countError("too_many_early_resets", ConnectionError(ErrCodeEnhanceYourCalm)) } sc.unstartedHandlers = append(sc.unstartedHandlers, unstartedHandler{ - streamID: streamID, + streamID: st.id, rw: rw, req: req, handler: handler, + body: st.body, }) return nil } @@ -2277,6 +2289,10 @@ func (sc *serverConn) handlerDone() { u := sc.unstartedHandlers[i] if sc.streams[u.streamID] == nil { // This stream was reset before its goroutine had a chance to start. + if u.body != nil { + u.body.BreakWithError(errClosedBody) + sc.sendWindowUpdate(nil, u.body.Len()) + } continue } if sc.curHandlers >= maxHandlers { @@ -2298,6 +2314,12 @@ func (sc *serverConn) runHandler(rw *responseWriter, req *http.Request, handler didPanic := true defer func() { rw.rws.stream.cancelCtx() + if b := rw.rws.stream.reqBody; b != nil { + // Closing the body refunds flow control credit for any unconsumed data. + // (reqBody is nil for Upgrade: h2c requests, but those do not use flow + // control for the request body.) + b.Close() + } if req.MultipartForm != nil { req.MultipartForm.RemoveAll() } @@ -2396,7 +2418,7 @@ func (sc *serverConn) noteBodyReadFromHandler(st *stream, n int, err error) { func (sc *serverConn) noteBodyRead(st *stream, n int) { sc.serveG.check() sc.sendWindowUpdate(nil, n) // conn-level - if st.state != stateHalfClosedRemote && st.state != stateClosed { + if st != nil && st.state != stateHalfClosedRemote && st.state != stateClosed { // Don't send this WINDOW_UPDATE if the stream is closed // remotely. sc.sendWindowUpdate(st, n) @@ -2444,6 +2466,9 @@ func (b *requestBody) Close() error { b.closeOnce.Do(func() { if b.pipe != nil { b.pipe.BreakWithError(errClosedBody) + if unread := b.pipe.Len(); unread > 0 { + b.conn.noteBodyReadFromHandler(nil, unread, errClosedBody) + } } }) return nil @@ -2461,9 +2486,6 @@ func (b *requestBody) Read(p []byte) (n int, err error) { if err == io.EOF { b.sawEOF = true } - if b.conn == nil { - return - } b.conn.noteBodyReadFromHandler(b.stream, n, err) return } diff --git a/vendor/golang.org/x/net/http2/transport.go b/vendor/golang.org/x/net/http2/transport.go index 088eb9bef8e1..d894c6019c50 100644 --- a/vendor/golang.org/x/net/http2/transport.go +++ b/vendor/golang.org/x/net/http2/transport.go @@ -27,6 +27,7 @@ import ( "net/http" "net/http/httptrace" "net/textproto" + "slices" "strconv" "sync" "sync/atomic" @@ -203,7 +204,7 @@ type ClientConn struct { mu sync.Mutex // guards following cond *sync.Cond // hold mu; broadcast on flow/closed changes - flow outflow // our conn-level flow control quota (cs.outflow is per stream) + flow connOutflow // our conn-level flow control quota (cs.outflow is per stream) inflow inflow // peer's conn-level flow control doNotReuse bool // whether conn is marked to not be reused for any future requests closing bool @@ -227,7 +228,6 @@ type ClientConn struct { maxConcurrentStreams uint32 peerMaxHeaderListSize uint64 peerMaxHeaderTableSize uint32 - initialWindowSize uint32 initialStreamRecvWindowSize int32 readIdleTimeout time.Duration pingTimeout time.Duration @@ -472,7 +472,6 @@ func (t *Transport) newClientConn(c net.Conn, singleUse bool, internalStateHook readerDone: make(chan struct{}), nextStreamID: 1, maxFrameSize: 16 << 10, // spec default - initialWindowSize: 65535, // spec default initialStreamRecvWindowSize: conf.MaxUploadBufferPerStream, maxConcurrentStreams: initialMaxConcurrentStreams, // "infinite", per spec. Use a smaller value until we have received server settings. strictMaxConcurrentStreams: conf.StrictMaxConcurrentRequests, @@ -497,7 +496,7 @@ func (t *Transport) newClientConn(c net.Conn, singleUse bool, internalStateHook } cc.cond = sync.NewCond(&cc.mu) - cc.flow.add(int32(initialWindowSize)) + cc.flow.init() // TODO: adjust this writer size to account for frame size + // MTU + crypto/tls record padding. @@ -876,6 +875,33 @@ func (cc *ClientConn) sendGoAway() error { return nil } +func (cc *ClientConn) goAwayAndClose(code ErrCode) { + cc.mu.Lock() + closed := cc.closed + cc.closing = true + cc.closed = true + cc.mu.Unlock() + if closed { + return + } + if f := cc.fr.countError; f != nil { + f(fmt.Sprintf("conn_close_error_%s", code.stringToken())) + } + done := make(chan struct{}) + go func() { + defer close(done) + cc.wmu.Lock() + cc.fr.WriteGoAway(0, code, nil) + cc.bw.Flush() + cc.wmu.Unlock() + }() + select { + case <-done: + case <-time.After(250 * time.Millisecond): + } + cc.closeForError(fmt.Errorf("http2: closing connection with %v", code)) +} + // closes the client connection immediately. In-flight requests are interrupted. // err is sent to streams. func (cc *ClientConn) closeForError(err error) { @@ -1319,7 +1345,9 @@ func (cs *clientStream) cleanupWriteRequest(err error) { } if err != nil { cs.abortStream(err) // possibly redundant, but harmless - if cs.sentHeaders { + if ce, ok := err.(ConnectionError); ok { + cc.goAwayAndClose(ErrCode(ce)) + } else if cs.sentHeaders { if se, ok := err.(StreamError); ok { if se.Cause != errFromPeer { cc.writeStreamReset(cs.ID, se.Code, false, err) @@ -1648,8 +1676,12 @@ func (cs *clientStream) awaitFlowControl(maxBytes int) (taken int32, err error) return 0, errRequestCanceled default: } - if a := cs.flow.available(); a > 0 { - take := a + avail, ok := cs.flow.available() + if !ok { + return 0, ConnectionError(ErrCodeFlowControl) + } + if avail > 0 { + take := avail if int(take) > maxBytes { take = int32(maxBytes) // can't truncate int; take is int32 @@ -1710,8 +1742,7 @@ type resAndError struct { // requires cc.mu be held. func (cc *ClientConn) addStreamLocked(cs *clientStream) { - cs.flow.add(int32(cc.initialWindowSize)) - cs.flow.setConnFlow(&cc.flow) + cs.flow.conn = &cc.flow cs.inflow.init(cc.initialStreamRecvWindowSize) cs.ID = cc.nextStreamID cc.nextStreamID += 2 @@ -2095,18 +2126,42 @@ func (rl *clientConnReadLoop) handleResponse(cs *clientStream, f *MetaHeadersFra return nil, nil } + // Delete various headers that might mess up framing for HTTP/1. This is + // not a problem for HTTP/2, but someone might use HTTP/2 transport as a + // reverse proxy which forwards the response to an HTTP/1 client. Our + // HTTP/1 client transport will properly reject improper headers such as + // multiple conflicting Content-Length headers, but other implementations + // might not. + // TODO: just reject such responses? We deleted them for compatibility + // since this was done in a security fix (go.dev/issue/81115). However, + // rejecting them seems entirely reasonable and relatively safe. + + // Connection-specific header fields must not appear in an HTTP/2 message, + // and any message containing them is malformed. RFC 9113, Section 8.2.2. + for _, k := range connHeaders { + delete(res.Header, k) + } res.ContentLength = -1 - if clens := res.Header["Content-Length"]; len(clens) == 1 { - if cl, err := strconv.ParseUint(clens[0], 10, 63); err == nil { - res.ContentLength = int64(cl) + if clens, ok := res.Header["Content-Length"]; ok { + // Repeated Content-Length values may be collapsed into one only if + // they are identical per RFC 9110 Section 8.6. + // No need to trim whitespace, HTTP/2 header values must not have + // extraneous whitespace per RFC 9113 Section 8.2.1. + // Non-canonical headers are already rejected by our framer at this + // point. + conflicting := slices.ContainsFunc(clens[1:], func(clen string) bool { + return clen != clens[0] + }) + cl, err := strconv.ParseUint(clens[0], 10, 63) + if conflicting || err != nil { + delete(res.Header, "Content-Length") } else { - // TODO: care? unlike http/1, it won't mess up our framing, so it's - // more safe smuggling-wise to ignore. + res.Header["Content-Length"] = clens[:1] + res.ContentLength = int64(cl) } - } else if len(clens) > 1 { - // TODO: care? unlike http/1, it won't mess up our framing, so it's - // more safe smuggling-wise to ignore. - } else if f.StreamEnded() && !cs.isHead { + } + + if res.ContentLength < 0 && f.StreamEnded() && !cs.isHead { res.ContentLength = 0 } @@ -2417,6 +2472,10 @@ const ( func (rl *clientConnReadLoop) streamByID(id uint32, headerOrData bool) *clientStream { rl.cc.mu.Lock() defer rl.cc.mu.Unlock() + return rl.streamByIDLocked(id, headerOrData) +} + +func (rl *clientConnReadLoop) streamByIDLocked(id uint32, headerOrData bool) *clientStream { if headerOrData { // Work around an unfortunate gRPC behavior. // See comment on ClientConn.rstStreamPingsBlocked for details. @@ -2499,24 +2558,10 @@ func (rl *clientConnReadLoop) processSettingsNoWrite(f *SettingsFrame) error { case SettingMaxHeaderListSize: cc.peerMaxHeaderListSize = uint64(s.Val) case SettingInitialWindowSize: - // Values above the maximum flow-control - // window size of 2^31-1 MUST be treated as a - // connection error (Section 5.4.1) of type - // FLOW_CONTROL_ERROR. - if s.Val > math.MaxInt32 { + if !cc.flow.changeInitialWindowSize(int64(s.Val)) { return ConnectionError(ErrCodeFlowControl) } - - // Adjust flow control of currently-open - // frames by the difference of the old initial - // window size and this one. - delta := int32(s.Val) - int32(cc.initialWindowSize) - for _, cs := range cc.streams { - cs.flow.add(delta) - } cc.cond.Broadcast() - - cc.initialWindowSize = s.Val case SettingHeaderTableSize: cc.henc.SetMaxDynamicTableSize(s.Val) cc.peerMaxHeaderTableSize = s.Val @@ -2558,29 +2603,30 @@ func (rl *clientConnReadLoop) processSettingsNoWrite(f *SettingsFrame) error { func (rl *clientConnReadLoop) processWindowUpdate(f *WindowUpdateFrame) error { cc := rl.cc - cs := rl.streamByID(f.StreamID, notHeaderOrDataFrame) - if f.StreamID != 0 && cs == nil { - return nil - } - cc.mu.Lock() defer cc.mu.Unlock() - - fl := &cc.flow - if cs != nil { - fl = &cs.flow - } - if !fl.add(int32(f.Increment)) { - // For stream, the sender sends RST_STREAM with an error code of FLOW_CONTROL_ERROR - if cs != nil { + if f.StreamID == 0 { + if !cc.flow.add(int32(f.Increment)) { + return ConnectionError(ErrCodeFlowControl) + } + } else { + cs := rl.streamByIDLocked(f.StreamID, notHeaderOrDataFrame) + if cs == nil { + return nil + } + if !cs.flow.add(int32(f.Increment)) { + if cs.flow.conn.flowErr { + // This is a lazily-detected connection-level flow control error. + return ConnectionError(ErrCodeFlowControl) + } + // For stream, the sender sends RST_STREAM with + // an error code of FLOW_CONTROL_ERROR. rl.endStreamErrorLocked(cs, StreamError{ StreamID: f.StreamID, Code: ErrCodeFlowControl, }) return nil } - - return ConnectionError(ErrCodeFlowControl) } cc.cond.Broadcast() return nil diff --git a/vendor/golang.org/x/net/http2/transport_wrap.go b/vendor/golang.org/x/net/http2/transport_wrap.go index 34dfb3cca624..eeaa95f03cb6 100644 --- a/vendor/golang.org/x/net/http2/transport_wrap.go +++ b/vendor/golang.org/x/net/http2/transport_wrap.go @@ -119,33 +119,27 @@ type http2TransportContextKey struct{} // DialFromContext dials a new connection using the http2.Transport's DialTLS/DialTLSContext. func (t transportConfig) DialFromContext(ctx context.Context, network, address string) (net.Conn, error) { - if ctx.Value(http2TransportContextKey{}) == nil { + dial, _ := ctx.Value(http2TransportContextKey{}).(*transportRoundTripState) + if dial == nil { // We're being called from a RoundTrip that did not start with an http2.Transport. // Use the http.Transport's dialer. return nil, errors.ErrUnsupported } - - tlsConf := t.t.TLSClientConfig - if tlsConf == nil { - tlsConf = &tls.Config{} - } else { - tlsConf = tlsConf.Clone() - } - if !slices.Contains(tlsConf.NextProtos, "h2") { - tlsConf.NextProtos = append([]string{"h2"}, tlsConf.NextProtos...) - } - if tlsConf.ServerName == "" { - host, _, err := net.SplitHostPort(address) - if err == nil { - tlsConf.ServerName = host - } - } - return t.t.dialTLS(ctx, network, address, tlsConf) + return t.t.dialForRoundTrip(ctx, dial, network, address) } type transportInternal struct { initOnce sync.Once lazyt1 *http.Transport + + dialMu sync.Mutex + dials map[string]*transportDialState +} + +type transportDialState struct { + dialc chan struct{} // close when dial returns + dialErr error // dial result, set before dialc is closed + rtdonec chan struct{} // closed when RoundTrip initiating the dial returns } func (t *Transport) init() *http.Transport { @@ -167,6 +161,52 @@ func (t *Transport) configure(t1 *http.Transport) { } } +// transportRoundTripState is the state of the dial for an http2.Transport.RoundTrip. +type transportRoundTripState struct { + mu sync.Mutex + rtdone bool // set when RoundTrip returns + rtdonec chan struct{} // closed when RoundTrip returns + gotconnc chan struct{} // closed when GotConn hook is called +} + +// startDial is called when a dial starts. +// +// It returns done=true if the RoundTrip has already returned, +// in which case we should skip dialing. +func (dial *transportRoundTripState) startDial() (gotconnc chan struct{}, done bool) { + dial.mu.Lock() + defer dial.mu.Unlock() + if dial.rtdone { + return nil, true + } + if dial.rtdonec == nil { + dial.rtdonec = make(chan struct{}) + } + dial.gotconnc = make(chan struct{}) + return dial.gotconnc, false +} + +// gotConn is called when RoundTrip gets a connection. +// This may happen multiple times per RoundTrip, if a request is retried. +func (dial *transportRoundTripState) gotConn() { + dial.mu.Lock() + defer dial.mu.Unlock() + if dial.gotconnc != nil { + close(dial.gotconnc) + dial.gotconnc = nil + } +} + +// roundTripDone is called when RoundTrip returns. +func (dial *transportRoundTripState) roundTripDone() { + dial.mu.Lock() + defer dial.mu.Unlock() + dial.rtdone = true + if dial.rtdonec != nil { + close(dial.rtdonec) + } +} + func (t *Transport) roundTripOpt(req *http.Request, opt RoundTripOpt) (*http.Response, error) { t1 := t.init() @@ -186,12 +226,146 @@ func (t *Transport) roundTripOpt(req *http.Request, opt RoundTripOpt) (*http.Res // Both http.Transport and http2.Transport allow the user to provide a custom // dial function, and historically you only get the dial function from the // Transport you're calling RoundTrip on. - ctx := context.WithValue(req.Context(), http2TransportContextKey{}, t) + // + // In addition, http2.Transport coalesces dials, which http.Transport historically + // has not. + dial := &transportRoundTripState{} + defer dial.roundTripDone() + ctx := context.WithValue(req.Context(), http2TransportContextKey{}, dial) + ctx = httptrace.WithClientTrace(ctx, &httptrace.ClientTrace{ + GotConn: func(httptrace.GotConnInfo) { + // Tell dialForRoundTrip that we have received a connection. + dial.gotConn() + }, + }) req = req.WithContext(ctx) - return t1.RoundTrip(req) } +var errCoalescedDialAbandoned = errors.New("http2: abandoned coalesced dial") + +const coalescedDialRetryTimeout = 50 * time.Millisecond + +// dialForRoundTrip is called (indirectly) by net/http when dialing a new connection +// for a request call which originated as an http2.Transport.RoundTrip call. +// +// It uses the http2.Transport's Dial hooks and coalesces dials. +func (t *Transport) dialForRoundTrip(ctx context.Context, dial *transportRoundTripState, network, address string) (net.Conn, error) { + gotconnc, rtdone := dial.startDial() + if rtdone { + // RoundTrip returned, no need for this dial to proceed. + return nil, errCoalescedDialAbandoned + } + + var state *transportDialState + for { + // The first dial to an address registers itself with t.dials. + // When the dial finishes, it records the outcome and removes itself from t.dials + // so future dials will start a new round of coalescing. + // + // Subsequent dials observe that a dial is in progress, and skip dialing. + t.dialMu.Lock() + if t.dials == nil { + t.dials = make(map[string]*transportDialState) + } + leader := t.dials[address] + if leader == nil { + // No entry in t.dials to coalesce with. Add ourselves as leader. + state = &transportDialState{ + dialc: make(chan struct{}), + rtdonec: dial.rtdonec, + } + t.dials[address] = state + } + t.dialMu.Unlock() + if leader == nil { + // We are the leader, so we should dial for real. + break + } + + // Coalesce with a previous dial. + var ( + done = false + leaderdialc = leader.dialc + leaderrtdonec = leader.rtdonec + timerc <-chan time.Time + ) + for !done { + if leaderdialc == nil && leaderrtdonec == nil && timerc == nil { + // The leader finished dialing, and its RoundTrip finished. + // + // There are several possibilities, which reduce to: + // - The RoundTrip we are dialing for is about to receive + // a connection (possibly the one the leader just dialed), + // but we haven't observed it yet. + // - There's something wrong with the leader's connection, + // such as a TLS handshake failure. + // + // We have no good way to distinguish between these cases. + // + // Wait a short time for the leader's connection (or some + // other usable connection) to be delivered to our RoundTrip. + // If it is not, start a dial of our own. + timerc = time.After(coalescedDialRetryTimeout) + } + select { + case <-leaderdialc: + // The leader's dial completed. + if leader.dialErr != nil { + return nil, leader.dialErr + } + leaderdialc = nil + case <-leaderrtdonec: + // The leader's RoundTrip completed. + leaderrtdonec = nil + case <-timerc: + // The leader's dial and RoundTrip completed, and some time + // has passed. Assume we're not getting a connection and + // restart the dial process. + done = true + case <-dial.rtdonec: + // Our RoundTrip returned. + return nil, errCoalescedDialAbandoned + case <-gotconnc: + // Our RoundTrip got a connection. + return nil, errCoalescedDialAbandoned + case <-ctx.Done(): + // net/http abandoned the dial. + return nil, ctx.Err() + } + } + } + + // Dial for real. + tlsConf := t.TLSClientConfig + if tlsConf == nil { + tlsConf = &tls.Config{} + } else { + tlsConf = tlsConf.Clone() + } + if !slices.Contains(tlsConf.NextProtos, "h2") { + tlsConf.NextProtos = append([]string{"h2"}, tlsConf.NextProtos...) + } + if tlsConf.ServerName == "" { + host, _, err := net.SplitHostPort(address) + if err == nil { + tlsConf.ServerName = host + } + } + nc, err := t.dialTLS(ctx, network, address, tlsConf) + + // Remove this dial from t.dials, so future dial attempts will not coalesce with it. + t.dialMu.Lock() + delete(t.dials, address) + t.dialMu.Unlock() + + // Notify any dial coalesced with this one that we are done. + state.dialErr = err + close(state.dialc) + + return nc, err +} + func (t *Transport) closeIdleConnections() { t1 := t.init() t1.CloseIdleConnections() diff --git a/vendor/golang.org/x/net/http2/writesched.go b/vendor/golang.org/x/net/http2/writesched.go index 36ad5e32b4d8..f1903a8c66af 100644 --- a/vendor/golang.org/x/net/http2/writesched.go +++ b/vendor/golang.org/x/net/http2/writesched.go @@ -79,10 +79,11 @@ func (wr FrameWriteRequest) Consume(n int32) (FrameWriteRequest, FrameWriteReque } // Might need to split after applying limits. - allowed := wr.stream.flow.available() - if n < allowed { - allowed = n + avail, ok := wr.stream.flow.available() + if !ok { + return empty, empty, 0 } + allowed := min(n, avail) if wr.stream.sc.maxFrameSize < allowed { allowed = wr.stream.sc.maxFrameSize } diff --git a/vendor/golang.org/x/net/internal/httpcommon/gzip.go b/vendor/golang.org/x/net/internal/httpcommon/gzip.go new file mode 100644 index 000000000000..819dd650cb6e --- /dev/null +++ b/vendor/golang.org/x/net/internal/httpcommon/gzip.go @@ -0,0 +1,132 @@ +// Copyright 2026 The Go Authors. All rights reserved. +// Use of this source code is governed by a BSD-style +// license that can be found in the LICENSE file. + +package httpcommon + +import ( + "compress/flate" + "compress/gzip" + "errors" + "io" + "io/fs" + "sync" +) + +var errConcurrentRead = errors.New("http: concurrent read on response body") + +// incomparable is a zero-width, non-comparable type. Adding it to a struct +// makes that struct also non-comparable, and generally doesn't add +// any size (as long as it's first). +type incomparable [0]func() + +// GzipReader wraps a response body so it can lazily +// get gzip.Reader from the pool on the first call to Read. +// After Close is called it puts gzip.Reader to the pool immediately +// if there is no Read in progress or later when Read completes. +type GzipReader struct { + _ incomparable + Body io.ReadCloser // underlying Response.Body + mu sync.Mutex // guards zr and zerr + zr *gzip.Reader // stores gzip reader from the pool between reads + zerr error // sticky gzip reader init error or sentinel value to detect concurrent read and read after close +} + +type eofReader struct{} + +func (eofReader) Read([]byte) (int, error) { return 0, io.EOF } +func (eofReader) ReadByte() (byte, error) { return 0, io.EOF } + +var gzipPool = sync.Pool{New: func() any { return new(gzip.Reader) }} + +// gzipPoolGet gets a gzip.Reader from the pool and resets it to read from r. +func gzipPoolGet(r io.Reader) (*gzip.Reader, error) { + zr := gzipPool.Get().(*gzip.Reader) + if err := zr.Reset(r); err != nil { + gzipPoolPut(zr) + return nil, err + } + return zr, nil +} + +// gzipPoolPut puts a gzip.Reader back into the pool. +func gzipPoolPut(zr *gzip.Reader) { + // Reset will allocate bufio.Reader if we pass it anything + // other than a flate.Reader, so ensure that it's getting one. + var r flate.Reader = eofReader{} + zr.Reset(r) + gzipPool.Put(zr) +} + +// acquire returns a gzip.Reader for reading response body. +// The reader must be released after use. +func (gz *GzipReader) acquire() (*gzip.Reader, error) { + gz.mu.Lock() + defer gz.mu.Unlock() + if gz.zerr != nil { + return nil, gz.zerr + } + if gz.zr == nil { + // gzipPoolGet might block indefinitely since it reads the gzip header. + // Therefore, drop mu temporarily when using gzipPoolGet. + // We set zerr to errConcurrentRead to prevent concurrent read + // even when mu is temporarily dropped. + gz.zerr = errConcurrentRead + gz.mu.Unlock() + zr, err := gzipPoolGet(gz.Body) + gz.mu.Lock() + // Guard against Close being called while gzipPoolGet is running. + if gz.zerr != errConcurrentRead { + if zr != nil { + gzipPoolPut(zr) + } + return nil, gz.zerr + } + gz.zr, gz.zerr = zr, err + if gz.zerr != nil { + return nil, gz.zerr + } + } + ret := gz.zr + gz.zr, gz.zerr = nil, errConcurrentRead + return ret, nil +} + +// release returns the gzip.Reader to the pool if Close was called during Read. +func (gz *GzipReader) release(zr *gzip.Reader) { + gz.mu.Lock() + defer gz.mu.Unlock() + if gz.zerr == errConcurrentRead { + gz.zr, gz.zerr = zr, nil + } else { // fs.ErrClosed + gzipPoolPut(zr) + } +} + +// close returns the gzip.Reader to the pool immediately or +// signals release to do so after Read completes. +func (gz *GzipReader) close() { + gz.mu.Lock() + defer gz.mu.Unlock() + if gz.zerr == nil && gz.zr != nil { + gzipPoolPut(gz.zr) + gz.zr = nil + } + gz.zerr = fs.ErrClosed +} + +func (gz *GzipReader) Read(p []byte) (n int, err error) { + zr, err := gz.acquire() + if err != nil { + return 0, err + } + defer gz.release(zr) + + return zr.Read(p) +} + +func (gz *GzipReader) Close() error { + gz.close() + + return gz.Body.Close() +} diff --git a/vendor/modules.txt b/vendor/modules.txt index dc910389f09d..f23979c86012 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -397,7 +397,7 @@ golang.org/x/mod/internal/lazyregexp golang.org/x/mod/modfile golang.org/x/mod/module golang.org/x/mod/semver -# golang.org/x/net v0.59.0 +# golang.org/x/net v0.60.0 ## explicit; go 1.26.0 golang.org/x/net/http/httpguts golang.org/x/net/http2