| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
There was a problem hiding this comment.
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.
Sorry, something went wrong.
|
/gemini review |
Sorry, something went wrong.
There was a problem hiding this comment.
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.
Sorry, something went wrong.
…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.
|
/gemini review |
Sorry, something went wrong.
There was a problem hiding this comment.
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.
Sorry, something went wrong.
| Back | FazBrowse Home | New Git URL |
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:
Solution