Skip to content

Commit 68620ef

Browse files
committed
TUN-10557: Bump quic-go v0.59.1
Bumps quic-go to v0.59.1 (chungthuang fork rebased from upstream v0.45 onto v0.59.1). Upstream removed the `logging` package and replaced its callback-based ConnectionTracer with the structured `qlog`/`qlogwriter` event API, which required migrating cloudflared's QUIC metrics collection. Migrations: - quic/tracing.go: connTracer no longer fills a logging.ConnectionTracer callback struct. It implements qlogwriter.Trace + qlogwriter.Recorder and dispatches qlog events (PacketSent, PacketReceived, MetricsUpdated, ...) to the collector through RecordEvent. NewClientTracer now returns a function compatible with quic.Config.Tracer. - quic/metrics.go: collector methods take qlog types (qlog.Frame, qlog.PacketType, qlog.MetricsUpdated, ...) and plain int64 in place of the removed logging.ByteCount/Frame/RTTStats/TransportParameters. - quic/conversion.go: PacketType, PacketDropReason and PacketLossReason are strings upstream rather than numeric iotas, so the converters become pass-through allowlists. CongestionState is also a string; congestionStateToFloat maps it back to the numeric gauge values cloudflared exports. - quic.Connection/quic.Stream became *quic.Conn/*quic.Stream; updated ConnWithCloser, SafeStreamCloser and the connection package accordingly. Tests and generated mocks (mocks/mock_quic_connection.go) were adapted to the new pointer-based API. Closes TUN-10557
1 parent 4d95ab7 commit 68620ef

367 files changed

Lines changed: 8721 additions & 76558 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

carrier/websocket_test.go

Lines changed: 27 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -2,10 +2,11 @@ package carrier
22

