diff --git a/go.mod b/go.mod index e1736c31..166496d9 100644 --- a/go.mod +++ b/go.mod @@ -4,7 +4,7 @@ go 1.26.0 require ( github.com/IBM/sarama v1.50.3 - github.com/pion/dtls/v2 v2.2.12 + github.com/pion/dtls/v3 v3.1.8 github.com/spf13/cobra v1.10.2 github.com/spf13/pflag v1.0.10 github.com/stretchr/testify v1.12.1 @@ -36,9 +36,8 @@ require ( github.com/modern-go/reflect2 v1.0.3-0.20250322232337-35a7c28c31ee // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/pierrec/lz4/v4 v4.1.27 // indirect - github.com/pion/logging v0.2.2 // indirect - github.com/pion/transport/v2 v2.2.10 // indirect - github.com/pion/transport/v3 v3.0.7 // indirect + github.com/pion/logging v0.2.4 // indirect + github.com/pion/transport/v4 v4.0.2 // indirect github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect github.com/prometheus/client_golang v1.23.2 // indirect github.com/prometheus/client_model v0.6.2 // indirect diff --git a/go.sum b/go.sum index de5a42e2..7cbbcd11 100644 --- a/go.sum +++ b/go.sum @@ -61,15 +61,12 @@ github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= github.com/pierrec/lz4/v4 v4.1.27 h1:+PhzhWDrjRj89TH2sw43nE3+4+W8lSxIuQadEHZyjUk= github.com/pierrec/lz4/v4 v4.1.27/go.mod h1:EoQMVJgeeEOMsCqCzqFm2O0cJvljX2nGZjcRIPL34O4= -github.com/pion/dtls/v2 v2.2.12 h1:KP7H5/c1EiVAAKUmXyCzPiQe5+bCJrpOeKg/L05dunk= -github.com/pion/dtls/v2 v2.2.12/go.mod h1:d9SYc9fch0CqK90mRk1dC7AkzzpwJj6u2GU3u+9pqFE= -github.com/pion/logging v0.2.2 h1:M9+AIj/+pxNsDfAT64+MAVgJO0rsyLnoJKCqf//DoeY= -github.com/pion/logging v0.2.2/go.mod h1:k0/tDVsRCX2Mb2ZEmTqNa7CWsQPc+YYCB7Q+5pahoms= -github.com/pion/transport/v2 v2.2.4/go.mod h1:q2U/tf9FEfnSBGSW6w5Qp5PFWRLRj3NjLhCCgpRK4p0= -github.com/pion/transport/v2 v2.2.10 h1:ucLBLE8nuxiHfvkFKnkDQRYWYfp8ejf4YBOPfaQpw6Q= -github.com/pion/transport/v2 v2.2.10/go.mod h1:sq1kSLWs+cHW9E+2fJP95QudkzbK7wscs8yYgQToO5E= -github.com/pion/transport/v3 v3.0.7 h1:iRbMH05BzSNwhILHoBoAPxoB9xQgOaJk+591KC9P1o0= -github.com/pion/transport/v3 v3.0.7/go.mod h1:YleKiTZ4vqNxVwh77Z0zytYi7rXHl7j6uPLGhhz9rwo= +github.com/pion/dtls/v3 v3.1.8 h1:aLcgjZqzrYn5AbjSds4LvK2WI5VzJc1PencExyDjYis= +github.com/pion/dtls/v3 v3.1.8/go.mod h1:gz1K4jg6c+fq86oQMH4pilpCEOEPwmEr2jY+VcF/mkU= +github.com/pion/logging v0.2.4 h1:tTew+7cmQ+Mc1pTBLKH2puKsOvhm32dROumOZ655zB8= +github.com/pion/logging v0.2.4/go.mod h1:DffhXTKYdNZU+KtJ5pyQDjvOAh/GsNSyv1lbkFbe3so= +github.com/pion/transport/v4 v4.0.2 h1:ifYlPqNwsy6aKQ9y8yzxXlHae5431ZrH2avkD/Rn6Tk= +github.com/pion/transport/v4 v4.0.2/go.mod h1:06hFI+jCFcok2X2MekVufNZ/uzNZXivGBPfviSVcjgM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U= github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= @@ -99,10 +96,8 @@ github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81P github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= -github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo= github.com/stretchr/testify v1.12.1 h1:EuwCh5fleGS7H32xRwO3wRGT7DxrDhLAT6FF8MpWDWE= github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUmoNlxoYthg= -github.com/wlynxg/anet v0.0.3/go.mod h1:eay5PRQr7fIVAMbTbchTnO9gG65Hg/uYGdc7mguHxoA= github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY= @@ -126,26 +121,19 @@ go.yaml.in/yaml/v3 v3.0.5/go.mod h1:HVTZu1O7/Vkt2N+BFy8Zza+lnLsABggaTM2ZpNIGuKg= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= golang.org/x/crypto v0.6.0/go.mod h1:OFC/31mSvZgRz0V1QTNCzfAI1aIRzbiufJtkMIlEp58= -golang.org/x/crypto v0.12.0/go.mod h1:NF0Gs7EO5K4qLn+Ylc+fih8BSTeIjAP05siRnAh98yw= -golang.org/x/crypto v0.18.0/go.mod h1:R0j02AL6hcrfOiy9T4ZYp/rcWeMxM3L6QYxlOuEG1mg= golang.org/x/crypto v0.53.0 h1:QZ4Muo8THX6CizN2vPPd5fBGHyogrdK9fG4wLPFUsto= golang.org/x/crypto v0.53.0/go.mod h1:DNLU434OwVakk9PzuwV8w62mAJpRJL3vsgcfp4Qnsio= golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= -golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20200114155413-6afb5195e5aa/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c= golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs= golang.org/x/net v0.7.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs= -golang.org/x/net v0.10.0/go.mod h1:0qNGK6F8kojg2nk9dLZ2mShWaEBan6FAoqfSigmmuDg= -golang.org/x/net v0.14.0/go.mod h1:PpSgVXXLK0OxS0F31C1/tv6XNguvCrnXIDrFMspZIUI= -golang.org/x/net v0.20.0/go.mod h1:z8BVo6PvndSri0LbOE3hAn0apkU+1YvI6E70E9jsnvY= golang.org/x/net v0.56.0 h1:Rw8j/hFzGvJUZwNBXnAtf5sVDVt+65SK2C7IxCxZt5o= golang.org/x/net v0.56.0/go.mod h1:D3Ku6r+V6JROoZK144D2XfMHFcMq/0zSfLelVTCFKec= golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM= golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= @@ -154,30 +142,20 @@ golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBc golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.11.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.16.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw= golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= golang.org/x/term v0.5.0/go.mod h1:jMB1sMXY+tzblOD4FWmEbocvup2/aLOaQEp7JmGp78k= -golang.org/x/term v0.8.0/go.mod h1:xPskH00ivmX89bAKVGSKKtLOWNx2+17Eiy94tnKShWo= -golang.org/x/term v0.11.0/go.mod h1:zC9APTIj3jG3FdV/Ons+XE1riIZXG4aZ4GTHiPZJPIU= -golang.org/x/term v0.16.0/go.mod h1:yn7UURbUtPyrVJPGPq404EukNFxcm/foM+bV/bfcDsY= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= -golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8= -golang.org/x/text v0.12.0/go.mod h1:TvPlkZtksWOMsz7fbANvkp4WM8x/WCo/om8BMLbz+aE= -golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= golang.org/x/text v0.38.0 h1:sXmwo9DwP3OK9EZ7PqAdaooSGozfl/3a6/xJcbzPRhE= golang.org/x/text v0.38.0/go.mod h1:YXZt3QhHUKYT53r2lLKFIVi6Ao1jdzrTR/KQ09qyxF4= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= -golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= google.golang.org/protobuf v1.36.12 h1:pJOKDDOyeXErUroCihFAd5LQuwXBSpVnKGrj5o/fwxc= google.golang.org/protobuf v1.36.12/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= diff --git a/pkg/collector/process.go b/pkg/collector/process.go index 35f9b228..ac60daeb 100644 --- a/pkg/collector/process.go +++ b/pkg/collector/process.go @@ -116,7 +116,7 @@ type CollectorInput struct { // List of supported cipher suites. // From https://pkg.go.dev/crypto/tls#pkg-constants // The order of the list is ignored.Note that TLS 1.3 ciphersuites are not configurable. - // For DTLS, cipher suites are from https://pkg.go.dev/github.com/pion/dtls/v2@v2.2.12/internal/ciphersuite#ID. + // For DTLS, cipher suites are from https://pkg.go.dev/github.com/pion/dtls/v3@v3.1.8/internal/ciphersuite#ID. TLSCipherSuites []uint16 // Min TLS version. // From https://pkg.go.dev/crypto/tls#pkg-constants diff --git a/pkg/collector/process_test.go b/pkg/collector/process_test.go index 5eb8a6f4..6b8b3346 100644 --- a/pkg/collector/process_test.go +++ b/pkg/collector/process_test.go @@ -29,7 +29,8 @@ import ( "testing" "time" - "github.com/pion/dtls/v2" + "github.com/pion/dtls/v3" + dtlsnet "github.com/pion/dtls/v3/pkg/net" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "k8s.io/apimachinery/pkg/util/wait" @@ -757,6 +758,41 @@ func TestTLSCollectingProcess(t *testing.T) { assert.Error(t, err) } +// dialDTLS connects to the DTLS collector listening at addr, trusting rootCertPEM, and +// completes the handshake. Failing to create the connection fails the test; a failed +// handshake does not, it is returned as an error, so that callers can assert on it. +// The connection is closed via t.Cleanup, which runs after the test function's own +// defers: do not rely on it being closed before, say, a deferred cp.Stop(). +func dialDTLS(t *testing.T, addr *net.UDPAddr, rootCertPEM []byte) (*dtls.Conn, error) { + t.Helper() + roots := x509.NewCertPool() + require.True(t, roots.AppendCertsFromPEM(rootCertPEM), "Failed to parse root certificate") + // We dial the socket ourselves and hand it to pion as a net.PacketConn, instead of + // using dtls.Dial which would leave the socket unconnected. This is what the exporter + // does: see InitExportingProcess. It matters here because the collector keys sessions + // by the source address of the datagrams it receives, and only a connected socket has + // a local address that is guaranteed to match that source address. + udpConn, err := net.DialUDP(addr.Network(), nil, addr) + require.NoError(t, err) + conn, err := dtls.ClientWithOptions( + dtlsnet.PacketConnFromConn(udpConn), addr, + dtls.WithRootCAs(roots), + dtls.WithExtendedMasterSecret(dtls.RequireExtendedMasterSecret), + ) + if err != nil { + udpConn.Close() + } + require.NoError(t, err) + t.Cleanup(func() { conn.Close() }) + // As of pion/dtls v3, the handshake is not performed when the connection is created, + // but lazily on the first Read / Write. We trigger it explicitly so that handshake + // errors are reported here. We bound it, so that a collector which stalls mid-handshake + // fails this test instead of hanging until the package-level test timeout. + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + return conn, conn.HandshakeContext(ctx) +} + func TestDTLSCollectingProcess(t *testing.T) { input := getCollectorInput(udpTransport, true, false) cp, err := InitCollectingProcess(input) @@ -767,16 +803,40 @@ func TestDTLSCollectingProcess(t *testing.T) { defer cp.Stop() collectorAddr, _ := net.ResolveUDPAddr("udp", cp.GetAddress().String()) - roots := x509.NewCertPool() - ok := roots.AppendCertsFromPEM(testcerts.FakeCert2) - if !ok { - t.Error("Failed to parse root certificate") - } - config := &dtls.Config{RootCAs: roots, - ExtendedMasterSecret: dtls.RequireExtendedMasterSecret} - conn, err := dtls.Dial("udp", collectorAddr, config) + conn, err := dialDTLS(t, collectorAddr, testcerts.FakeCert2) + require.NoError(t, err) + _, err = conn.Write(validTemplatePacket) + require.NoError(t, err) + <-cp.GetMsgChan() + template, _ := cp.getTemplateIEs(localConnSessionID(conn), 1, 256) + assert.NotNil(t, template, "DTLS Collecting Process should receive and store the received template") +} + +// TestDTLSCollectingProcessFailedHandshake checks that a client which fails the DTLS +// handshake does not prevent subsequent clients from connecting: otherwise any peer could +// take the collector down by starting a handshake it cannot complete. +func TestDTLSCollectingProcessFailedHandshake(t *testing.T) { + input := getCollectorInput(udpTransport, true, false) + cp, err := InitCollectingProcess(input) + require.NoError(t, err) + go cp.Start() + // wait until collector is ready + waitForCollectorReady(t, cp) + defer cp.Stop() + + collectorAddr, _ := net.ResolveUDPAddr("udp", cp.GetAddress().String()) + // This client does not trust the collector's certificate, so it aborts the handshake. + _, _, unrelatedCACertPEM, err := testcerts.GenerateCACert() + require.NoError(t, err) + _, err = dialDTLS(t, collectorAddr, unrelatedCACertPEM) + require.Error(t, err, "Handshake should fail when the collector certificate is not trusted") + // The collector must have rejected us, not just gone silent: a timeout here would mean + // the accept path stalled, which is what the second dial below is meant to detect. + require.NotErrorIs(t, err, context.DeadlineExceeded) + + // The collector must still be accepting connections. + conn, err := dialDTLS(t, collectorAddr, testcerts.FakeCert2) require.NoError(t, err) - defer conn.Close() _, err = conn.Write(validTemplatePacket) require.NoError(t, err) <-cp.GetMsgChan() diff --git a/pkg/collector/udp.go b/pkg/collector/udp.go index 55684dc4..09b6472a 100644 --- a/pkg/collector/udp.go +++ b/pkg/collector/udp.go @@ -16,15 +16,29 @@ package collector import ( "bytes" + "context" "crypto/tls" "crypto/x509" "fmt" "net" + "time" - "github.com/pion/dtls/v2" + "github.com/pion/dtls/v3" "k8s.io/klog/v2" ) +// dtlsHandshakeTimeout bounds the duration of the DTLS handshake with a client. pion/dtls +// v2 applied a default 30s timeout to the handshake performed by Accept, while in v3 +// HandshakeContext honors only the context it is given. Without a bound, a peer which +// starts a handshake and then stops responding would block this goroutine for the +// lifetime of the process. +const dtlsHandshakeTimeout = 30 * time.Second + +// dtlsHandshaker is the part of *dtls.Conn we need from the net.Conn returned by Accept. +type dtlsHandshaker interface { + HandshakeContext(ctx context.Context) error +} + func (cp *CollectingProcess) startUDPServer() { var listener net.Listener var err error @@ -35,12 +49,12 @@ func (cp *CollectingProcess) startUDPServer() { return } if cp.isEncrypted { // use DTLS - config, err := cp.createServerDTLSConfig() + options, err := cp.createServerDTLSOptions() if err != nil { klog.Error(err) return } - listener, err = dtls.Listen("udp", address, config) + listener, err = dtls.ListenWithOptions("udp", address, options...) if err != nil { klog.Error(err) return @@ -48,10 +62,41 @@ func (cp *CollectingProcess) startUDPServer() { defer listener.Close() cp.updateAddress(listener.Addr()) klog.Infof("Start dtls collecting process on %s", cp.netAddress) - conn, err = listener.Accept() - if err != nil { - klog.Error(err) - return + // As of pion/dtls v3, Accept no longer performs the handshake: it happens lazily + // on the first Read / Write. We trigger it explicitly so that handshake failures + // (e.g. client certificate validation errors) are reported clearly here, instead + // of surfacing as an opaque read error below. A handshake failure invalidates + // that connection only: we go back to accepting, so that a peer which fails to + // authenticate cannot permanently stop the collector from serving a legitimate + // exporter. Unlike the plaintext UDP path below, which demultiplexes datagrams + // from any number of exporters by source address, this loop serves a single + // exporter: it stops accepting once one client has completed the handshake. + // Because the handshake runs inline, a peer which starts one and then stops + // responding still blocks the loop for dtlsHandshakeTimeout. + for { + conn, err = listener.Accept() + if err != nil { + klog.Error(err) + return + } + handshaker, ok := conn.(dtlsHandshaker) + if !ok { + // Only reachable if pion stops returning *dtls.Conn from Accept. Drop the + // connection and keep serving: the listener is still usable. + klog.ErrorS(nil, "DTLS listener returned an unexpected connection type", "type", fmt.Sprintf("%T", conn)) + conn.Close() + continue + } + ctx, cancel := context.WithTimeout(context.Background(), dtlsHandshakeTimeout) + err = handshaker.HandshakeContext(ctx) + cancel() + if err == nil { + break + } + // Remotely triggerable, so this can be noisy: it is logged per failed + // handshake, like the per-connection errors on the TCP path. + klog.ErrorS(err, "Error during DTLS handshake with client, dropping connection", "client", conn.RemoteAddr()) + conn.Close() } defer conn.Close() cp.wg.Add(1) @@ -173,7 +218,7 @@ func (cp *CollectingProcess) createUDPClient(addr string) *transportSession { return session } -func (cp *CollectingProcess) createServerDTLSConfig() (*dtls.Config, error) { +func (cp *CollectingProcess) createServerDTLSOptions() ([]dtls.ServerOption, error) { if cp.tlsMinVersion != 0 && cp.tlsMinVersion != tls.VersionTLS12 { return nil, fmt.Errorf("DTLS 1.2 is the only supported version") } @@ -181,25 +226,30 @@ func (cp *CollectingProcess) createServerDTLSConfig() (*dtls.Config, error) { if err != nil { return nil, err } - // If tlsConfig.CipherSuites is nil, cipherSuites should also be nil! - var cipherSuites []dtls.CipherSuiteID - for _, cipherSuite := range cp.tlsCipherSuites { - cipherSuites = append(cipherSuites, dtls.CipherSuiteID(cipherSuite)) + options := []dtls.ServerOption{ + dtls.WithCertificates(cert), + dtls.WithExtendedMasterSecret(dtls.RequireExtendedMasterSecret), } - config := &dtls.Config{ - Certificates: []tls.Certificate{cert}, - ExtendedMasterSecret: dtls.RequireExtendedMasterSecret, - CipherSuites: cipherSuites, + // If cp.tlsCipherSuites is empty, we must not set the option at all, so that the pion + // defaults are used. + if len(cp.tlsCipherSuites) > 0 { + cipherSuites := make([]dtls.CipherSuiteID, 0, len(cp.tlsCipherSuites)) + for _, cipherSuite := range cp.tlsCipherSuites { + cipherSuites = append(cipherSuites, dtls.CipherSuiteID(cipherSuite)) + } + options = append(options, dtls.WithCipherSuites(cipherSuites...)) } if cp.caCert == nil { - return config, nil + return options, nil } clientCAs := x509.NewCertPool() ok := clientCAs.AppendCertsFromPEM(cp.caCert) if !ok { return nil, fmt.Errorf("failed to parse client CA certificate") } - config.ClientAuth = dtls.RequireAndVerifyClientCert - config.ClientCAs = clientCAs - return config, nil + options = append(options, + dtls.WithClientAuth(dtls.RequireAndVerifyClientCert), + dtls.WithClientCAs(clientCAs), + ) + return options, nil } diff --git a/pkg/exporter/process.go b/pkg/exporter/process.go index 38575806..6f4e55ec 100644 --- a/pkg/exporter/process.go +++ b/pkg/exporter/process.go @@ -16,6 +16,7 @@ package exporter import ( "bytes" + "context" "crypto/tls" "crypto/x509" "encoding/json" @@ -26,16 +27,23 @@ import ( "sync/atomic" "time" - "github.com/pion/dtls/v2" + "github.com/pion/dtls/v3" + dtlsnet "github.com/pion/dtls/v3/pkg/net" "k8s.io/klog/v2" "github.com/vmware/go-ipfix/pkg/entities" ) const startTemplateID uint16 = 255 + const defaultCheckConnInterval = 10 * time.Second const defaultJSONBufferLen = 5000 +// dtlsHandshakeTimeout bounds the duration of the DTLS handshake with the collector. +// pion/dtls v2 applied a default 30s timeout to the handshake performed by Dial, while in +// v3 HandshakeContext honors only the context it is given. +const dtlsHandshakeTimeout = 30 * time.Second + type templateValue struct { elements []*entities.InfoElement minDataRecLen uint16 @@ -79,7 +87,7 @@ type ExporterTLSClientConfig struct { // List of supported cipher suites. // From https://pkg.go.dev/crypto/tls#pkg-constants // The order of the list is ignored.Note that TLS 1.3 ciphersuites are not configurable. - // For DTLS, cipher suites are from https://pkg.go.dev/github.com/pion/dtls/v2@v2.2.12/internal/ciphersuite#ID. + // For DTLS, cipher suites are from https://pkg.go.dev/github.com/pion/dtls/v3@v3.1.8/internal/ciphersuite#ID. CipherSuites []uint16 // Min TLS version. // From https://pkg.go.dev/crypto/tls#pkg-constants @@ -191,25 +199,58 @@ func InitExportingProcess(input ExporterInput) (*ExportingProcess, error) { if !ok { return nil, fmt.Errorf("failed to parse root certificate") } - // If tlsConfig.CipherSuites is nil, cipherSuites should also be nil! - var cipherSuites []dtls.CipherSuiteID - for _, cipherSuite := range tlsConfig.CipherSuites { - cipherSuites = append(cipherSuites, dtls.CipherSuiteID(cipherSuite)) + options := []dtls.ClientOption{ + dtls.WithRootCAs(roots), + dtls.WithExtendedMasterSecret(dtls.RequireExtendedMasterSecret), + dtls.WithServerName(tlsConfig.ServerName), } - config := &dtls.Config{ - RootCAs: roots, - ExtendedMasterSecret: dtls.RequireExtendedMasterSecret, - ServerName: tlsConfig.ServerName, - CipherSuites: cipherSuites, + // If tlsConfig.CipherSuites is empty, we must not set the option at all, so + // that the pion defaults are used. + if len(tlsConfig.CipherSuites) > 0 { + cipherSuites := make([]dtls.CipherSuiteID, 0, len(tlsConfig.CipherSuites)) + for _, cipherSuite := range tlsConfig.CipherSuites { + cipherSuites = append(cipherSuites, dtls.CipherSuiteID(cipherSuite)) + } + options = append(options, dtls.WithCipherSuites(cipherSuites...)) } udpAddr, err := net.ResolveUDPAddr(input.CollectorProtocol, input.CollectorAddress) if err != nil { return nil, err } - conn, err = dtls.Dial(udpAddr.Network(), udpAddr, config) + // We do not use dtls.Dial, because as of pion/dtls v3 it creates the socket + // with net.ListenUDP instead of net.DialUDP: the socket is no longer + // connected, which pion needs in order to support net.PacketConn.WriteTo. We + // have no use for WriteTo, and an unconnected socket would change our + // behavior in 3 ways: the source address of exported packets would be chosen + // by a route lookup for every datagram, instead of being fixed when we dial, + // the kernel would no longer discard datagrams from sources other than the + // collector, and Write would no longer report ICMP errors such as + // ECONNREFUSED. It would also make DTLS behave differently from plaintext + // UDP, which goes through net.Dial and is therefore connected. We dial the + // socket ourselves and pass it to pion as a net.PacketConn: the + // PacketConnFromConn wrapper ignores the address given to WriteTo and simply + // writes to the connected socket. + udpConn, err := net.DialUDP(udpAddr.Network(), nil, udpAddr) if err != nil { + return nil, fmt.Errorf("cannot create the UDP connection to the Collector %q: %w", udpAddr.String(), err) + } + dtlsConn, err := dtls.ClientWithOptions(dtlsnet.PacketConnFromConn(udpConn), udpAddr, options...) + if err != nil { + udpConn.Close() return nil, fmt.Errorf("cannot create the DTLS connection to the Collector %q: %w", udpAddr.String(), err) } + // As of pion/dtls v3, Dial no longer performs the handshake: it happens lazily + // on the first Read / Write. We trigger it explicitly so that connection + // failures (e.g. certificate validation errors) are reported to the caller + // here, instead of when the first IPFIX message is sent. + handshakeCtx, cancelHandshake := context.WithTimeout(context.Background(), dtlsHandshakeTimeout) + err = dtlsConn.HandshakeContext(handshakeCtx) + cancelHandshake() + if err != nil { + dtlsConn.Close() + return nil, fmt.Errorf("error during DTLS handshake with the Collector %q: %w", udpAddr.String(), err) + } + conn = dtlsConn } } else { conn, err = net.Dial(input.CollectorProtocol, input.CollectorAddress) diff --git a/pkg/exporter/process_test.go b/pkg/exporter/process_test.go index 6bb8f5e1..a05bef11 100644 --- a/pkg/exporter/process_test.go +++ b/pkg/exporter/process_test.go @@ -25,7 +25,7 @@ import ( "testing" "time" - "github.com/pion/dtls/v2" + "github.com/pion/dtls/v3" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -636,11 +636,11 @@ func TestExportingProcessWithDTLS(t *testing.T) { t.Error(err) return } - config := &dtls.Config{ - Certificates: []tls.Certificate{cert}, - ExtendedMasterSecret: dtls.RequireExtendedMasterSecret, - } - listener, err := dtls.Listen("udp", address, config) + listener, err := dtls.ListenWithOptions( + "udp", address, + dtls.WithCertificates(cert), + dtls.WithExtendedMasterSecret(dtls.RequireExtendedMasterSecret), + ) if err != nil { t.Errorf("Cannot start dtls collecting process on %s: %v", listener.Addr().String(), err) return