Skip to content
This repository was archived by the owner on Apr 1, 2026. It is now read-only.

Commit 2683b50

Browse files
committed
added aclose test to read_rows_stream
1 parent 0f4ee8d commit 2683b50

2 files changed

Lines changed: 37 additions & 3 deletions

File tree

google/cloud/bigtable/data/_async/_read_rows.py

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@
1919

2020
import time
2121

22+
from grpc import StatusCode
23+
2224
from google.cloud.bigtable_v2.types import ReadRowsRequest as ReadRowsRequestPB
2325
from google.cloud.bigtable_v2.types import ReadRowsResponse as ReadRowsResponsePB
2426
from google.cloud.bigtable_v2.types import RowSet as RowSetPB
@@ -339,7 +341,12 @@ async def merge_rows(
339341
continue
340342
except CrossSync.StopIteration:
341343
raise InvalidChunk("premature end of stream")
344+
except GeneratorExit as close_exception:
345+
# handle aclose()
346+
self._operation_metric.end_with_status(StatusCode.CANCELLED)
347+
raise close_exception
342348
except Exception as generic_exception:
349+
# handle exceptions in retry wrapper
343350
raise generic_exception
344351
else:
345352
self._operation_metric.end_with_success()

tests/system/data/test_metrics_async.py

Lines changed: 30 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -397,13 +397,40 @@ async def test_read_rows_stream(self, table, temp_rows, handler, cluster_config)
397397
assert attempt.grpc_throttling_time_ns == 0 # TODO: confirm
398398

399399
@CrossSync.pytest
400-
async def test_read_rows_stream_failure_grpc(
400+
async def test_read_rows_stream_failure_closed(
401401
self, table, temp_rows, handler, error_injector
402402
):
403403
"""
404-
Test failure in grpc layer by injecting an error into an interceptor
404+
Test how metrics collection handles closed generator
405+
"""
406+
await temp_rows.add_row(b"row_key_1")
407+
await temp_rows.add_row(b"row_key_2")
408+
handler.clear()
409+
generator = await table.read_rows_stream(
410+
ReadRowsQuery()
411+
)
412+
await generator.__anext__()
413+
await generator.aclose()
414+
with pytest.raises(CrossSync.StopIteration):
415+
await generator.__anext__()
416+
# validate counts
417+
assert len(handler.completed_operations) == 1
418+
assert len(handler.completed_attempts) == 1
419+
assert len(handler.cancelled_operations) == 0
420+
# validate operation
421+
operation = handler.completed_operations[0]
422+
assert operation.final_status.name == "CANCELLED"
423+
assert operation.op_type.value == "ReadRows"
424+
assert operation.is_streaming is True
425+
assert len(operation.completed_attempts) == 1
426+
assert operation.cluster_id == "unspecified"
427+
assert operation.zone == "global"
428+
# validate attempt
429+
attempt = handler.completed_attempts[0]
430+
assert attempt.end_status.name == "CANCELLED"
431+
assert attempt.gfe_latency_ns is None
432+
405433

406-
No headers expected
407434
"""
408435
await temp_rows.add_row(b"row_key_1")
409436
handler.clear()

0 commit comments

Comments
 (0)