33
import (
44
"context"
5+
"crypto/rand"
56
"crypto/tls"
67
"crypto/x509"
78
"fmt"
8-
"math/rand"
9+
"math/big"
910
"testing"
1011
"time"
1112

@@ -23,28 +24,19 @@ import (
2324
func websocketClientTLSConfig(t *testing.T) *tls.Config {
2425
certPool := x509.NewCertPool()
2526
helloCert, err := tlsconfig.GetHelloCertificateX509()
26-
assert.NoError(t, err)
27+
require.NoError(t, err)
2728
certPool.AddCert(helloCert)
2829
assert.NotNil(t, certPool)
2930
return &tls.Config{RootCAs: certPool}
3031
}
3132

32-
func TestWebsocketHeaders(t *testing.T) {
33-
req := testRequest(t, "http://example.com", nil)
34-
wsHeaders := websocketHeaders(req)
35-
for _, header := range stripWebsocketHeaders {
36-
assert.Empty(t, wsHeaders[header])
37-
}
38-
assert.Equal(t, "curl/7.59.0", wsHeaders.Get("User-Agent"))
39-
}
40-
4133
func TestServe(t *testing.T) {
4234
log := zerolog.Nop()
4335
shutdownC := make(chan struct{})
4436
errC := make(chan error)
4537
listener, err := hello.CreateTLSListener("localhost:1111")
46-
assert.NoError(t, err)
47-
defer listener.Close()
38+
require.NoError(t, err)
39+
defer func() { _ = listener.Close() }()
4840

4941
go func() {
5042
errC <- hello.StartHelloWorldServer(&log, listener, shutdownC)
@@ -56,19 +48,25 @@ func TestServe(t *testing.T) {
5648
assert.NotNil(t, tlsConfig)
5749
d := gws.Dialer{TLSClientConfig: tlsConfig}
5850
conn, resp, err := clientConnect(req, &d)
59-
assert.NoError(t, err)
51+
require.NoError(t, err)
52+
defer func() { _ = resp.Body.Close() }()
6053
assert.Equal(t, "websocket", resp.Header.Get("Upgrade"))
6154

62-
for i := 0; i < 1000; i++ {
63-
messageSize := rand.Int()%2048 + 1
64-
clientMessage := make([]byte, messageSize)
65-
// rand.Read always returns len(clientMessage) and a nil error
66-
rand.Read(clientMessage)
55+
for range 1000 {
56+
messageSize, err := rand.Int(rand.Reader, big.NewInt(2048))
57+
require.NoError(t, err)
58+
clientMessage := make([]byte, messageSize.Int64()+1)
59+
for i := range clientMessage {
60+
n, err := rand.Int(rand.Reader, big.NewInt(256))
61+
n8 := uint8(n.Uint64()) //nolint:gosec // test-only
62+
require.NoError(t, err)
63+
clientMessage[i] = n8
64+
}
6765
err = conn.WriteMessage(websocket.BinaryFrame, clientMessage)
68-
assert.NoError(t, err)
66+
require.NoError(t, err)
6967

7068
messageType, message, err := conn.ReadMessage()
71-
assert.NoError(t, err)
69+
require.NoError(t, err)
7270
assert.Equal(t, websocket.BinaryFrame, messageType)
7371
assert.Equal(t, clientMessage, message)
7472
}
@@ -97,27 +95,30 @@ func TestWebsocketWrapper(t *testing.T) {
9795
req := testRequest(t, testAddr, nil)
9896
conn, resp, err := clientConnect(req, &d)
9997
require.NoError(t, err)
98+
defer func() { _ = resp.Body.Close() }()
10099
assert.Equal(t, "websocket", resp.Header.Get("Upgrade"))
101100

102101
// Websocket now connected to test server so lets check our wrapper
103102
wrapper := cfwebsocket.GorillaConn{Conn: conn}
104103
buf := make([]byte, 100)
105-
wrapper.Write([]byte("abc"))
104+
_, err = wrapper.Write([]byte("abc"))
105+
require.NoError(t, err)
106106
n, err := wrapper.Read(buf)
107107
require.NoError(t, err)
108-
require.Equal(t, n, 3)
108+
require.Equal(t, 3, n)
109109
require.Equal(t, "abc", string(buf[:n]))
110110

111111
// Test partial read, read 1 of 3 bytes in one read and the other 2 in another read
112-
wrapper.Write([]byte("abc"))
112+
_, err = wrapper.Write([]byte("abc"))
113+
require.NoError(t, err)
113114
buf = buf[:1]
114115
n, err = wrapper.Read(buf)
115116
require.NoError(t, err)
116-
require.Equal(t, n, 1)
117+
require.Equal(t, 1, n)
117118
require.Equal(t, "a", string(buf[:n]))
118119
buf = buf[:cap(buf)]
119120
n, err = wrapper.Read(buf)
120121
require.NoError(t, err)
121-
require.Equal(t, n, 2)
122+
require.Equal(t, 2, n)
122123
require.Equal(t, "bc", string(buf[:n]))
123124
}

connection/quic_connection.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -143,7 +143,7 @@ func (q *quicConnection) Serve(ctx context.Context) error {
143143
}
144144

145145
// serveControlStream will serve the RPC; blocking until the control plane is done.
146-
func (q *quicConnection) serveControlStream(ctx context.Context, controlStream quic.Stream) error {
146+
func (q *quicConnection) serveControlStream(ctx context.Context, controlStream *quic.Stream) error {
147147
return q.controlStreamHandler.ServeControlStream(ctx, controlStream, q.connOptions.ConnectionOptions(), q.orchestrator)
148148
}
149149

@@ -166,7 +166,7 @@ func (q *quicConnection) acceptStream(ctx context.Context) error {
166166
}
167167
}
168168

169-
func (q *quicConnection) runStream(quicStream quic.Stream) {
169+
func (q *quicConnection) runStream(quicStream *quic.Stream) {
170170
ctx := quicStream.Context()
171171
stream := cfdquic.NewSafeStreamCloser(quicStream, q.streamWriteTimeout, q.logger)
172172
defer func() { _ = stream.Close() }()

connection/quic_connection_test.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -530,7 +530,7 @@ func TestServeUDPSession(t *testing.T) {
530530
ctx, cancel := context.WithCancel(t.Context())
531531

532532
// Establish QUIC connection with edge
533-
edgeQUICSessionChan := make(chan quic.Connection)
533+
edgeQUICSessionChan := make(chan *quic.Conn)
534534
go func() {
535535
earlyListener, err := quic.Listen(udpListener, testTLSServerConfig, testQUICConfig)
536536
assert.NoError(t, err)
@@ -779,7 +779,7 @@ func TestDialQuicWithSkipPortReuse(t *testing.T) {
779779
<-serverDone
780780
}
781781

782-
func serveSession(ctx context.Context, datagramConn *datagramV2Connection, edgeQUICSession quic.Connection, closeType closeReason, expectedReason string, t *testing.T) {
782+
func serveSession(ctx context.Context, datagramConn *datagramV2Connection, edgeQUICSession cfdquic.QUICConnection, closeType closeReason, expectedReason string, t *testing.T) {
783783
payload := []byte(t.Name())
784784
sessionID := uuid.New()
785785
cfdConn, originConn := net.Pipe()
@@ -843,7 +843,7 @@ const (
843843
closedByTimeout
844844
)
845845

846-
func runRPCServer(ctx context.Context, session quic.Connection, sessionRPCServer pogs.SessionManager, configRPCServer pogs.ConfigurationManager, t *testing.T) {
846+
func runRPCServer(ctx context.Context, session cfdquic.QUICConnection, sessionRPCServer pogs.SessionManager, configRPCServer pogs.ConfigurationManager, t *testing.T) {
847847
stream, err := session.AcceptStream(ctx)
848848
require.NoError(t, err)
849849

connection/quic_datagram_v2.go

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ import (
88
"time"
99

1010
"github.com/google/uuid"
11-
pkgerrors "github.com/pkg/errors"
11+
"github.com/pkg/errors"
1212
"github.com/rs/zerolog"
1313
"go.opentelemetry.io/otel/attribute"
1414
"go.opentelemetry.io/otel/trace"
@@ -22,7 +22,7 @@ import (
2222
"github.com/cloudflare/cloudflared/packet"
2323
cfdquic "github.com/cloudflare/cloudflared/quic"
2424
"github.com/cloudflare/cloudflared/tracing"
25-
tunnelpogs "github.com/cloudflare/cloudflared/tunnelrpc/pogs"
25+
"github.com/cloudflare/cloudflared/tunnelrpc/pogs"
2626
rpcquic "github.com/cloudflare/cloudflared/tunnelrpc/quic"
2727
)
2828

@@ -31,14 +31,14 @@ const (
3131
demuxChanCapacity = 16
3232
)
3333

34-
var errInvalidDestinationIP = pkgerrors.New("unable to parse destination IP")
34+
var errInvalidDestinationIP = errors.New("unable to parse destination IP")
3535

3636
// DatagramSessionHandler is a service that can serve datagrams for a connection and handle sessions from incoming
3737
// connection streams.
3838
type DatagramSessionHandler interface {
3939
Serve(context.Context) error
4040

41-
tunnelpogs.SessionManager
41+
pogs.SessionManager
4242
}
4343

4444
type datagramV2Connection struct {
@@ -111,7 +111,7 @@ func (d *datagramV2Connection) Serve(ctx context.Context) error {
111111
}
112112

113113
// RegisterUdpSession is the RPC method invoked by edge to register and run a session
114-
func (q *datagramV2Connection) RegisterUdpSession(ctx context.Context, sessionID uuid.UUID, dstIP net.IP, dstPort uint16, closeAfterIdleHint time.Duration, traceContext string) (*tunnelpogs.RegisterUdpSessionResponse, error) {
114+
func (q *datagramV2Connection) RegisterUdpSession(ctx context.Context, sessionID uuid.UUID, dstIP net.IP, dstPort uint16, closeAfterIdleHint time.Duration, traceContext string) (*pogs.RegisterUdpSessionResponse, error) {
115115
traceCtx := tracing.NewTracedContext(ctx, traceContext, q.logger)
116116
ctx, registerSpan := traceCtx.Tracer().Start(traceCtx, "register-session", trace.WithAttributes(
117117
attribute.String("session-id", sessionID.String()),
@@ -123,7 +123,7 @@ func (q *datagramV2Connection) RegisterUdpSession(ctx context.Context, sessionID
123123
if err := q.flowLimiter.Acquire(management.UDP.String()); err != nil {
124124
log.Warn().Msgf("Too many concurrent sessions being handled, rejecting udp proxy to %s:%d", dstIP, dstPort)
125125

126-
err := pkgerrors.Wrap(err, "failed to start udp session due to rate limiting")
126+
err := errors.Wrap(err, "failed to start udp session due to rate limiting")
127127
tracing.EndWithErrorStatus(registerSpan, err)
128128
return nil, err
129129
}
@@ -180,7 +180,7 @@ func (q *datagramV2Connection) RegisterUdpSession(ctx context.Context, sessionID
180180
Msgf("Registered session")
181181
tracing.End(registerSpan)
182182

183-
resp := tunnelpogs.RegisterUdpSessionResponse{
183+
resp := pogs.RegisterUdpSessionResponse{
184184
Spans: traceCtx.GetProtoSpans(),
185185
}
186186

connection/quic_datagram_v2_test.go

Lines changed: 2 additions & 62 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,11 @@
11
package connection
22

33
import (
4-
"context"
54
"net"
65
"testing"
76
"time"
87

98
"github.com/google/uuid"
10-
"github.com/quic-go/quic-go"
119
"github.com/rs/zerolog"
1210
"github.com/stretchr/testify/require"
1311
"go.uber.org/mock/gomock"
@@ -16,73 +14,15 @@ import (
1614
"github.com/cloudflare/cloudflared/mocks"
1715
)
1816

19-
type mockQuicConnection struct{}
20-
21-
func (m *mockQuicConnection) AcceptStream(_ context.Context) (quic.Stream, error) {
22-
return nil, nil
23-
}
24-
25-
func (m *mockQuicConnection) AcceptUniStream(_ context.Context) (quic.ReceiveStream, error) {
26-
return nil, nil
27-
}
28-
29-
func (m *mockQuicConnection) OpenStream() (quic.Stream, error) {
30-
return nil, nil
31-
}
32-
33-
func (m *mockQuicConnection) OpenStreamSync(_ context.Context) (quic.Stream, error) {
34-
return nil, nil
35-
}
36-
37-
func (m *mockQuicConnection) OpenUniStream() (quic.SendStream, error) {
38-
return nil, nil
39-
}
40-
41-
func (m *mockQuicConnection) OpenUniStreamSync(_ context.Context) (quic.SendStream, error) {
42-
return nil, nil
43-
}
44-
45-
func (m *mockQuicConnection) LocalAddr() net.Addr {
46-
return nil
47-
}
48-
49-
func (m *mockQuicConnection) RemoteAddr() net.Addr {
50-
return nil
51-
}
52-
53-
func (m *mockQuicConnection) CloseWithError(_ quic.ApplicationErrorCode, s string) error {
54-
return nil
55-
}
56-
57-
func (m *mockQuicConnection) Context() context.Context {
58-
return nil
59-
}
60-
61-
func (m *mockQuicConnection) ConnectionState() quic.ConnectionState {
62-
panic("not meant to be called")
63-
}
64-
65-
func (m *mockQuicConnection) SendDatagram(_ []byte) error {
66-
return nil
67-
}
68-
69-
func (m *mockQuicConnection) ReceiveDatagram(_ context.Context) ([]byte, error) {
70-
return nil, nil
71-
}
72-
73-
func (m *mockQuicConnection) AddPath(*quic.Transport) (*quic.Path, error) {
74-
return nil, nil
75-
}
76-
7717
func TestRateLimitOnNewDatagramV2UDPSession(t *testing.T) {
7818
log := zerolog.Nop()
79-
conn := &mockQuicConnection{}
8019
ctrl := gomock.NewController(t)
8120
flowLimiterMock := mocks.NewMockLimiter(ctrl)
21+
connMock := mocks.NewMockQUICConnection(ctrl)
8222

8323
datagramConn := NewDatagramV2Connection(
8424
t.Context(),
85-
conn,
25+
connMock,
8626
nil,
8727
nil,
8828
0,

connection/quic_datagram_v3.go

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,7 @@ import (
1212
"github.com/cloudflare/cloudflared/ingress"
1313
"github.com/cloudflare/cloudflared/management"
1414
cfdquic "github.com/cloudflare/cloudflared/quic"
15-
cfdquicv3 "github.com/cloudflare/cloudflared/quic/v3"
15+
v3 "github.com/cloudflare/cloudflared/quic/v3"
1616
"github.com/cloudflare/cloudflared/tunnelrpc/pogs"
1717
)
1818

@@ -25,25 +25,25 @@ type datagramV3Connection struct {
2525
conn cfdquic.QUICConnection
2626
index uint8
2727
// datagramMuxer mux/demux datagrams from quic connection
28-
datagramMuxer cfdquicv3.DatagramConn
29-
metrics cfdquicv3.Metrics
28+
datagramMuxer v3.DatagramConn
29+
metrics v3.Metrics
3030
logger *zerolog.Logger
3131
}
3232

3333
func NewDatagramV3Connection(ctx context.Context,
3434
conn cfdquic.QUICConnection,
35-
sessionManager cfdquicv3.SessionManager,
35+
sessionManager v3.SessionManager,
3636
icmpRouter ingress.ICMPRouter,
3737
index uint8,
38-
metrics cfdquicv3.Metrics,
38+
metrics v3.Metrics,
3939
logger *zerolog.Logger,
4040
) DatagramSessionHandler {
4141
log := logger.
4242
With().
4343
Int(management.EventTypeKey, int(management.UDP)).
4444
Uint8(LogFieldConnIndex, index).
4545
Logger()
46-
datagramMuxer := cfdquicv3.NewDatagramConn(conn, sessionManager, icmpRouter, index, metrics, &log)
46+
datagramMuxer := v3.NewDatagramConn(conn, sessionManager, icmpRouter, index, metrics, &log)
4747

4848
return &datagramV3Connection{
4949
conn,

go.mod

Lines changed: 3 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ require (
2323
github.com/pkg/errors v0.9.1
2424
github.com/prometheus/client_golang v1.22.0
2525
github.com/prometheus/client_model v0.6.2
26-
github.com/quic-go/quic-go v0.52.0
26+
github.com/quic-go/quic-go v0.59.1
2727
github.com/rs/zerolog v1.20.0
2828
github.com/shirou/gopsutil/v4 v4.26.3
2929
github.com/stretchr/testify v1.11.1
@@ -35,7 +35,7 @@ require (
3535
go.opentelemetry.io/otel/trace v1.43.0
3636
go.opentelemetry.io/proto/otlp v1.10.0
3737
go.uber.org/automaxprocs v1.6.0
38-
go.uber.org/mock v0.5.1
38+
go.uber.org/mock v0.5.2
3939
golang.org/x/crypto v0.52.0
4040
golang.org/x/net v0.55.0
4141
golang.org/x/sync v0.20.0
@@ -65,10 +65,8 @@ require (
6565
github.com/go-logr/stdr v1.2.2 // indirect
6666
github.com/go-ole/go-ole v1.2.6 // indirect
6767
github.com/go-playground/validator/v10 v10.15.1 // indirect
68-
github.com/go-task/slim-sprig/v3 v3.0.0 // indirect
6968
github.com/gobwas/httphead v0.1.0 // indirect
7069
github.com/gobwas/pool v0.2.1 // indirect
71-
github.com/google/pprof v0.0.0-20250418163039-24c5476c6587 // indirect
7270
github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 // indirect
7371
github.com/klauspost/compress v1.18.0 // indirect
7472
github.com/klauspost/cpuid/v2 v2.2.5 // indirect
@@ -77,7 +75,6 @@ require (
7775
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
7876
github.com/modern-go/reflect2 v1.0.2 // indirect
7977
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
80-
github.com/onsi/ginkgo/v2 v2.23.4 // indirect
8178
github.com/pelletier/go-toml/v2 v2.0.9 // indirect
8279
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect
8380
github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect
@@ -108,5 +105,4 @@ replace github.com/prometheus/golang_client => github.com/prometheus/golang_clie
108105

109106
replace gopkg.in/yaml.v3 => gopkg.in/yaml.v3 v3.0.1
110107

111-
// This fork is based on quic-go v0.45
112-
replace github.com/quic-go/quic-go => github.com/chungthuang/quic-go v0.45.1-0.20250428085412-43229ad201fd
108+
replace github.com/quic-go/quic-go => github.com/chungthuang/quic-go v0.45.1-0.20260529212404-a9fddf436fc4 // This fork is based on quic-go v0.59.1

0 commit comments

Comments
 (0)