From ce4737ed9c35669206b0f433157aa37fd51e7051 Mon Sep 17 00:00:00 2001 From: Purnesh Dixit Date: Mon, 3 Feb 2025 14:24:30 +0530 Subject: [PATCH 01/10] grpc based transport --- .../clients/grpctransport/grpc_transport.go | 142 ++++++++++++++++ .../grpctransport/grpc_transport_test.go | 157 ++++++++++++++++++ 2 files changed, 299 insertions(+) create mode 100644 xds/internal/clients/grpctransport/grpc_transport.go create mode 100644 xds/internal/clients/grpctransport/grpc_transport_test.go diff --git a/xds/internal/clients/grpctransport/grpc_transport.go b/xds/internal/clients/grpctransport/grpc_transport.go new file mode 100644 index 000000000000..49d82c5b9ce3 --- /dev/null +++ b/xds/internal/clients/grpctransport/grpc_transport.go @@ -0,0 +1,142 @@ +/* + * + * Copyright 2025 gRPC authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +// Package grpctransport provides an implementation of the +// clients.TransportBuilder interface using gRPC. +package grpctransport + +import ( + "context" + "fmt" + "time" + + "google.golang.org/grpc" + "google.golang.org/grpc/credentials" + "google.golang.org/grpc/keepalive" + "google.golang.org/grpc/xds/internal/clients" +) + +// ServerConfigExtension allows extending the clients.ServerConfig for the +// gRPC-based transport builder. Any implementation of this must implement +// the ServerConfig() method. +// +// This interface must be implemented by the grpc transport builder provided +// in the Extensions field of the clients.ServerConfig. +type ServerConfigExtension interface { + ServerConfig() *ServerConfig +} + +// ServerConfig holds the settings for connecting to an xDS management server +// using gRPC. +type ServerConfig struct { + // Credentials is the credential bundle containing the gRPC credentials for + // connecting to the xDS management server. + Credentials credentials.Bundle +} + +// ServerConfig returns the ServerConfig itself. This method is designed +// to satisfy [ServerConfigExtension] interface requirement. +func (s *ServerConfig) ServerConfig() *ServerConfig { + return s +} + +// Builder provides a way to build a gRPC-based transport to an xDS management +// server. +type Builder struct{} + +// Build creates a new gRPC-based transport to an xDS management server using +// the provided clients.ServerConfig. This involves creating a +// grpc.ClientConn to the server using the provided credentials and server URI. +// +// If any of ServerURI or Extensions of `sc` are not present, Build() will return +// an error. +func (b *Builder) Build(sc clients.ServerConfig) (clients.Transport, error) { + if sc.ServerURI == "" { + return nil, fmt.Errorf("xds: ServerConfig's ServerURI field cannot be empty") + } + if sc.Extensions == nil { + return nil, fmt.Errorf("xds: ServerConfig's Extensions field cannot be nil for gRPC transport") + } + gtsce, ok := sc.Extensions.(ServerConfigExtension) + if !ok { + return nil, fmt.Errorf("xds: ServerConfig's Extensions field cannot be anything other than grpctransport.ServerConfigExtension for gRPC transport") + } + gtsc := gtsce.ServerConfig() + if gtsc.Credentials == nil { + return nil, fmt.Errorf("xsd: ServerConfigExtensions's Credentials field cannot be nil for gRPC transport") + } + + // TODO: Incorporate reference count map for existing transports and + // deduplicate transports based on server URI and credentials so that + // transport channel to same server can be shared between xDS and LRS + // client. + + // Dial the xDS management server with the provided credentials, server URI, + // and a static keepalive configuration that is common across gRPC language + // implementations. + kpCfg := grpc.WithKeepaliveParams(keepalive.ClientParameters{ + Time: 5 * time.Minute, + Timeout: 20 * time.Second, + }) + cc, err := grpc.NewClient(sc.ServerURI, kpCfg, grpc.WithCredentialsBundle(gtsc.Credentials)) + if err != nil { + return nil, fmt.Errorf("error creating grpc client for server uri %s, %v", sc.ServerURI, err) + } + cc.Connect() + + return &grpcTransport{cc: cc}, nil +} + +type grpcTransport struct { + cc *grpc.ClientConn +} + +// NewStream creates a new gRPC stream to the xDS management server for the +// specified method. The returned Stream interface can be used to send and +// receive messages on the stream. +func (g *grpcTransport) NewStream(ctx context.Context, method string) (clients.Stream, error) { + s, err := g.cc.NewStream(ctx, &grpc.StreamDesc{StreamName: method, ClientStreams: true, ServerStreams: true}, method) + if err != nil { + return nil, err + } + return &stream{stream: s}, nil +} + +type stream struct { + stream grpc.ClientStream +} + +// Send sends a message to the xDS management server. +func (s *stream) Send(msg []byte) error { + return s.stream.SendMsg(msg) +} + +// Recv receives a message from the xDS management server. +func (s *stream) Recv() ([]byte, error) { + var typedRes []byte + err := s.stream.RecvMsg(&typedRes) + if err != nil { + return typedRes, err + } + return typedRes, nil +} + +// Close closes the gRPC stream to the xDS management server. +func (g *grpcTransport) Close() error { + return g.cc.Close() +} diff --git a/xds/internal/clients/grpctransport/grpc_transport_test.go b/xds/internal/clients/grpctransport/grpc_transport_test.go new file mode 100644 index 000000000000..fc1635f08b0d --- /dev/null +++ b/xds/internal/clients/grpctransport/grpc_transport_test.go @@ -0,0 +1,157 @@ +/* + * + * Copyright 2025 gRPC authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package grpctransport + +import ( + "context" + "testing" + "time" + + "google.golang.org/grpc/credentials/insecure" + "google.golang.org/grpc/internal/grpctest" + "google.golang.org/grpc/internal/testutils/xds/e2e" + "google.golang.org/grpc/xds/internal/clients" +) + +const ( + defaultTestWatchExpiryTimeout = 500 * time.Millisecond + defaultTestTimeout = 10 * time.Second +) + +type s struct { + grpctest.Tester +} + +func Test(t *testing.T) { + grpctest.RunSubTests(t, s{}) +} + +// TestBuild verifies that the grpctransport.Builder creates a new +// grpc.ClientConn every time Build() is called. +// +// It covers the following scenarios: +// - ServerURI is empty. +// - Extensions is nil. +// - Extensions is not ServerConfigExtension. +// - Credentials are nil. +// - Success cases. +func (s) TestBuild(t *testing.T) { + tests := []struct { + name string + serverCfg clients.ServerConfig + wantErr bool + }{ + { + name: "ServerURI_empty", + serverCfg: clients.ServerConfig{ + ServerURI: "", + Extensions: &ServerConfig{Credentials: insecure.NewBundle()}, + }, + wantErr: true, + }, + { + name: "Extensions_nil", + serverCfg: clients.ServerConfig{ServerURI: "server-address"}, + wantErr: true, + }, + { + name: "Extensions_not_ServerConfigExtension", + serverCfg: clients.ServerConfig{ + ServerURI: "server-address", + Extensions: 1, + }, + wantErr: true, + }, + { + name: "ServerConfigExtension_Credentials_nil", + serverCfg: clients.ServerConfig{ + ServerURI: "server-address", + Extensions: &ServerConfig{}, + }, + wantErr: true, + }, + { + name: "success", + serverCfg: clients.ServerConfig{ + ServerURI: "server-address", + Extensions: &ServerConfig{Credentials: insecure.NewBundle()}, + }, + wantErr: false, + }, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + b := &Builder{} + tr, err := b.Build(test.serverCfg) + if (err != nil) != test.wantErr { + t.Fatalf("Build() error = %v, wantErr %v", err, test.wantErr) + } + if tr != nil { + defer tr.Close() + } + if !test.wantErr && tr == nil { + t.Fatalf("got non-nil transport from Build(), want nil") + } + }) + } +} + +// TestNewStream verifies that grpcTransport.NewStream() successfully creates a +// new client stream for the xDS management server. +func (s) TestNewStream(t *testing.T) { + mgmtServer := e2e.StartManagementServer(t, e2e.ManagementServerOptions{}) + + tests := []struct { + name string + serverURI string + wantErr bool + }{ + { + name: "success", + serverURI: mgmtServer.Address, + wantErr: false, + }, + { + name: "error", + serverURI: "invalid-server-uri", + wantErr: true, + }, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + serverCfg := clients.ServerConfig{ + ServerURI: test.serverURI, + Extensions: &ServerConfig{Credentials: insecure.NewBundle()}, + } + builder := Builder{} + transport, err := builder.Build(serverCfg) + if err != nil { + t.Fatalf("failed to build transport: %v", err) + } + defer transport.Close() + + ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) + defer cancel() + _, err = transport.NewStream(ctx, "test-method") + if (err != nil) != test.wantErr { + t.Fatalf("transport.NewStream() error = %v, wantErr %v", err, test.wantErr) + } + }) + } +} From eb5e8bddff2a1bbc55f6bef5bf123e738965b989 Mon Sep 17 00:00:00 2001 From: Purnesh Dixit Date: Wed, 5 Feb 2025 22:09:58 +0530 Subject: [PATCH 02/10] remove server config extension interface --- .../clients/grpctransport/grpc_transport.go | 25 +++---------------- .../grpctransport/grpc_transport_test.go | 8 +++--- 2 files changed, 8 insertions(+), 25 deletions(-) diff --git a/xds/internal/clients/grpctransport/grpc_transport.go b/xds/internal/clients/grpctransport/grpc_transport.go index 49d82c5b9ce3..e8341ba80176 100644 --- a/xds/internal/clients/grpctransport/grpc_transport.go +++ b/xds/internal/clients/grpctransport/grpc_transport.go @@ -31,30 +31,14 @@ import ( "google.golang.org/grpc/xds/internal/clients" ) -// ServerConfigExtension allows extending the clients.ServerConfig for the -// gRPC-based transport builder. Any implementation of this must implement -// the ServerConfig() method. -// -// This interface must be implemented by the grpc transport builder provided -// in the Extensions field of the clients.ServerConfig. -type ServerConfigExtension interface { - ServerConfig() *ServerConfig -} - -// ServerConfig holds the settings for connecting to an xDS management server +// ServerConfigExtension holds the settings for connecting to an xDS management server // using gRPC. -type ServerConfig struct { +type ServerConfigExtension struct { // Credentials is the credential bundle containing the gRPC credentials for // connecting to the xDS management server. Credentials credentials.Bundle } -// ServerConfig returns the ServerConfig itself. This method is designed -// to satisfy [ServerConfigExtension] interface requirement. -func (s *ServerConfig) ServerConfig() *ServerConfig { - return s -} - // Builder provides a way to build a gRPC-based transport to an xDS management // server. type Builder struct{} @@ -76,8 +60,7 @@ func (b *Builder) Build(sc clients.ServerConfig) (clients.Transport, error) { if !ok { return nil, fmt.Errorf("xds: ServerConfig's Extensions field cannot be anything other than grpctransport.ServerConfigExtension for gRPC transport") } - gtsc := gtsce.ServerConfig() - if gtsc.Credentials == nil { + if gtsce.Credentials == nil { return nil, fmt.Errorf("xsd: ServerConfigExtensions's Credentials field cannot be nil for gRPC transport") } @@ -93,7 +76,7 @@ func (b *Builder) Build(sc clients.ServerConfig) (clients.Transport, error) { Time: 5 * time.Minute, Timeout: 20 * time.Second, }) - cc, err := grpc.NewClient(sc.ServerURI, kpCfg, grpc.WithCredentialsBundle(gtsc.Credentials)) + cc, err := grpc.NewClient(sc.ServerURI, kpCfg, grpc.WithCredentialsBundle(gtsce.Credentials)) if err != nil { return nil, fmt.Errorf("error creating grpc client for server uri %s, %v", sc.ServerURI, err) } diff --git a/xds/internal/clients/grpctransport/grpc_transport_test.go b/xds/internal/clients/grpctransport/grpc_transport_test.go index fc1635f08b0d..82f21b1c94be 100644 --- a/xds/internal/clients/grpctransport/grpc_transport_test.go +++ b/xds/internal/clients/grpctransport/grpc_transport_test.go @@ -61,7 +61,7 @@ func (s) TestBuild(t *testing.T) { name: "ServerURI_empty", serverCfg: clients.ServerConfig{ ServerURI: "", - Extensions: &ServerConfig{Credentials: insecure.NewBundle()}, + Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, }, wantErr: true, }, @@ -82,7 +82,7 @@ func (s) TestBuild(t *testing.T) { name: "ServerConfigExtension_Credentials_nil", serverCfg: clients.ServerConfig{ ServerURI: "server-address", - Extensions: &ServerConfig{}, + Extensions: ServerConfigExtension{}, }, wantErr: true, }, @@ -90,7 +90,7 @@ func (s) TestBuild(t *testing.T) { name: "success", serverCfg: clients.ServerConfig{ ServerURI: "server-address", - Extensions: &ServerConfig{Credentials: insecure.NewBundle()}, + Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, }, wantErr: false, }, @@ -137,7 +137,7 @@ func (s) TestNewStream(t *testing.T) { t.Run(test.name, func(t *testing.T) { serverCfg := clients.ServerConfig{ ServerURI: test.serverURI, - Extensions: &ServerConfig{Credentials: insecure.NewBundle()}, + Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, } builder := Builder{} transport, err := builder.Build(serverCfg) From dfcd6866e3536b4a25123e7a1c42131d48f9e0aa Mon Sep 17 00:00:00 2001 From: Purnesh Dixit Date: Wed, 5 Feb 2025 23:17:34 +0530 Subject: [PATCH 03/10] add byte codec --- .../clients/grpctransport/grpc_transport.go | 32 ++++++++++++++- .../grpctransport/grpc_transport_test.go | 41 +++++++++++++++++++ 2 files changed, 71 insertions(+), 2 deletions(-) diff --git a/xds/internal/clients/grpctransport/grpc_transport.go b/xds/internal/clients/grpctransport/grpc_transport.go index e8341ba80176..3c18d7110525 100644 --- a/xds/internal/clients/grpctransport/grpc_transport.go +++ b/xds/internal/clients/grpctransport/grpc_transport.go @@ -27,6 +27,7 @@ import ( "google.golang.org/grpc" "google.golang.org/grpc/credentials" + "google.golang.org/grpc/encoding" "google.golang.org/grpc/keepalive" "google.golang.org/grpc/xds/internal/clients" ) @@ -43,9 +44,15 @@ type ServerConfigExtension struct { // server. type Builder struct{} +func init() { + encoding.RegisterCodec(&byteCodec{}) +} + // Build creates a new gRPC-based transport to an xDS management server using // the provided clients.ServerConfig. This involves creating a -// grpc.ClientConn to the server using the provided credentials and server URI. +// grpc.ClientConn to the server using the provided credentials and server URI, +// and byteCodec which is byte-based implementation of encoding.Codec to send +// and receive messages as bytes on the stream. // // If any of ServerURI or Extensions of `sc` are not present, Build() will return // an error. @@ -76,7 +83,7 @@ func (b *Builder) Build(sc clients.ServerConfig) (clients.Transport, error) { Time: 5 * time.Minute, Timeout: 20 * time.Second, }) - cc, err := grpc.NewClient(sc.ServerURI, kpCfg, grpc.WithCredentialsBundle(gtsce.Credentials)) + cc, err := grpc.NewClient(sc.ServerURI, kpCfg, grpc.WithCredentialsBundle(gtsce.Credentials), grpc.WithDefaultCallOptions(grpc.ForceCodec(&byteCodec{}))) if err != nil { return nil, fmt.Errorf("error creating grpc client for server uri %s, %v", sc.ServerURI, err) } @@ -123,3 +130,24 @@ func (s *stream) Recv() ([]byte, error) { func (g *grpcTransport) Close() error { return g.cc.Close() } + +type byteCodec struct{} + +func (c *byteCodec) Marshal(v any) ([]byte, error) { + if b, ok := v.([]byte); ok { + return b, nil + } + return nil, fmt.Errorf("message must be a byte slice") +} + +func (c *byteCodec) Unmarshal(data []byte, v any) error { + if b, ok := v.(*[]byte); ok { + *b = data + return nil + } + return fmt.Errorf("target must be a pointer to a byte slice") +} + +func (c *byteCodec) Name() string { + return "byteCodec" +} diff --git a/xds/internal/clients/grpctransport/grpc_transport_test.go b/xds/internal/clients/grpctransport/grpc_transport_test.go index 82f21b1c94be..bdd9c737bbbe 100644 --- a/xds/internal/clients/grpctransport/grpc_transport_test.go +++ b/xds/internal/clients/grpctransport/grpc_transport_test.go @@ -27,6 +27,9 @@ import ( "google.golang.org/grpc/internal/grpctest" "google.golang.org/grpc/internal/testutils/xds/e2e" "google.golang.org/grpc/xds/internal/clients" + "google.golang.org/protobuf/proto" + + v3discoverypb "github.com/envoyproxy/go-control-plane/envoy/service/discovery/v3" ) const ( @@ -155,3 +158,41 @@ func (s) TestNewStream(t *testing.T) { }) } } + +// TestStream_Send verifies that grpcTransport.Stream.Send() successfully sends +// a message on the stream. It starts a management server to create a stream +// and sends a marshalled DiscoveryRequest proto on it. +func (s) TestStream_Send(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) + defer cancel() + + // Start an xDS management server. + mgmtServer := e2e.StartManagementServer(t, e2e.ManagementServerOptions{}) + + // Build a grpc-based transport to the above xDS management server. + serverCfg := clients.ServerConfig{ + ServerURI: mgmtServer.Address, + Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, + } + builder := Builder{} + transport, err := builder.Build(serverCfg) + if err != nil { + t.Fatalf("failed to build transport: %v", err) + } + defer transport.Close() + + // Create a new stream to the xDS management server. + stream, err := transport.NewStream(ctx, "test-method") + if err != nil { + t.Fatalf("failed to create stream: %v", err) + } + + // Send a discovery request message on the stream. + req, err := proto.Marshal(&v3discoverypb.DiscoveryRequest{}) + if err != nil { + t.Fatalf("failed to marshal DiscoveryRequest: %v", err) + } + if err := stream.Send(req); err != nil { + t.Fatalf("failed to send message: %v", err) + } +} From 56b0d41fc77d39cb56144cabce9490cec48ed200 Mon Sep 17 00:00:00 2001 From: Purnesh Dixit Date: Fri, 7 Feb 2025 17:05:55 +0530 Subject: [PATCH 04/10] dfawley review 1 --- .../clients/grpctransport/grpc_transport.go | 44 +++++++------------ .../grpctransport/grpc_transport_test.go | 8 ++-- 2 files changed, 20 insertions(+), 32 deletions(-) diff --git a/xds/internal/clients/grpctransport/grpc_transport.go b/xds/internal/clients/grpctransport/grpc_transport.go index 3c18d7110525..aaf2ecd25089 100644 --- a/xds/internal/clients/grpctransport/grpc_transport.go +++ b/xds/internal/clients/grpctransport/grpc_transport.go @@ -27,48 +27,38 @@ import ( "google.golang.org/grpc" "google.golang.org/grpc/credentials" - "google.golang.org/grpc/encoding" "google.golang.org/grpc/keepalive" "google.golang.org/grpc/xds/internal/clients" ) -// ServerConfigExtension holds the settings for connecting to an xDS management server -// using gRPC. +// ServerConfigExtension holds settings for connecting to a gRPC server, +// such as an xDS or LRS server. type ServerConfigExtension struct { - // Credentials is the credential bundle containing the gRPC credentials for - // connecting to the xDS management server. + // Credentials will be used for all gRPC transports. If is unset, transport + // creation will fail. Credentials credentials.Bundle } -// Builder provides a way to build a gRPC-based transport to an xDS management -// server. +// Builder creates gRPC-based Transports. It must be paired with ServerConfigs +// that contain its ServerConfigExtension. type Builder struct{} -func init() { - encoding.RegisterCodec(&byteCodec{}) -} - -// Build creates a new gRPC-based transport to an xDS management server using -// the provided clients.ServerConfig. This involves creating a -// grpc.ClientConn to the server using the provided credentials and server URI, -// and byteCodec which is byte-based implementation of encoding.Codec to send -// and receive messages as bytes on the stream. +// Build returns a gRPC-based clients.Transport. // -// If any of ServerURI or Extensions of `sc` are not present, Build() will return -// an error. +// The Extension field of the ServerConfig must be a ServerConfigExtension. func (b *Builder) Build(sc clients.ServerConfig) (clients.Transport, error) { if sc.ServerURI == "" { - return nil, fmt.Errorf("xds: ServerConfig's ServerURI field cannot be empty") + return nil, fmt.Errorf("ServerConfig's ServerURI field cannot be empty") } if sc.Extensions == nil { - return nil, fmt.Errorf("xds: ServerConfig's Extensions field cannot be nil for gRPC transport") + return nil, fmt.Errorf("ServerConfig's Extensions field cannot be nil for gRPC transport") } gtsce, ok := sc.Extensions.(ServerConfigExtension) if !ok { - return nil, fmt.Errorf("xds: ServerConfig's Extensions field cannot be anything other than grpctransport.ServerConfigExtension for gRPC transport") + return nil, fmt.Errorf("ServerConfig Extensions field is %T, but must be %T", sc.Extensions, ServerConfigExtension{}) } if gtsce.Credentials == nil { - return nil, fmt.Errorf("xsd: ServerConfigExtensions's Credentials field cannot be nil for gRPC transport") + return nil, fmt.Errorf("ServerConfigExtensions's Credentials field cannot be nil for gRPC transport") } // TODO: Incorporate reference count map for existing transports and @@ -96,9 +86,7 @@ type grpcTransport struct { cc *grpc.ClientConn } -// NewStream creates a new gRPC stream to the xDS management server for the -// specified method. The returned Stream interface can be used to send and -// receive messages on the stream. +// NewStream creates a new gRPC stream to the server for the specified method. func (g *grpcTransport) NewStream(ctx context.Context, method string) (clients.Stream, error) { s, err := g.cc.NewStream(ctx, &grpc.StreamDesc{StreamName: method, ClientStreams: true, ServerStreams: true}, method) if err != nil { @@ -111,12 +99,12 @@ type stream struct { stream grpc.ClientStream } -// Send sends a message to the xDS management server. +// Send sends a message to the server. func (s *stream) Send(msg []byte) error { return s.stream.SendMsg(msg) } -// Recv receives a message from the xDS management server. +// Recv receives a message from the server. func (s *stream) Recv() ([]byte, error) { var typedRes []byte err := s.stream.RecvMsg(&typedRes) @@ -126,7 +114,7 @@ func (s *stream) Recv() ([]byte, error) { return typedRes, nil } -// Close closes the gRPC stream to the xDS management server. +// Close closes all the gRPC streams to the server. func (g *grpcTransport) Close() error { return g.cc.Close() } diff --git a/xds/internal/clients/grpctransport/grpc_transport_test.go b/xds/internal/clients/grpctransport/grpc_transport_test.go index bdd9c737bbbe..f546e681bd27 100644 --- a/xds/internal/clients/grpctransport/grpc_transport_test.go +++ b/xds/internal/clients/grpctransport/grpc_transport_test.go @@ -61,7 +61,7 @@ func (s) TestBuild(t *testing.T) { wantErr bool }{ { - name: "ServerURI_empty", + name: "ServerURI_is_empty", serverCfg: clients.ServerConfig{ ServerURI: "", Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, @@ -69,12 +69,12 @@ func (s) TestBuild(t *testing.T) { wantErr: true, }, { - name: "Extensions_nil", + name: "Extensions is nil", serverCfg: clients.ServerConfig{ServerURI: "server-address"}, wantErr: true, }, { - name: "Extensions_not_ServerConfigExtension", + name: "Extensions is not a ServerConfigExtension", serverCfg: clients.ServerConfig{ ServerURI: "server-address", Extensions: 1, @@ -82,7 +82,7 @@ func (s) TestBuild(t *testing.T) { wantErr: true, }, { - name: "ServerConfigExtension_Credentials_nil", + name: "ServerConfigExtension Credentials is nil", serverCfg: clients.ServerConfig{ ServerURI: "server-address", Extensions: ServerConfigExtension{}, From 61bfddad33e8513843cc6c646bf4561054e99378 Mon Sep 17 00:00:00 2001 From: Purnesh Dixit Date: Sun, 9 Feb 2025 18:11:15 +0530 Subject: [PATCH 05/10] send and recv tests with byte based test server --- .../clients/grpctransport/grpc_transport.go | 1 - .../grpctransport/grpc_transport_test.go | 154 +++++++++++++++--- 2 files changed, 135 insertions(+), 20 deletions(-) diff --git a/xds/internal/clients/grpctransport/grpc_transport.go b/xds/internal/clients/grpctransport/grpc_transport.go index aaf2ecd25089..ea9e244be3c1 100644 --- a/xds/internal/clients/grpctransport/grpc_transport.go +++ b/xds/internal/clients/grpctransport/grpc_transport.go @@ -77,7 +77,6 @@ func (b *Builder) Build(sc clients.ServerConfig) (clients.Transport, error) { if err != nil { return nil, fmt.Errorf("error creating grpc client for server uri %s, %v", sc.ServerURI, err) } - cc.Connect() return &grpcTransport{cc: cc}, nil } diff --git a/xds/internal/clients/grpctransport/grpc_transport_test.go b/xds/internal/clients/grpctransport/grpc_transport_test.go index f546e681bd27..638399b5522d 100644 --- a/xds/internal/clients/grpctransport/grpc_transport_test.go +++ b/xds/internal/clients/grpctransport/grpc_transport_test.go @@ -20,21 +20,29 @@ package grpctransport import ( "context" + "io" + "net" "testing" "time" + "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/internal/grpctest" - "google.golang.org/grpc/internal/testutils/xds/e2e" "google.golang.org/grpc/xds/internal/clients" "google.golang.org/protobuf/proto" + "google.golang.org/protobuf/testing/protocmp" v3discoverypb "github.com/envoyproxy/go-control-plane/envoy/service/discovery/v3" + "github.com/google/go-cmp/cmp" ) const ( - defaultTestWatchExpiryTimeout = 500 * time.Millisecond - defaultTestTimeout = 10 * time.Second + defaultTestTimeout = 10 * time.Second +) + +var ( + testDiscoverRequest = &v3discoverypb.DiscoveryRequest{VersionInfo: "1"} + testDiscoverResponse = &v3discoverypb.DiscoveryResponse{VersionInfo: "1"} ) type s struct { @@ -45,6 +53,77 @@ func Test(t *testing.T) { grpctest.RunSubTests(t, s{}) } +type testServer struct { + lis net.Listener // listener used by the test gRPC server + requestChan chan []byte // channel to send received requests on to verify +} + +// setupTestServer starts a gRPC test server that uses the same byteCodec as +// grpcTransport. It registers a streaming handler for the "test.Service/Stream" +// method and returns a testServer struct that contains the listener and a +// channel for received requests from the client on the stream. +func setupTestServer(t *testing.T) *testServer { + lis, err := net.Listen("tcp", "localhost:0") + if err != nil { + t.Fatalf("Failed to listen on localhost:0: %v", err) + } + ts := &testServer{ + requestChan: make(chan []byte, 100), + lis: lis, + } + + s := grpc.NewServer(grpc.ForceServerCodec(&byteCodec{})) + s.RegisterService(&grpc.ServiceDesc{ + ServiceName: "test.Service", + HandlerType: (*any)(nil), + Streams: []grpc.StreamDesc{ + { + StreamName: "Stream", + Handler: ts.streamHandler, + ServerStreams: true, + ClientStreams: true, + }, + }, + }, struct{}{}) + go func() { + if err := s.Serve(lis); err != nil { + t.Logf("Server exited with error: %v", err) + } + }() + t.Cleanup(s.Stop) + + return ts +} + +// streamHandler is the handler for the "test.Service/Stream" method. It waits +// for a message from the client on the stream, and then sends a discovery +// response message back to the client. It also put the received message in +// requestChan for client to verify if the correct request was received. It +// continues until the client closes the stream. +func (s *testServer) streamHandler(_ any, stream grpc.ServerStream) error { + for { + var msg []byte + err := stream.RecvMsg(&msg) + if err == io.EOF { + // read done. + return nil + } + if err != nil { + return err + } + s.requestChan <- msg + + // Send a discovery response message on the stream. + res, err := proto.Marshal(testDiscoverResponse) + if err != nil { + return err + } + if err := stream.SendMsg(res); err != nil { + return err + } + } +} + // TestBuild verifies that the grpctransport.Builder creates a new // grpc.ClientConn every time Build() is called. // @@ -116,9 +195,9 @@ func (s) TestBuild(t *testing.T) { } // TestNewStream verifies that grpcTransport.NewStream() successfully creates a -// new client stream for the xDS management server. +// new client stream for the server. func (s) TestNewStream(t *testing.T) { - mgmtServer := e2e.StartManagementServer(t, e2e.ManagementServerOptions{}) + ts := setupTestServer(t) tests := []struct { name string @@ -127,7 +206,7 @@ func (s) TestNewStream(t *testing.T) { }{ { name: "success", - serverURI: mgmtServer.Address, + serverURI: ts.lis.Addr().String(), wantErr: false, }, { @@ -151,7 +230,7 @@ func (s) TestNewStream(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) defer cancel() - _, err = transport.NewStream(ctx, "test-method") + _, err = transport.NewStream(ctx, "/test.Service/Stream") if (err != nil) != test.wantErr { t.Fatalf("transport.NewStream() error = %v, wantErr %v", err, test.wantErr) } @@ -159,19 +238,24 @@ func (s) TestNewStream(t *testing.T) { } } -// TestStream_Send verifies that grpcTransport.Stream.Send() successfully sends -// a message on the stream. It starts a management server to create a stream -// and sends a marshalled DiscoveryRequest proto on it. -func (s) TestStream_Send(t *testing.T) { +// TestStream_SendAndRecv verifies that grpcTransport.Stream.Send() +// and grpcTransport.Stream.Recv() successfully send and receive messages +// on the stream. +// +// It starts a gRPC test server using setupTestServer(). The test then sends a +// testDiscoverRequest on the stream and verifies that the received discovery +// on the server is same as sent. It then wait to receive a +// testDiscoverResponse from the server and verifies that the received +// discovery response is same as sent from the server. +func (s) TestStream_SendAndRecv(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) defer cancel() - // Start an xDS management server. - mgmtServer := e2e.StartManagementServer(t, e2e.ManagementServerOptions{}) + ts := setupTestServer(t) - // Build a grpc-based transport to the above xDS management server. + // Build a grpc-based transport to the above server. serverCfg := clients.ServerConfig{ - ServerURI: mgmtServer.Address, + ServerURI: ts.lis.Addr().String(), Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, } builder := Builder{} @@ -181,18 +265,50 @@ func (s) TestStream_Send(t *testing.T) { } defer transport.Close() - // Create a new stream to the xDS management server. - stream, err := transport.NewStream(ctx, "test-method") + // Create a new stream to the server. + stream, err := transport.NewStream(ctx, "/test.Service/Stream") if err != nil { t.Fatalf("failed to create stream: %v", err) } // Send a discovery request message on the stream. - req, err := proto.Marshal(&v3discoverypb.DiscoveryRequest{}) + msg, err := proto.Marshal(testDiscoverRequest) if err != nil { t.Fatalf("failed to marshal DiscoveryRequest: %v", err) } - if err := stream.Send(req); err != nil { + if err := stream.Send(msg); err != nil { t.Fatalf("failed to send message: %v", err) } + // Verify that the DiscoveryRequest received on the server was same as + // sent. + select { + case res, ok := <-ts.requestChan: + if !ok { + t.Fatalf("ts.requestChan is closed") + } + var gotReq v3discoverypb.DiscoveryRequest + if err := proto.Unmarshal(res, &gotReq); err != nil { + t.Fatalf("failed to unmarshal response from ts.requestChan to DiscoveryRequest: %v", err) + } + if !cmp.Equal(testDiscoverRequest, &gotReq, protocmp.Transform()) { + t.Fatalf("<-ts.requestChan = %v, want %v", &gotReq, &testDiscoverRequest) + } + case <-ctx.Done(): + t.Fatalf("timeout waiting for request to reach server") + } + + // Wait until response message is received from the server. + res, err := stream.Recv() + if err != nil { + t.Fatalf("failed to receive message: %v", err) + } + // Verify that the DiscoveryResponse received was same as sent from the + // server. + var gotRes v3discoverypb.DiscoveryResponse + if err := proto.Unmarshal(res, &gotRes); err != nil { + t.Fatalf("failed to unmarshal response from ts.requestChan to DiscoveryRequest: %v", err) + } + if !cmp.Equal(testDiscoverResponse, &gotRes, protocmp.Transform()) { + t.Fatalf("proto.Unmarshal(res, &gotRes) = %v, want %v", &gotRes, testDiscoverResponse) + } } From 84b4f2bde6724dc46321657f9dc48120611d6d0e Mon Sep 17 00:00:00 2001 From: Purnesh Dixit Date: Wed, 12 Feb 2025 00:54:11 +0700 Subject: [PATCH 06/10] change to proto based server --- .../clients/grpctransport/grpc_transport.go | 15 ++-- .../grpctransport/grpc_transport_test.go | 87 ++++++++----------- 2 files changed, 45 insertions(+), 57 deletions(-) diff --git a/xds/internal/clients/grpctransport/grpc_transport.go b/xds/internal/clients/grpctransport/grpc_transport.go index ea9e244be3c1..53a5a8e17756 100644 --- a/xds/internal/clients/grpctransport/grpc_transport.go +++ b/xds/internal/clients/grpctransport/grpc_transport.go @@ -32,10 +32,10 @@ import ( ) // ServerConfigExtension holds settings for connecting to a gRPC server, -// such as an xDS or LRS server. +// such as an xDS management or an LRS server. type ServerConfigExtension struct { - // Credentials will be used for all gRPC transports. If is unset, transport - // creation will fail. + // Credentials will be used for all gRPC transports. If it is unset, + // transport creation will fail. Credentials credentials.Bundle } @@ -62,13 +62,14 @@ func (b *Builder) Build(sc clients.ServerConfig) (clients.Transport, error) { } // TODO: Incorporate reference count map for existing transports and - // deduplicate transports based on server URI and credentials so that + // deduplicate transports based on the provided ServerConfig so that // transport channel to same server can be shared between xDS and LRS // client. - // Dial the xDS management server with the provided credentials, server URI, - // and a static keepalive configuration that is common across gRPC language - // implementations. + // Create a new gRPC client/channel for the xDS management server with the + // provided credentials, server URI, and a byte codec to send and receive + // messages. Aso set a static keepalive configuration that is common across + // gRPC language implementations. kpCfg := grpc.WithKeepaliveParams(keepalive.ClientParameters{ Time: 5 * time.Minute, Timeout: 20 * time.Second, diff --git a/xds/internal/clients/grpctransport/grpc_transport_test.go b/xds/internal/clients/grpctransport/grpc_transport_test.go index 638399b5522d..5aea70cc118f 100644 --- a/xds/internal/clients/grpctransport/grpc_transport_test.go +++ b/xds/internal/clients/grpctransport/grpc_transport_test.go @@ -25,6 +25,7 @@ import ( "testing" "time" + "github.com/google/go-cmp/cmp" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/internal/grpctest" @@ -32,8 +33,8 @@ import ( "google.golang.org/protobuf/proto" "google.golang.org/protobuf/testing/protocmp" + v3discoverygrpc "github.com/envoyproxy/go-control-plane/envoy/service/discovery/v3" v3discoverypb "github.com/envoyproxy/go-control-plane/envoy/service/discovery/v3" - "github.com/google/go-cmp/cmp" ) const ( @@ -53,38 +54,30 @@ func Test(t *testing.T) { grpctest.RunSubTests(t, s{}) } +// testServer implements the AggregatedDiscoveryServiceServer interface to test +// the gRPC transport implementation. type testServer struct { - lis net.Listener // listener used by the test gRPC server - requestChan chan []byte // channel to send received requests on to verify + v3discoverygrpc.UnimplementedAggregatedDiscoveryServiceServer + + lis net.Listener // listener used by the test server + requestChan chan *v3discoverypb.DiscoveryRequest // channel to send the received requests on for verification } -// setupTestServer starts a gRPC test server that uses the same byteCodec as -// grpcTransport. It registers a streaming handler for the "test.Service/Stream" -// method and returns a testServer struct that contains the listener and a -// channel for received requests from the client on the stream. +// setupTestServer set up the gRPC server for AggregatedDiscoveryService. It +// creates an instance of testServer and registers it with a gRPC server. func setupTestServer(t *testing.T) *testServer { lis, err := net.Listen("tcp", "localhost:0") if err != nil { t.Fatalf("Failed to listen on localhost:0: %v", err) } ts := &testServer{ - requestChan: make(chan []byte, 100), + requestChan: make(chan *v3discoverypb.DiscoveryRequest), lis: lis, } - s := grpc.NewServer(grpc.ForceServerCodec(&byteCodec{})) - s.RegisterService(&grpc.ServiceDesc{ - ServiceName: "test.Service", - HandlerType: (*any)(nil), - Streams: []grpc.StreamDesc{ - { - StreamName: "Stream", - Handler: ts.streamHandler, - ServerStreams: true, - ClientStreams: true, - }, - }, - }, struct{}{}) + s := grpc.NewServer() + + v3discoverygrpc.RegisterAggregatedDiscoveryServiceServer(s, ts) go func() { if err := s.Serve(lis); err != nil { t.Logf("Server exited with error: %v", err) @@ -95,30 +88,26 @@ func setupTestServer(t *testing.T) *testServer { return ts } -// streamHandler is the handler for the "test.Service/Stream" method. It waits -// for a message from the client on the stream, and then sends a discovery -// response message back to the client. It also put the received message in -// requestChan for client to verify if the correct request was received. It -// continues until the client closes the stream. -func (s *testServer) streamHandler(_ any, stream grpc.ServerStream) error { +// StreamAggregatedResources handles bidirectional streaming of +// DiscoveryRequest and DiscoveryResponse. It waits for a message from the +// client on the stream, and then sends a discovery response message back to +// the client. It also put the received message in requestChan for client to +// verify if the correct request was received. It continues until the client +// closes the stream. +func (s *testServer) StreamAggregatedResources(stream v3discoverygrpc.AggregatedDiscoveryService_StreamAggregatedResourcesServer) error { for { - var msg []byte - err := stream.RecvMsg(&msg) + // Receive a DiscoveryRequest from the client + req, err := stream.Recv() if err == io.EOF { - // read done. - return nil + return nil // Stream closed by client } if err != nil { - return err + return err // Handle other errors } - s.requestChan <- msg + s.requestChan <- req - // Send a discovery response message on the stream. - res, err := proto.Marshal(testDiscoverResponse) - if err != nil { - return err - } - if err := stream.SendMsg(res); err != nil { + // Send the response back to the client + if err := stream.Send(testDiscoverResponse); err != nil { return err } } @@ -140,7 +129,7 @@ func (s) TestBuild(t *testing.T) { wantErr bool }{ { - name: "ServerURI_is_empty", + name: "ServerURI is empty", serverCfg: clients.ServerConfig{ ServerURI: "", Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, @@ -230,7 +219,7 @@ func (s) TestNewStream(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) defer cancel() - _, err = transport.NewStream(ctx, "/test.Service/Stream") + _, err = transport.NewStream(ctx, "/envoy.service.discovery.v3.AggregatedDiscoveryService/StreamAggregatedResources") if (err != nil) != test.wantErr { t.Fatalf("transport.NewStream() error = %v, wantErr %v", err, test.wantErr) } @@ -240,7 +229,7 @@ func (s) TestNewStream(t *testing.T) { // TestStream_SendAndRecv verifies that grpcTransport.Stream.Send() // and grpcTransport.Stream.Recv() successfully send and receive messages -// on the stream. +// on the stream to and from the gRPC server. // // It starts a gRPC test server using setupTestServer(). The test then sends a // testDiscoverRequest on the stream and verifies that the received discovery @@ -248,7 +237,7 @@ func (s) TestNewStream(t *testing.T) { // testDiscoverResponse from the server and verifies that the received // discovery response is same as sent from the server. func (s) TestStream_SendAndRecv(t *testing.T) { - ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) + ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout*2000) defer cancel() ts := setupTestServer(t) @@ -266,7 +255,7 @@ func (s) TestStream_SendAndRecv(t *testing.T) { defer transport.Close() // Create a new stream to the server. - stream, err := transport.NewStream(ctx, "/test.Service/Stream") + stream, err := transport.NewStream(ctx, "/envoy.service.discovery.v3.AggregatedDiscoveryService/StreamAggregatedResources") if err != nil { t.Fatalf("failed to create stream: %v", err) } @@ -279,18 +268,15 @@ func (s) TestStream_SendAndRecv(t *testing.T) { if err := stream.Send(msg); err != nil { t.Fatalf("failed to send message: %v", err) } + // Verify that the DiscoveryRequest received on the server was same as // sent. select { - case res, ok := <-ts.requestChan: + case gotReq, ok := <-ts.requestChan: if !ok { t.Fatalf("ts.requestChan is closed") } - var gotReq v3discoverypb.DiscoveryRequest - if err := proto.Unmarshal(res, &gotReq); err != nil { - t.Fatalf("failed to unmarshal response from ts.requestChan to DiscoveryRequest: %v", err) - } - if !cmp.Equal(testDiscoverRequest, &gotReq, protocmp.Transform()) { + if !cmp.Equal(testDiscoverRequest, gotReq, protocmp.Transform()) { t.Fatalf("<-ts.requestChan = %v, want %v", &gotReq, &testDiscoverRequest) } case <-ctx.Done(): @@ -302,6 +288,7 @@ func (s) TestStream_SendAndRecv(t *testing.T) { if err != nil { t.Fatalf("failed to receive message: %v", err) } + // Verify that the DiscoveryResponse received was same as sent from the // server. var gotRes v3discoverypb.DiscoveryResponse From 9049d97d3ff5fb433064d179a6940bcbcd8d4735 Mon Sep 17 00:00:00 2001 From: Purnesh Dixit Date: Thu, 20 Feb 2025 00:50:04 +0530 Subject: [PATCH 07/10] easwar review 1 --- .../clients/grpctransport/grpc_transport.go | 24 ++-- .../grpctransport/grpc_transport_test.go | 127 +++++++++--------- 2 files changed, 76 insertions(+), 75 deletions(-) diff --git a/xds/internal/clients/grpctransport/grpc_transport.go b/xds/internal/clients/grpctransport/grpc_transport.go index 53a5a8e17756..06eef206b2f7 100644 --- a/xds/internal/clients/grpctransport/grpc_transport.go +++ b/xds/internal/clients/grpctransport/grpc_transport.go @@ -40,7 +40,7 @@ type ServerConfigExtension struct { } // Builder creates gRPC-based Transports. It must be paired with ServerConfigs -// that contain its ServerConfigExtension. +// that contain Extension field of type ServerConfigExtension. type Builder struct{} // Build returns a gRPC-based clients.Transport. @@ -53,11 +53,11 @@ func (b *Builder) Build(sc clients.ServerConfig) (clients.Transport, error) { if sc.Extensions == nil { return nil, fmt.Errorf("ServerConfig's Extensions field cannot be nil for gRPC transport") } - gtsce, ok := sc.Extensions.(ServerConfigExtension) + sce, ok := sc.Extensions.(ServerConfigExtension) if !ok { return nil, fmt.Errorf("ServerConfig Extensions field is %T, but must be %T", sc.Extensions, ServerConfigExtension{}) } - if gtsce.Credentials == nil { + if sce.Credentials == nil { return nil, fmt.Errorf("ServerConfigExtensions's Credentials field cannot be nil for gRPC transport") } @@ -66,15 +66,15 @@ func (b *Builder) Build(sc clients.ServerConfig) (clients.Transport, error) { // transport channel to same server can be shared between xDS and LRS // client. - // Create a new gRPC client/channel for the xDS management server with the - // provided credentials, server URI, and a byte codec to send and receive - // messages. Aso set a static keepalive configuration that is common across - // gRPC language implementations. + // Create a new gRPC client/channel for the server with the provided + // credentials, server URI, and a byte codec to send and receive messages. + // Also set a static keepalive configuration that is common across gRPC + // language implementations. kpCfg := grpc.WithKeepaliveParams(keepalive.ClientParameters{ Time: 5 * time.Minute, Timeout: 20 * time.Second, }) - cc, err := grpc.NewClient(sc.ServerURI, kpCfg, grpc.WithCredentialsBundle(gtsce.Credentials), grpc.WithDefaultCallOptions(grpc.ForceCodec(&byteCodec{}))) + cc, err := grpc.NewClient(sc.ServerURI, kpCfg, grpc.WithCredentialsBundle(sce.Credentials), grpc.WithDefaultCallOptions(grpc.ForceCodec(&byteCodec{}))) if err != nil { return nil, fmt.Errorf("error creating grpc client for server uri %s, %v", sc.ServerURI, err) } @@ -107,14 +107,14 @@ func (s *stream) Send(msg []byte) error { // Recv receives a message from the server. func (s *stream) Recv() ([]byte, error) { var typedRes []byte - err := s.stream.RecvMsg(&typedRes) - if err != nil { + + if err := s.stream.RecvMsg(&typedRes); err != nil { return typedRes, err } return typedRes, nil } -// Close closes all the gRPC streams to the server. +// Close closes the gRPC channel to the server. func (g *grpcTransport) Close() error { return g.cc.Close() } @@ -137,5 +137,5 @@ func (c *byteCodec) Unmarshal(data []byte, v any) error { } func (c *byteCodec) Name() string { - return "byteCodec" + return "grpc.xds.internal.clients.grpctransport.byte_codec" } diff --git a/xds/internal/clients/grpctransport/grpc_transport_test.go b/xds/internal/clients/grpctransport/grpc_transport_test.go index 5aea70cc118f..b78be35e197d 100644 --- a/xds/internal/clients/grpctransport/grpc_transport_test.go +++ b/xds/internal/clients/grpctransport/grpc_transport_test.go @@ -41,11 +41,6 @@ const ( defaultTestTimeout = 10 * time.Second ) -var ( - testDiscoverRequest = &v3discoverypb.DiscoveryRequest{VersionInfo: "1"} - testDiscoverResponse = &v3discoverypb.DiscoveryResponse{VersionInfo: "1"} -) - type s struct { grpctest.Tester } @@ -59,20 +54,23 @@ func Test(t *testing.T) { type testServer struct { v3discoverygrpc.UnimplementedAggregatedDiscoveryServiceServer - lis net.Listener // listener used by the test server + address string // address of the server requestChan chan *v3discoverypb.DiscoveryRequest // channel to send the received requests on for verification + response *v3discoverypb.DiscoveryResponse // response to send back to the client from handler } // setupTestServer set up the gRPC server for AggregatedDiscoveryService. It -// creates an instance of testServer and registers it with a gRPC server. -func setupTestServer(t *testing.T) *testServer { +// creates an instance of testServer that returns the provided response from +// the StreamAggregatedResources() handler and registers it with a gRPC server. +func setupTestServer(t *testing.T, response *v3discoverypb.DiscoveryResponse) *testServer { lis, err := net.Listen("tcp", "localhost:0") if err != nil { t.Fatalf("Failed to listen on localhost:0: %v", err) } ts := &testServer{ requestChan: make(chan *v3discoverypb.DiscoveryRequest), - lis: lis, + address: lis.Addr().String(), + response: response, } s := grpc.NewServer() @@ -107,7 +105,7 @@ func (s *testServer) StreamAggregatedResources(stream v3discoverygrpc.Aggregated s.requestChan <- req // Send the response back to the client - if err := stream.Send(testDiscoverResponse); err != nil { + if err := stream.Send(s.response); err != nil { return err } } @@ -179,51 +177,56 @@ func (s) TestBuild(t *testing.T) { if !test.wantErr && tr == nil { t.Fatalf("got non-nil transport from Build(), want nil") } + if test.wantErr && tr != nil { + t.Fatalf("got nil transport from Build(), want non-nil") + } }) } } -// TestNewStream verifies that grpcTransport.NewStream() successfully creates a -// new client stream for the server. -func (s) TestNewStream(t *testing.T) { - ts := setupTestServer(t) +// TestNewStream_Success verifies that grpcTransport.NewStream() successfully +// creates a new client stream for the server when provided a valid server URI. +func (s) TestNewStream_Success(t *testing.T) { + ts := setupTestServer(t, &v3discoverypb.DiscoveryResponse{VersionInfo: "1"}) - tests := []struct { - name string - serverURI string - wantErr bool - }{ - { - name: "success", - serverURI: ts.lis.Addr().String(), - wantErr: false, - }, - { - name: "error", - serverURI: "invalid-server-uri", - wantErr: true, - }, + serverCfg := clients.ServerConfig{ + ServerURI: ts.address, + Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, } - for _, test := range tests { - t.Run(test.name, func(t *testing.T) { - serverCfg := clients.ServerConfig{ - ServerURI: test.serverURI, - Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, - } - builder := Builder{} - transport, err := builder.Build(serverCfg) - if err != nil { - t.Fatalf("failed to build transport: %v", err) - } - defer transport.Close() + builder := Builder{} + transport, err := builder.Build(serverCfg) + if err != nil { + t.Fatalf("Failed to build transport: %v", err) + } + defer transport.Close() - ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) - defer cancel() - _, err = transport.NewStream(ctx, "/envoy.service.discovery.v3.AggregatedDiscoveryService/StreamAggregatedResources") - if (err != nil) != test.wantErr { - t.Fatalf("transport.NewStream() error = %v, wantErr %v", err, test.wantErr) - } - }) + ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) + defer cancel() + _, err = transport.NewStream(ctx, "/envoy.service.discovery.v3.AggregatedDiscoveryService/StreamAggregatedResources") + if err != nil { + t.Fatalf("transport.NewStream() failed: %v", err) + } +} + +// TestNewStream_Error verifies that grpcTransport.NewStream() returns an error +// when attempting to create a stream with an invalid server URI. +func (s) TestNewStream_Error(t *testing.T) { + serverCfg := clients.ServerConfig{ + ServerURI: "invalid-server-uri", + Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, + } + builder := Builder{} + transport, err := builder.Build(serverCfg) + if err != nil { + t.Fatalf("Failed to build transport: %v", err) + } + defer transport.Close() + + ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) + defer cancel() + _, err = transport.NewStream(ctx, "/envoy.service.discovery.v3.AggregatedDiscoveryService/StreamAggregatedResources") + if err == nil { + t.Fatal("transport.NewStream() succeeded, want failure") } } @@ -240,43 +243,41 @@ func (s) TestStream_SendAndRecv(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout*2000) defer cancel() - ts := setupTestServer(t) + ts := setupTestServer(t, &v3discoverypb.DiscoveryResponse{VersionInfo: "1"}) // Build a grpc-based transport to the above server. serverCfg := clients.ServerConfig{ - ServerURI: ts.lis.Addr().String(), + ServerURI: ts.address, Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, } builder := Builder{} transport, err := builder.Build(serverCfg) if err != nil { - t.Fatalf("failed to build transport: %v", err) + t.Fatalf("Failed to build transport: %v", err) } defer transport.Close() // Create a new stream to the server. stream, err := transport.NewStream(ctx, "/envoy.service.discovery.v3.AggregatedDiscoveryService/StreamAggregatedResources") if err != nil { - t.Fatalf("failed to create stream: %v", err) + t.Fatalf("Failed to create stream: %v", err) } // Send a discovery request message on the stream. + testDiscoverRequest := &v3discoverypb.DiscoveryRequest{VersionInfo: "1"} msg, err := proto.Marshal(testDiscoverRequest) if err != nil { - t.Fatalf("failed to marshal DiscoveryRequest: %v", err) + t.Fatalf("Failed to marshal DiscoveryRequest: %v", err) } if err := stream.Send(msg); err != nil { - t.Fatalf("failed to send message: %v", err) + t.Fatalf("Failed to send message: %v", err) } // Verify that the DiscoveryRequest received on the server was same as // sent. select { - case gotReq, ok := <-ts.requestChan: - if !ok { - t.Fatalf("ts.requestChan is closed") - } - if !cmp.Equal(testDiscoverRequest, gotReq, protocmp.Transform()) { + case gotReq := <-ts.requestChan: + if !cmp.Equal(gotReq, testDiscoverRequest, protocmp.Transform()) { t.Fatalf("<-ts.requestChan = %v, want %v", &gotReq, &testDiscoverRequest) } case <-ctx.Done(): @@ -286,16 +287,16 @@ func (s) TestStream_SendAndRecv(t *testing.T) { // Wait until response message is received from the server. res, err := stream.Recv() if err != nil { - t.Fatalf("failed to receive message: %v", err) + t.Fatalf("Failed to receive message: %v", err) } // Verify that the DiscoveryResponse received was same as sent from the // server. var gotRes v3discoverypb.DiscoveryResponse if err := proto.Unmarshal(res, &gotRes); err != nil { - t.Fatalf("failed to unmarshal response from ts.requestChan to DiscoveryRequest: %v", err) + t.Fatalf("Failed to unmarshal response from ts.requestChan to DiscoveryResponse: %v", err) } - if !cmp.Equal(testDiscoverResponse, &gotRes, protocmp.Transform()) { - t.Fatalf("proto.Unmarshal(res, &gotRes) = %v, want %v", &gotRes, testDiscoverResponse) + if !cmp.Equal(&gotRes, ts.response, protocmp.Transform()) { + t.Fatalf("proto.Unmarshal(res, &gotRes) = %v, want %v", &gotRes, ts.response) } } From 505036003b901a92dcde45af56f85f00cb5020a5 Mon Sep 17 00:00:00 2001 From: Purnesh Dixit Date: Fri, 21 Feb 2025 14:03:06 +0530 Subject: [PATCH 08/10] easwar review 3 --- .../clients/grpctransport/grpc_transport.go | 16 ++-- .../grpctransport/grpc_transport_test.go | 80 ++++++++++--------- 2 files changed, 50 insertions(+), 46 deletions(-) diff --git a/xds/internal/clients/grpctransport/grpc_transport.go b/xds/internal/clients/grpctransport/grpc_transport.go index 06eef206b2f7..be1f75e747cb 100644 --- a/xds/internal/clients/grpctransport/grpc_transport.go +++ b/xds/internal/clients/grpctransport/grpc_transport.go @@ -48,17 +48,17 @@ type Builder struct{} // The Extension field of the ServerConfig must be a ServerConfigExtension. func (b *Builder) Build(sc clients.ServerConfig) (clients.Transport, error) { if sc.ServerURI == "" { - return nil, fmt.Errorf("ServerConfig's ServerURI field cannot be empty") + return nil, fmt.Errorf("grpctransport: ServerURI is not set in ServerConfig") } if sc.Extensions == nil { - return nil, fmt.Errorf("ServerConfig's Extensions field cannot be nil for gRPC transport") + return nil, fmt.Errorf("grpctransport: Extensions is not set in ServerConfig") } sce, ok := sc.Extensions.(ServerConfigExtension) if !ok { - return nil, fmt.Errorf("ServerConfig Extensions field is %T, but must be %T", sc.Extensions, ServerConfigExtension{}) + return nil, fmt.Errorf("grpctransport: Extensions field is %T, but must be %T in ServerConfig", sc.Extensions, ServerConfigExtension{}) } if sce.Credentials == nil { - return nil, fmt.Errorf("ServerConfigExtensions's Credentials field cannot be nil for gRPC transport") + return nil, fmt.Errorf("grptransport: Credentials field is not set in ServerConfigExtension") } // TODO: Incorporate reference count map for existing transports and @@ -76,7 +76,7 @@ func (b *Builder) Build(sc clients.ServerConfig) (clients.Transport, error) { }) cc, err := grpc.NewClient(sc.ServerURI, kpCfg, grpc.WithCredentialsBundle(sce.Credentials), grpc.WithDefaultCallOptions(grpc.ForceCodec(&byteCodec{}))) if err != nil { - return nil, fmt.Errorf("error creating grpc client for server uri %s, %v", sc.ServerURI, err) + return nil, fmt.Errorf("grpctransport: failed to create transport to server %q: %v", sc.ServerURI, err) } return &grpcTransport{cc: cc}, nil @@ -125,7 +125,7 @@ func (c *byteCodec) Marshal(v any) ([]byte, error) { if b, ok := v.([]byte); ok { return b, nil } - return nil, fmt.Errorf("message must be a byte slice") + return nil, fmt.Errorf("message is %T, but must be a byte slice ", v) } func (c *byteCodec) Unmarshal(data []byte, v any) error { @@ -133,9 +133,9 @@ func (c *byteCodec) Unmarshal(data []byte, v any) error { *b = data return nil } - return fmt.Errorf("target must be a pointer to a byte slice") + return fmt.Errorf("target is %T, but must be a pointer to a byte slice", v) } func (c *byteCodec) Name() string { - return "grpc.xds.internal.clients.grpctransport.byte_codec" + return "grpctransport.byteCodec" } diff --git a/xds/internal/clients/grpctransport/grpc_transport_test.go b/xds/internal/clients/grpctransport/grpc_transport_test.go index b78be35e197d..4ad2f39e92d8 100644 --- a/xds/internal/clients/grpctransport/grpc_transport_test.go +++ b/xds/internal/clients/grpctransport/grpc_transport_test.go @@ -63,6 +63,8 @@ type testServer struct { // creates an instance of testServer that returns the provided response from // the StreamAggregatedResources() handler and registers it with a gRPC server. func setupTestServer(t *testing.T, response *v3discoverypb.DiscoveryResponse) *testServer { + t.Helper() + lis, err := net.Listen("tcp", "localhost:0") if err != nil { t.Fatalf("Failed to listen on localhost:0: %v", err) @@ -111,20 +113,41 @@ func (s *testServer) StreamAggregatedResources(stream v3discoverygrpc.Aggregated } } -// TestBuild verifies that the grpctransport.Builder creates a new -// grpc.ClientConn every time Build() is called. +// TestBuild_Success verifies that the Builder successfully creates a new +// Transport with a non-nil grpc.ClientConn. +func (s) TestBuild_Success(t *testing.T) { + serverCfg := clients.ServerConfig{ + ServerURI: "server-address", + Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, + } + + b := &Builder{} + tr, err := b.Build(serverCfg) + if err != nil { + t.Fatalf("Build() failed: %v", err) + } + defer tr.Close() + + if tr == nil { + t.Fatalf("Got nil transport from Build(), want non-nil") + } + if tr.(*grpcTransport).cc == nil { + t.Fatalf("Got nil grpc.ClientConn in transport, want non-nil") + } +} + +// TestBuild_Failure verifies that the Builder returns error when incorrect +// ServerConfig is provided. // // It covers the following scenarios: // - ServerURI is empty. // - Extensions is nil. // - Extensions is not ServerConfigExtension. // - Credentials are nil. -// - Success cases. -func (s) TestBuild(t *testing.T) { +func (s) TestBuild_Failure(t *testing.T) { tests := []struct { name string serverCfg clients.ServerConfig - wantErr bool }{ { name: "ServerURI is empty", @@ -132,12 +155,10 @@ func (s) TestBuild(t *testing.T) { ServerURI: "", Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, }, - wantErr: true, }, { name: "Extensions is nil", serverCfg: clients.ServerConfig{ServerURI: "server-address"}, - wantErr: true, }, { name: "Extensions is not a ServerConfigExtension", @@ -145,7 +166,6 @@ func (s) TestBuild(t *testing.T) { ServerURI: "server-address", Extensions: 1, }, - wantErr: true, }, { name: "ServerConfigExtension Credentials is nil", @@ -153,39 +173,24 @@ func (s) TestBuild(t *testing.T) { ServerURI: "server-address", Extensions: ServerConfigExtension{}, }, - wantErr: true, - }, - { - name: "success", - serverCfg: clients.ServerConfig{ - ServerURI: "server-address", - Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, - }, - wantErr: false, }, } for _, test := range tests { t.Run(test.name, func(t *testing.T) { b := &Builder{} tr, err := b.Build(test.serverCfg) - if (err != nil) != test.wantErr { - t.Fatalf("Build() error = %v, wantErr %v", err, test.wantErr) + if err == nil { + t.Fatalf("Build() succeeded, want error") } if tr != nil { - defer tr.Close() - } - if !test.wantErr && tr == nil { - t.Fatalf("got non-nil transport from Build(), want nil") - } - if test.wantErr && tr != nil { - t.Fatalf("got nil transport from Build(), want non-nil") + t.Fatalf("Got non-nil transport from Build(), want nil") } }) } } -// TestNewStream_Success verifies that grpcTransport.NewStream() successfully -// creates a new client stream for the server when provided a valid server URI. +// TestNewStream_Success verifies that NewStream() successfully creates a new +// client stream for the server when provided a valid server URI. func (s) TestNewStream_Success(t *testing.T) { ts := setupTestServer(t, &v3discoverypb.DiscoveryResponse{VersionInfo: "1"}) @@ -208,7 +213,7 @@ func (s) TestNewStream_Success(t *testing.T) { } } -// TestNewStream_Error verifies that grpcTransport.NewStream() returns an error +// TestNewStream_Error verifies that NewStream() returns an error // when attempting to create a stream with an invalid server URI. func (s) TestNewStream_Error(t *testing.T) { serverCfg := clients.ServerConfig{ @@ -230,13 +235,12 @@ func (s) TestNewStream_Error(t *testing.T) { } } -// TestStream_SendAndRecv verifies that grpcTransport.Stream.Send() -// and grpcTransport.Stream.Recv() successfully send and receive messages -// on the stream to and from the gRPC server. +// TestStream_SendAndRecv verifies that Send() and Recv() successfully send +// and receive messages on the stream to and from the gRPC server. // // It starts a gRPC test server using setupTestServer(). The test then sends a // testDiscoverRequest on the stream and verifies that the received discovery -// on the server is same as sent. It then wait to receive a +// request on the server is same as sent. It then wait to receive a // testDiscoverResponse from the server and verifies that the received // discovery response is same as sent from the server. func (s) TestStream_SendAndRecv(t *testing.T) { @@ -277,11 +281,11 @@ func (s) TestStream_SendAndRecv(t *testing.T) { // sent. select { case gotReq := <-ts.requestChan: - if !cmp.Equal(gotReq, testDiscoverRequest, protocmp.Transform()) { - t.Fatalf("<-ts.requestChan = %v, want %v", &gotReq, &testDiscoverRequest) + if diff := cmp.Diff(gotReq, testDiscoverRequest, protocmp.Transform()); diff != "" { + t.Fatalf("Unexpected diff in request received on server (-want +got):\n%s", diff) } case <-ctx.Done(): - t.Fatalf("timeout waiting for request to reach server") + t.Fatalf("Timeout waiting for request to reach server") } // Wait until response message is received from the server. @@ -296,7 +300,7 @@ func (s) TestStream_SendAndRecv(t *testing.T) { if err := proto.Unmarshal(res, &gotRes); err != nil { t.Fatalf("Failed to unmarshal response from ts.requestChan to DiscoveryResponse: %v", err) } - if !cmp.Equal(&gotRes, ts.response, protocmp.Transform()) { - t.Fatalf("proto.Unmarshal(res, &gotRes) = %v, want %v", &gotRes, ts.response) + if diff := cmp.Diff(&gotRes, ts.response, protocmp.Transform()); diff != "" { + t.Fatalf("proto.Unmarshal(res, &gotRes) returned unexpected diff (-want +got):\n%s", diff) } } From f7be854d47a3690636fadc408b05ce045ad8c877 Mon Sep 17 00:00:00 2001 From: Purnesh Dixit Date: Wed, 26 Feb 2025 10:27:44 +0530 Subject: [PATCH 09/10] doug review 3 --- .../clients/grpctransport/grpc_transport.go | 20 +++++++++---------- .../grpctransport/grpc_transport_test.go | 19 ++++++++++++------ 2 files changed, 23 insertions(+), 16 deletions(-) diff --git a/xds/internal/clients/grpctransport/grpc_transport.go b/xds/internal/clients/grpctransport/grpc_transport.go index be1f75e747cb..c56252c172af 100644 --- a/xds/internal/clients/grpctransport/grpc_transport.go +++ b/xds/internal/clients/grpctransport/grpc_transport.go @@ -40,7 +40,7 @@ type ServerConfigExtension struct { } // Builder creates gRPC-based Transports. It must be paired with ServerConfigs -// that contain Extension field of type ServerConfigExtension. +// that contain an Extension field of type ServerConfigExtension. type Builder struct{} // Build returns a gRPC-based clients.Transport. @@ -88,13 +88,18 @@ type grpcTransport struct { // NewStream creates a new gRPC stream to the server for the specified method. func (g *grpcTransport) NewStream(ctx context.Context, method string) (clients.Stream, error) { - s, err := g.cc.NewStream(ctx, &grpc.StreamDesc{StreamName: method, ClientStreams: true, ServerStreams: true}, method) + s, err := g.cc.NewStream(ctx, &grpc.StreamDesc{ClientStreams: true, ServerStreams: true}, method) if err != nil { return nil, err } return &stream{stream: s}, nil } +// Close closes the gRPC channel to the server. +func (g *grpcTransport) Close() error { + return g.cc.Close() +} + type stream struct { stream grpc.ClientStream } @@ -109,23 +114,18 @@ func (s *stream) Recv() ([]byte, error) { var typedRes []byte if err := s.stream.RecvMsg(&typedRes); err != nil { - return typedRes, err + return nil, err } return typedRes, nil } -// Close closes the gRPC channel to the server. -func (g *grpcTransport) Close() error { - return g.cc.Close() -} - type byteCodec struct{} func (c *byteCodec) Marshal(v any) ([]byte, error) { if b, ok := v.([]byte); ok { return b, nil } - return nil, fmt.Errorf("message is %T, but must be a byte slice ", v) + return nil, fmt.Errorf("message is %T, but must be a []byte", v) } func (c *byteCodec) Unmarshal(data []byte, v any) error { @@ -133,7 +133,7 @@ func (c *byteCodec) Unmarshal(data []byte, v any) error { *b = data return nil } - return fmt.Errorf("target is %T, but must be a pointer to a byte slice", v) + return fmt.Errorf("target is %T, but must be *[]byte", v) } func (c *byteCodec) Name() string { diff --git a/xds/internal/clients/grpctransport/grpc_transport_test.go b/xds/internal/clients/grpctransport/grpc_transport_test.go index 4ad2f39e92d8..823421cc26a5 100644 --- a/xds/internal/clients/grpctransport/grpc_transport_test.go +++ b/xds/internal/clients/grpctransport/grpc_transport_test.go @@ -27,7 +27,9 @@ import ( "github.com/google/go-cmp/cmp" "google.golang.org/grpc" + "google.golang.org/grpc/credentials" "google.golang.org/grpc/credentials/insecure" + "google.golang.org/grpc/credentials/local" "google.golang.org/grpc/internal/grpctest" "google.golang.org/grpc/xds/internal/clients" "google.golang.org/protobuf/proto" @@ -78,11 +80,7 @@ func setupTestServer(t *testing.T, response *v3discoverypb.DiscoveryResponse) *t s := grpc.NewServer() v3discoverygrpc.RegisterAggregatedDiscoveryServiceServer(s, ts) - go func() { - if err := s.Serve(lis); err != nil { - t.Logf("Server exited with error: %v", err) - } - }() + go s.Serve(lis) t.Cleanup(s.Stop) return ts @@ -113,12 +111,21 @@ func (s *testServer) StreamAggregatedResources(stream v3discoverygrpc.Aggregated } } +type testCredentials struct { + credentials.Bundle + transportCredentials credentials.TransportCredentials +} + +func (tc *testCredentials) TransportCredentials() credentials.TransportCredentials { + return tc.transportCredentials +} + // TestBuild_Success verifies that the Builder successfully creates a new // Transport with a non-nil grpc.ClientConn. func (s) TestBuild_Success(t *testing.T) { serverCfg := clients.ServerConfig{ ServerURI: "server-address", - Extensions: ServerConfigExtension{Credentials: insecure.NewBundle()}, + Extensions: ServerConfigExtension{Credentials: &testCredentials{transportCredentials: local.NewCredentials()}}, } b := &Builder{} From b175bc0489c928bc5a01a6a1648beed9914bbfa9 Mon Sep 17 00:00:00 2001 From: Purnesh Dixit Date: Fri, 28 Feb 2025 14:03:58 +0530 Subject: [PATCH 10/10] easwar nits --- .../clients/grpctransport/grpc_transport.go | 4 ++-- .../grpctransport/grpc_transport_test.go | 21 ++++++++++++------- 2 files changed, 15 insertions(+), 10 deletions(-) diff --git a/xds/internal/clients/grpctransport/grpc_transport.go b/xds/internal/clients/grpctransport/grpc_transport.go index c56252c172af..c5c1f99694ba 100644 --- a/xds/internal/clients/grpctransport/grpc_transport.go +++ b/xds/internal/clients/grpctransport/grpc_transport.go @@ -125,7 +125,7 @@ func (c *byteCodec) Marshal(v any) ([]byte, error) { if b, ok := v.([]byte); ok { return b, nil } - return nil, fmt.Errorf("message is %T, but must be a []byte", v) + return nil, fmt.Errorf("grpctransport: message is %T, but must be a []byte", v) } func (c *byteCodec) Unmarshal(data []byte, v any) error { @@ -133,7 +133,7 @@ func (c *byteCodec) Unmarshal(data []byte, v any) error { *b = data return nil } - return fmt.Errorf("target is %T, but must be *[]byte", v) + return fmt.Errorf("grpctransport: target is %T, but must be *[]byte", v) } func (c *byteCodec) Name() string { diff --git a/xds/internal/clients/grpctransport/grpc_transport_test.go b/xds/internal/clients/grpctransport/grpc_transport_test.go index 823421cc26a5..0ab48707c05e 100644 --- a/xds/internal/clients/grpctransport/grpc_transport_test.go +++ b/xds/internal/clients/grpctransport/grpc_transport_test.go @@ -93,6 +93,8 @@ func setupTestServer(t *testing.T, response *v3discoverypb.DiscoveryResponse) *t // verify if the correct request was received. It continues until the client // closes the stream. func (s *testServer) StreamAggregatedResources(stream v3discoverygrpc.AggregatedDiscoveryService_StreamAggregatedResourcesServer) error { + ctx := stream.Context() + for { // Receive a DiscoveryRequest from the client req, err := stream.Recv() @@ -102,7 +104,12 @@ func (s *testServer) StreamAggregatedResources(stream v3discoverygrpc.Aggregated if err != nil { return err // Handle other errors } - s.requestChan <- req + + select { + case s.requestChan <- req: + case <-ctx.Done(): + return ctx.Err() + } // Send the response back to the client if err := stream.Send(s.response); err != nil { @@ -214,8 +221,7 @@ func (s) TestNewStream_Success(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) defer cancel() - _, err = transport.NewStream(ctx, "/envoy.service.discovery.v3.AggregatedDiscoveryService/StreamAggregatedResources") - if err != nil { + if _, err = transport.NewStream(ctx, "/envoy.service.discovery.v3.AggregatedDiscoveryService/StreamAggregatedResources"); err != nil { t.Fatalf("transport.NewStream() failed: %v", err) } } @@ -236,8 +242,7 @@ func (s) TestNewStream_Error(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout) defer cancel() - _, err = transport.NewStream(ctx, "/envoy.service.discovery.v3.AggregatedDiscoveryService/StreamAggregatedResources") - if err == nil { + if _, err = transport.NewStream(ctx, "/envoy.service.discovery.v3.AggregatedDiscoveryService/StreamAggregatedResources"); err == nil { t.Fatal("transport.NewStream() succeeded, want failure") } } @@ -288,7 +293,7 @@ func (s) TestStream_SendAndRecv(t *testing.T) { // sent. select { case gotReq := <-ts.requestChan: - if diff := cmp.Diff(gotReq, testDiscoverRequest, protocmp.Transform()); diff != "" { + if diff := cmp.Diff(testDiscoverRequest, gotReq, protocmp.Transform()); diff != "" { t.Fatalf("Unexpected diff in request received on server (-want +got):\n%s", diff) } case <-ctx.Done(): @@ -305,9 +310,9 @@ func (s) TestStream_SendAndRecv(t *testing.T) { // server. var gotRes v3discoverypb.DiscoveryResponse if err := proto.Unmarshal(res, &gotRes); err != nil { - t.Fatalf("Failed to unmarshal response from ts.requestChan to DiscoveryResponse: %v", err) + t.Fatalf("Failed to unmarshal response from server to DiscoveryResponse: %v", err) } - if diff := cmp.Diff(&gotRes, ts.response, protocmp.Transform()); diff != "" { + if diff := cmp.Diff(ts.response, &gotRes, protocmp.Transform()); diff != "" { t.Fatalf("proto.Unmarshal(res, &gotRes) returned unexpected diff (-want +got):\n%s", diff) } }