Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 3 additions & 4 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
34 changes: 6 additions & 28 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down Expand Up @@ -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=
Expand All @@ -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=
Expand All @@ -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=
Expand Down
2 changes: 1 addition & 1 deletion pkg/collector/process.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
80 changes: 70 additions & 10 deletions pkg/collector/process_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
Expand All @@ -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()
Expand Down
90 changes: 70 additions & 20 deletions pkg/collector/udp.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -35,23 +49,54 @@ 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
}
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)
Expand Down Expand Up @@ -173,33 +218,38 @@ 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")
}
cert, err := tls.X509KeyPair(cp.serverCert, cp.serverKey)
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
}
Loading
Loading