Skip to content

Commit f80658f

Browse files
mbissaPranjali-2501
authored andcommitted
client: process RPC stats/tracing only when a handler is configured (grpc#8874)
This resolves an issue where below log statement keeps printing repeatedly: "ERROR: [otel-plugin] ctx passed into client side stats handler metrics event handling has no client attempt data present". When new stream is created for health/orca producers, stats and tracing is not setup. However, this fact is ignored during RPC and an error logs is printed to denote that stats cannot be handled. We will enable stream to have its own reference to the stats handler and only process per RPC implementation when it is present (like in case of regular data streams). Internal issue: b/385685802 RELEASE NOTES: * stats: only process RPC stats/tracing in health and ORCA producers if a handler is configured, preventing unnecessary error logging
1 parent 3ee3896 commit f80658f

7 files changed

Lines changed: 111 additions & 59 deletions

File tree

internal/transport/client_stream.go

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ import (
2424
"golang.org/x/net/http2"
2525
"google.golang.org/grpc/mem"
2626
"google.golang.org/grpc/metadata"
27+
"google.golang.org/grpc/stats"
2728
"google.golang.org/grpc/status"
2829
)
2930

@@ -46,10 +47,11 @@ type ClientStream struct {
4647
// meaningful after headerChan is closed (always call waitOnHeader() before
4748
// reading its value).
4849
headerValid bool
49-
noHeaders bool // set if the client never received headers (set only after the stream is done).
50-
headerChanClosed uint32 // set when headerChan is closed. Used to avoid closing headerChan multiple times.
51-
bytesReceived atomic.Bool // indicates whether any bytes have been received on this stream
52-
unprocessed atomic.Bool // set if the server sends a refused stream or GOAWAY including this stream
50+
noHeaders bool // set if the client never received headers (set only after the stream is done).
51+
headerChanClosed uint32 // set when headerChan is closed. Used to avoid closing headerChan multiple times.
52+
bytesReceived atomic.Bool // indicates whether any bytes have been received on this stream
53+
unprocessed atomic.Bool // set if the server sends a refused stream or GOAWAY including this stream
54+
statsHandler stats.Handler // nil for internal streams (e.g., health check, ORCA) where telemetry is not supported.
5355
}
5456

5557
// Read reads an n byte message from the input stream.

internal/transport/http2_client.go

Lines changed: 13 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -478,18 +478,19 @@ func NewHTTP2Client(connectCtx, ctx context.Context, addr resolver.Address, opts
478478
return t, nil
479479
}
480480

481-
func (t *http2Client) newStream(ctx context.Context, callHdr *CallHdr) *ClientStream {
481+
func (t *http2Client) newStream(ctx context.Context, callHdr *CallHdr, handler stats.Handler) *ClientStream {
482482
// TODO(zhaoq): Handle uint32 overflow of Stream.id.
483483
s := &ClientStream{
484484
Stream: Stream{
485485
method: callHdr.Method,
486486
sendCompress: callHdr.SendCompress,
487487
contentSubtype: callHdr.ContentSubtype,
488488
},
489-
ct: t,
490-
done: make(chan struct{}),
491-
headerChan: make(chan struct{}),
492-
doneFunc: callHdr.DoneFunc,
489+
ct: t,
490+
done: make(chan struct{}),
491+
headerChan: make(chan struct{}),
492+
doneFunc: callHdr.DoneFunc,
493+
statsHandler: handler,
493494
}
494495
s.Stream.buf.init()
495496
s.Stream.wq.init(defaultWriteQuota, s.done)
@@ -744,7 +745,7 @@ func (e NewStreamError) Error() string {
744745

745746
// NewStream creates a stream and registers it into the transport as "active"
746747
// streams. All non-nil errors returned will be *NewStreamError.
747-
func (t *http2Client) NewStream(ctx context.Context, callHdr *CallHdr) (*ClientStream, error) {
748+
func (t *http2Client) NewStream(ctx context.Context, callHdr *CallHdr, handler stats.Handler) (*ClientStream, error) {
748749
ctx = peer.NewContext(ctx, t.Peer())
749750

750751
// ServerName field of the resolver returned address takes precedence over
@@ -781,7 +782,7 @@ func (t *http2Client) NewStream(ctx context.Context, callHdr *CallHdr) (*ClientS
781782
if err != nil {
782783
return nil, &NewStreamError{Err: err, AllowTransparentRetry: false}
783784
}
784-
s := t.newStream(ctx, callHdr)
785+
s := t.newStream(ctx, callHdr, handler)
785786
cleanup := func(err error) {
786787
if s.swapState(streamDone) == streamDone {
787788
// If it was already done, return.
@@ -906,7 +907,7 @@ func (t *http2Client) NewStream(ctx context.Context, callHdr *CallHdr) (*ClientS
906907
return nil, &NewStreamError{Err: ErrConnClosing, AllowTransparentRetry: true}
907908
}
908909
}
909-
if t.statsHandler != nil {
910+
if s.statsHandler != nil {
910911
header, ok := metadata.FromOutgoingContext(ctx)
911912
if ok {
912913
header.Set("user-agent", t.userAgent)
@@ -915,7 +916,7 @@ func (t *http2Client) NewStream(ctx context.Context, callHdr *CallHdr) (*ClientS
915916
}
916917
// Note: The header fields are compressed with hpack after this call returns.
917918
// No WireLength field is set here.
918-
t.statsHandler.HandleRPC(s.ctx, &stats.OutHeader{
919+
s.statsHandler.HandleRPC(s.ctx, &stats.OutHeader{
919920
Client: true,
920921
FullMethod: callHdr.Method,
921922
RemoteAddr: t.remoteAddr,
@@ -1591,16 +1592,16 @@ func (t *http2Client) operateHeaders(frame *http2.MetaHeadersFrame) {
15911592
}
15921593
}
15931594

1594-
if t.statsHandler != nil {
1595+
if s.statsHandler != nil {
15951596
if !endStream {
1596-
t.statsHandler.HandleRPC(s.ctx, &stats.InHeader{
1597+
s.statsHandler.HandleRPC(s.ctx, &stats.InHeader{
15971598
Client: true,
15981599
WireLength: int(frame.Header().Length),
15991600
Header: metadata.MD(mdata).Copy(),
16001601
Compression: s.recvCompress,
16011602
})
16021603
} else {
1603-
t.statsHandler.HandleRPC(s.ctx, &stats.InTrailer{
1604+
s.statsHandler.HandleRPC(s.ctx, &stats.InTrailer{
16041605
Client: true,
16051606
WireLength: int(frame.Header().Length),
16061607
Trailer: metadata.MD(mdata).Copy(),

internal/transport/keepalive_test.go

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -69,7 +69,7 @@ func (s) TestMaxConnectionIdle(t *testing.T) {
6969

7070
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
7171
defer cancel()
72-
stream, err := client.NewStream(ctx, &CallHdr{})
72+
stream, err := client.NewStream(ctx, &CallHdr{}, nil)
7373
if err != nil {
7474
t.Fatalf("client.NewStream() failed: %v", err)
7575
}
@@ -111,7 +111,7 @@ func (s) TestMaxConnectionIdleBusyClient(t *testing.T) {
111111

112112
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
113113
defer cancel()
114-
_, err := client.NewStream(ctx, &CallHdr{})
114+
_, err := client.NewStream(ctx, &CallHdr{}, nil)
115115
if err != nil {
116116
t.Fatalf("client.NewStream() failed: %v", err)
117117
}
@@ -150,7 +150,7 @@ func (s) TestMaxConnectionAge(t *testing.T) {
150150

151151
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
152152
defer cancel()
153-
if _, err := client.NewStream(ctx, &CallHdr{}); err != nil {
153+
if _, err := client.NewStream(ctx, &CallHdr{}, nil); err != nil {
154154
t.Fatalf("client.NewStream() failed: %v", err)
155155
}
156156

@@ -373,7 +373,7 @@ func (s) TestKeepaliveClientClosesWithActiveStreams(t *testing.T) {
373373
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
374374
defer cancel()
375375
// Create a stream, but send no data on it.
376-
if _, err := client.NewStream(ctx, &CallHdr{}); err != nil {
376+
if _, err := client.NewStream(ctx, &CallHdr{}, nil); err != nil {
377377
t.Fatalf("Stream creation failed: %v", err)
378378
}
379379

@@ -519,7 +519,7 @@ func (s) TestKeepaliveServerEnforcementWithAbusiveClientWithRPC(t *testing.T) {
519519

520520
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
521521
defer cancel()
522-
if _, err := client.NewStream(ctx, &CallHdr{}); err != nil {
522+
if _, err := client.NewStream(ctx, &CallHdr{}, nil); err != nil {
523523
t.Fatalf("Stream creation failed: %v", err)
524524
}
525525

@@ -748,7 +748,7 @@ func (s) TestTCPUserTimeout(t *testing.T) {
748748

749749
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
750750
defer cancel()
751-
stream, err := client.NewStream(ctx, &CallHdr{})
751+
stream, err := client.NewStream(ctx, &CallHdr{}, nil)
752752
if err != nil {
753753
t.Fatalf("client.NewStream() failed: %v", err)
754754
}
@@ -815,7 +815,7 @@ func makeTLSCreds(t *testing.T, certPath, keyPath, rootsPath string) credentials
815815
func checkForHealthyStream(client *http2Client) error {
816816
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
817817
defer cancel()
818-
stream, err := client.NewStream(ctx, &CallHdr{})
818+
stream, err := client.NewStream(ctx, &CallHdr{}, nil)
819819
stream.Close(err)
820820
return err
821821
}
@@ -824,7 +824,7 @@ func pollForStreamCreationError(client *http2Client) error {
824824
ctx, cancel := context.WithTimeout(context.Background(), defaultTestTimeout)
825825
defer cancel()
826826
for {
827-
if _, err := client.NewStream(ctx, &CallHdr{}); err != nil {
827+
if _, err := client.NewStream(ctx, &CallHdr{}, nil); err != nil {
828828
break
829829
}
830830
time.Sleep(50 * time.Millisecond)
@@ -850,7 +850,7 @@ func waitForGoAwayTooManyPings(client *http2Client) error {
850850
return fmt.Errorf("test timed out before getting GoAway with reason:GoAwayTooManyPings from server")
851851
}
852852

853-
if _, err := client.NewStream(ctx, &CallHdr{}); err == nil {
853+
if _, err := client.NewStream(ctx, &CallHdr{}, nil); err == nil {
854854
return fmt.Errorf("stream creation succeeded after receiving a GoAway from the server")
855855
}
856856
return nil

internal/transport/transport.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -617,7 +617,7 @@ type ClientTransport interface {
617617
GracefulClose()
618618

619619
// NewStream creates a Stream for an RPC.
620-
NewStream(ctx context.Context, callHdr *CallHdr) (*ClientStream, error)
620+
NewStream(ctx context.Context, callHdr *CallHdr, handler stats.Handler) (*ClientStream, error)
621621

622622
// Error returns a channel that is closed when some I/O error
623623
// happens. Typically the caller should have a goroutine to monitor

0 commit comments

Comments
 (0)