FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

fix(spanner): correctly wrap async streaming calls in metrics interceptor by olavloite · Pull Request #18581 · googleapis/google-cloud-python · GitHub

Repository navigation

fix(spanner): correctly wrap async streaming calls in metrics interceptor - #18581

Merged
olavloite merged 1 commit into
mainfrom
spanner-correctly-wrap-streaming-calls
Oct 6, 2026
Merged

olavloite merged 1 commit into
mainfrom
spanner-correctly-wrap-streaming-calls

Conversation

Copy link
Copy Markdown
Contributor

Fix an issue where async streaming RPCs were incorrectly wrapped in _AsyncUnaryResponseWrapper instead of _AsyncStreamingResponseWrapper, preventing metrics from being properly recorded.

Problem

In grpc.aio, streaming call objects (like UnaryStreamCall) implement __aiter__ (returning an async generator) but do not define __anext__ directly on the call object itself. Because AsyncMetricsInterceptor checked hasattr(response, "__anext__"), all async streaming calls (such as ExecuteStreamingSql and StreamingRead) were misclassified as unary calls and wrapped in _AsyncUnaryResponseWrapper.

Because _AsyncUnaryResponseWrapper only records metrics when awaited (__await__), and callers iterate over streams chunk-by-chunk rather than awaiting the stream object itself:

  1. Frontend metrics (gfe_latencies, afe_latencies, and missing header counts) were never recorded for any async queries or streaming reads.
  2. Attempt completion was never recorded when the stream finished. Instead, it was only recorded much later when the wrapper was garbage-collected (__del__), resulting in wildly inflated attempt latency measurements or attempts incorrectly reporting as CANCELLED.

Solution

  • Explicitly dispatch async interceptor wrappers based on the gRPC call type (is_streaming=True for intercept_unary_stream and intercept_stream_stream) rather than fragile runtime dunder checks.
  • Realign class inheritance so StreamUnaryCall is handled by _AsyncUnaryResponseWrapper, adding write() and done_writing() delegation.
  • Support sending real HTTP/2 response headers in the in-memory mock Spanner server via an add_header() helper, removing monkeypatching of internal gRPC classes in integration tests.
  • Add unit tests for direct __anext__() consumption (matching StreamedResultSet), delegation methods, error handling, and sync/async response wrappers.

olavloite requested a review from a team as a code owner October 6, 2026 14:09

gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

Code Review

This pull request introduces async metrics collection for Cloud Spanner operations by implementing AsyncMetricsInterceptor and wrapping async unary and streaming responses to defer metrics recording. It also enhances the mock Spanner server to support mock headers and initial metadata, and adds comprehensive integration and unit tests. Feedback on these changes suggests avoiding fragile runtime stack-frame inspection in the mock server by explicitly passing the method name to send_initial_metadata, and properly closing gRPC channels in the async integration tests using async with blocks to prevent resource leaks.

Copy link
Copy Markdown
Contributor Author

/gemini review

gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

Code Review

This pull request updates the AsyncMetricsInterceptor to explicitly track streaming RPCs using an is_streaming flag, adjusts the base classes of the response wrappers, and adds comprehensive unit and integration tests for async frontend metrics. The reviewer feedback highlights two important issues: first, _AsyncStreamingResponseWrapper needs explicit implementations of write and done_writing to properly delegate calls because they are defined as abstract methods in its base class; second, the unit tests for streaming response wrapper delegation should be updated to assert that these methods are actually called on the underlying mock response.

…ptor

Fix an issue where async streaming RPCs were incorrectly wrapped in
_AsyncUnaryResponseWrapper instead of _AsyncStreamingResponseWrapper,
preventing metrics from being properly recorded.

In grpc.aio, streaming call objects (like UnaryStreamCall) implement
__aiter__ (returning an async generator) but do not define __anext__
directly on the call object itself. Because AsyncMetricsInterceptor
checked hasattr(response, "__anext__"), all async streaming calls
(such as ExecuteStreamingSql and StreamingRead) were misclassified as
unary calls and wrapped in _AsyncUnaryResponseWrapper.

Because _AsyncUnaryResponseWrapper only records metrics when awaited
(__await__), and callers iterate over streams chunk-by-chunk rather than
awaiting the stream object itself:
1. Frontend metrics (gfe_latencies, afe_latencies, and missing header
   counts) were never recorded for any async queries or streaming reads.
2. Attempt completion was never recorded when the stream finished. Instead,
   it was only recorded much later when the wrapper was garbage-collected
   (__del__), resulting in wildly inflated attempt latency measurements
   or attempts incorrectly reporting as CANCELLED.

- Explicitly dispatch async interceptor wrappers based on the gRPC call
  type (is_streaming=True for intercept_unary_stream and
  intercept_stream_stream) rather than fragile runtime dunder checks.
- Realign class inheritance so StreamUnaryCall is handled by
  _AsyncUnaryResponseWrapper, adding write() and done_writing() delegation.
- Support sending real HTTP/2 response headers in the in-memory mock Spanner
  server via an add_header() helper, removing monkeypatching of internal
  gRPC classes in integration tests.
- Add unit tests for direct __anext__() consumption (matching
  StreamedResultSet), delegation methods, error handling, and sync/async
  response wrappers.
olavloite force-pushed the spanner-correctly-wrap-streaming-calls branch from 33b5000 to b43876b Compare October 6, 2026 14:33

Copy link
Copy Markdown
Contributor Author

/gemini review

gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Choose a reason Spam Abuse Off Topic Outdated Duplicate Resolved Low Quality

Code Review

This pull request implements and tests async metrics interception for Cloud Spanner operations. It updates AsyncMetricsInterceptor to correctly identify streaming vs. unary calls using an explicit is_streaming flag instead of inspecting __anext__. It also updates response wrappers to inherit from the correct gRPC call classes and delegate necessary methods like write and done_writing. Additionally, mock Spanner and test suites are enhanced to support sending initial metadata (such as server-timing headers) and verify that frontend metrics are correctly recorded and exported in both synchronous and asynchronous contexts. There are no review comments, so we have no feedback to provide.

olavloite merged commit 2d19b0a into main Oct 6, 2026
54 checks passed
olavloite deleted the spanner-correctly-wrap-streaming-calls branch October 6, 2026 16:30
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants


Back | FazBrowse Home | New Git URL