From 98b4b16820ae6824ce80446f1c06966428b35eb4 Mon Sep 17 00:00:00 2001 From: raj-prince Date: Tue, 21 Jul 2026 03:01:33 +0000 Subject: [PATCH 1/8] feat: adding read/write support for gRPC stall --- README.md | 2 +- gcs/upload.py | 36 +++++++ testbench/common.py | 5 +- testbench/grpc_server.py | 51 +++++++++ tests/test_testbench_retry.py | 193 ++++++++++++++++++++++++++++++++++ 5 files changed, 285 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index d2dc8a91..7494f0cd 100644 --- a/README.md +++ b/README.md @@ -274,7 +274,7 @@ curl -H "x-retry-test-id: 1d05c20627844214a9ff7cbcf696317d" "http://localhost:91 | return-broken-stream | [HTTP] Testbench will fail after a few downloaded bytes
[GRPC] Testbench will fail with `UNAVAILABLE` after a few downloaded bytes | return-broken-stream-after-YK | [HTTP] Testbench will fail after YKiB of downloaded data
[GRPC] Testbench will fail with `UNAVAILABLE` after YKiB of downloaded data | return-reset-connection | [HTTP] Testbench will fail with a reset connection
[GRPC] Testbench will fail the RPC with `UNAVAILABLE` -| stall-for-Ts-after-YK | [HTTP] Testbench will stall for T second after reading YKiB of downloaded/uploaded data, e.g. stall-for-10s-after-12K stalls after reading/writing 12KiB of data
[GRPC] Not supported +| stall-for-Ts-after-YK | [HTTP] Testbench will stall for T second after reading YKiB of downloaded/uploaded data, e.g. stall-for-10s-after-12K stalls after reading/writing 12KiB of data
[GRPC] Supported for `storage.objects.get` and `storage.objects.insert` | redirect-send-token-T | [HTTP] Unsupported [GRPC] Testbench will fail the RPC with `ABORTED` and include appropriate redirection error details. | redirect-send-handle-and-token-T | [HTTP] Unsupported [GRPC] Testbench will fail the RPC with `ABORTED` and include appropriate redirection error details. | return-X-if-dp-enforced | [HTTP] Unsupported [GRPC] Testbench will fail with the equivalent gRPC error to the HTTP code provided for X, but only if DirectPath is enforced. diff --git a/gcs/upload.py b/gcs/upload.py index f7a44541..573b4497 100644 --- a/gcs/upload.py +++ b/gcs/upload.py @@ -269,6 +269,24 @@ def init_write_object_grpc(cls, db, request_iterator, context): test_id=test_id, ) + # Handle retry test stall-for-Xs-after-YK instructions if applicable. + ( + stall_time, + after_bytes, + test_id, + ) = testbench.common.get_stall_uploads_after_bytes( + db, request, context=context, transport="GRPC" + ) + if stall_time: + testbench.common.handle_stall_uploads_after_bytes( + upload, + content, + db, + stall_time, + after_bytes, + test_id=test_id, + ) + upload.media += content if request.finish_write: upload.complete = True @@ -582,6 +600,24 @@ def update_upload_checksums(upload_metadata, object_checksums): test_id=test_id, ) + # Handle retry test stall-for-Xs-after-YK instructions if applicable. + ( + stall_time, + after_bytes, + test_id, + ) = testbench.common.get_stall_uploads_after_bytes( + db, request, context=context, transport="GRPC" + ) + if stall_time: + testbench.common.handle_stall_uploads_after_bytes( + upload, + content, + db, + stall_time, + after_bytes, + test_id=test_id, + ) + # Currently, the testbench will always checkpoint and flush data for testing purposes, # instead of the 15 seconds interval used in the GCS server. # TODO(#592): Refactor testbench checkpointing to more closely follow GCS server behavior. diff --git a/testbench/common.py b/testbench/common.py index 97378ddd..580afef7 100644 --- a/testbench/common.py +++ b/testbench/common.py @@ -929,7 +929,10 @@ def wrapper(*args, **kwargs): def get_stall_uploads_after_bytes(database, request, context=None, transport="HTTP"): """Retrieve stall time and #bytes corresponding to uploads from retry test instructions.""" method = "storage.objects.insert" - test_id = request.headers.get("x-retry-test-id", None) + if context is not None: + test_id = get_retry_test_id_from_context(context) + else: + test_id = request.headers.get("x-retry-test-id", None) if not test_id: return 0, 0, "" next_instruction = None diff --git a/testbench/grpc_server.py b/testbench/grpc_server.py index 4f8bbdcf..c4875e4f 100644 --- a/testbench/grpc_server.py +++ b/testbench/grpc_server.py @@ -20,6 +20,7 @@ import json import re import sys +import time import types import uuid from collections.abc import Iterable @@ -607,6 +608,9 @@ def ReadObject(self, request, context): # Check retry test broken-stream instructions. test_id = testbench.common.get_retry_test_id_from_context(context) broken_stream_after_bytes = 0 + stall_time = 0 + stall_after_bytes = 0 + stall_applied = False method = "storage.objects.get" if test_id and self.db.has_instructions_retry_test( test_id, method, transport="GRPC" @@ -615,12 +619,34 @@ def ReadObject(self, request, context): broken_stream_after_bytes = testbench.common.get_broken_stream_after_bytes( next_instruction ) + retry_stall_after_bytes_matches = ( + testbench.common.retry_stall_after_bytes.match(next_instruction) + ) + if retry_stall_after_bytes_matches: + items = list(retry_stall_after_bytes_matches.groups()) + stall_time = int(items[0]) + stall_after_bytes = int(items[1]) * 1024 + bytes_yielded = 0 while start <= read_end: end = min(start + size, read_end) + chunk_len = end - start + + # Apply stall once when the configured byte threshold is reached. + if ( + stall_time > 0 + and not stall_applied + and bytes_yielded < stall_after_bytes + and (bytes_yielded + chunk_len) >= stall_after_bytes + ): + self.db.dequeue_next_instruction(test_id, method) + time.sleep(stall_time) + stall_applied = True + # Handle retry test broken-stream failures if applicable. if broken_stream_after_bytes and end >= broken_stream_after_bytes: chunk = blob.media[start:broken_stream_after_bytes] + bytes_yielded += len(chunk) yield storage_pb2.ReadObjectResponse( checksummed_data={ "content": chunk, @@ -636,6 +662,7 @@ def ReadObject(self, request, context): "Injected 'broken stream' fault", ) chunk = blob.media[start:end] + bytes_yielded += len(chunk) yield storage_pb2.ReadObjectResponse( checksummed_data={ "content": chunk, @@ -670,6 +697,9 @@ def BidiReadObject(self, request_iterator, context): # Check retry test broken-stream instructions. test_id = testbench.common.get_retry_test_id_from_context(context) broken_stream_after_bytes = 0 + stall_time = 0 + stall_after_bytes = 0 + stall_applied = False method = "storage.objects.get" if test_id and self.db.has_instructions_retry_test( test_id, method, transport="GRPC" @@ -678,6 +708,13 @@ def BidiReadObject(self, request_iterator, context): broken_stream_after_bytes = testbench.common.get_broken_stream_after_bytes( next_instruction ) + retry_stall_after_bytes_matches = ( + testbench.common.retry_stall_after_bytes.match(next_instruction) + ) + if retry_stall_after_bytes_matches: + items = list(retry_stall_after_bytes_matches.groups()) + stall_time = int(items[0]) + stall_after_bytes = int(items[1]) * 1024 return_redirect_token = ( testbench.common.get_return_read_handle_and_redirect_token(self.db, context) ) @@ -763,6 +800,7 @@ def terminate_w_timer(): returnable = ( broken_stream_after_bytes if broken_stream_after_bytes else sys.maxsize ) + bytes_yielded = 0 # We don't want to have a thread blocking on request_iterator, so we # have to handle results in batches rather than concurrently. This @@ -789,6 +827,19 @@ def read_results(): chunk = chunk[:returnable] range_end = False read_range["read_length"] -= excess + + # Apply stall once when the configured byte threshold is reached. + if ( + stall_time > 0 + and not stall_applied + and bytes_yielded < stall_after_bytes + and (bytes_yielded + len(chunk)) >= stall_after_bytes + ): + self.db.dequeue_next_instruction(test_id, method) + time.sleep(stall_time) + stall_applied = True + + bytes_yielded += len(chunk) returnable -= count yield response( storage_pb2.BidiReadObjectResponse( diff --git a/tests/test_testbench_retry.py b/tests/test_testbench_retry.py index 6ea3c19e..db25a205 100644 --- a/tests/test_testbench_retry.py +++ b/tests/test_testbench_retry.py @@ -1062,6 +1062,199 @@ def test_grpc_retry_reset_connection(self): "Injected 'socket closed, connection reset by peer' fault", ) + def test_grpc_retry_stall_read_after_bytes(self): + media = self._create_block(2 * UPLOAD_QUANTUM) + response = self.rest_client.put( + "/bucket-name/512k.txt", + content_type="text/plain", + data=media, + ) + self.assertEqual(response.status_code, 200) + + response = self.rest_client.post( + "/retry_test", + data=json.dumps( + { + "instructions": {"storage.objects.get": ["stall-for-1s-after-128K"]}, + "transport": "GRPC", + }, + ), + ) + self.assertEqual(response.status_code, 200) + create_rest = json.loads(response.data) + self.assertIn("id", create_rest) + + context = unittest.mock.Mock() + context.invocation_metadata = unittest.mock.Mock( + return_value=(("x-retry-test-id", create_rest.get("id")),) + ) + + start_time = time.perf_counter() + response = self.grpc.ReadObject( + storage_pb2.ReadObjectRequest( + bucket="projects/_/buckets/bucket-name", object="512k.txt" + ), + context, + ) + list(response) + elapsed = time.perf_counter() - start_time + self.assertGreater(elapsed, 1) + + def test_grpc_retry_stall_write_after_bytes(self): + response = self.rest_client.post( + "/retry_test", + data=json.dumps( + { + "instructions": { + "storage.objects.insert": ["stall-for-1s-after-250K"] + }, + "transport": "GRPC", + } + ), + ) + self.assertEqual(response.status_code, 200) + create_rest = json.loads(response.data) + self.assertIn("id", create_rest) + id = create_rest.get("id") + + context = unittest.mock.Mock() + context.invocation_metadata = unittest.mock.Mock( + return_value=(("x-retry-test-id", id),) + ) + start = self.grpc.StartResumableWrite( + storage_pb2.StartResumableWriteRequest( + write_object_spec=storage_pb2.WriteObjectSpec( + resource=storage_pb2.Object( + name="object-name-stall", bucket="projects/_/buckets/bucket-name" + ) + ) + ), + context=context, + ) + self.assertIsNotNone(start.upload_id) + + content = self._create_block(UPLOAD_QUANTUM).encode("utf-8") + r1 = storage_pb2.WriteObjectRequest( + upload_id=start.upload_id, + write_offset=0, + checksummed_data=storage_pb2.ChecksummedData( + content=content, crc32c=crc32c.crc32c(content) + ), + finish_write=False, + ) + start_time = time.perf_counter() + _ = self.grpc.WriteObject([r1], context) + elapsed = time.perf_counter() - start_time + self.assertGreater(elapsed, 1) + + # Instruction consumed; finishing write should be fast. + r2 = storage_pb2.WriteObjectRequest( + upload_id=start.upload_id, + write_offset=len(content), + checksummed_data=storage_pb2.ChecksummedData(content=b"", crc32c=crc32c.crc32c(b"")), + finish_write=True, + ) + start_time = time.perf_counter() + _ = self.grpc.WriteObject([r2], context) + elapsed = time.perf_counter() - start_time + self.assertLess(elapsed, 1) + + def test_grpc_retry_stall_bidiwrite_after_bytes(self): + response = self.rest_client.post( + "/retry_test", + data=json.dumps( + { + "instructions": { + "storage.objects.insert": ["stall-for-1s-after-250K"] + }, + "transport": "GRPC", + } + ), + ) + self.assertEqual(response.status_code, 200) + create_rest = json.loads(response.data) + self.assertIn("id", create_rest) + id = create_rest.get("id") + + context = unittest.mock.Mock() + context.invocation_metadata = unittest.mock.Mock( + return_value=(("x-retry-test-id", id),) + ) + start = self.grpc.StartResumableWrite( + storage_pb2.StartResumableWriteRequest( + write_object_spec=storage_pb2.WriteObjectSpec( + resource=storage_pb2.Object( + name="object-name-bidi-stall", + bucket="projects/_/buckets/bucket-name", + ) + ) + ), + context=context, + ) + self.assertIsNotNone(start.upload_id) + + content = self._create_block(UPLOAD_QUANTUM).encode("utf-8") + r1 = storage_pb2.BidiWriteObjectRequest( + upload_id=start.upload_id, + write_offset=0, + checksummed_data=storage_pb2.ChecksummedData( + content=content, crc32c=crc32c.crc32c(content) + ), + finish_write=False, + ) + + start_time = time.perf_counter() + _ = list(self.grpc.BidiWriteObject([r1], context)) + elapsed = time.perf_counter() - start_time + self.assertGreater(elapsed, 1) + + def test_grpc_bidiread_retry_stall_after_bytes(self): + media = self._create_block(5 * 1024 * 1024) + response = self.rest_client.put( + "/bucket-name/512k.txt", + content_type="text/plain", + data=media, + ) + self.assertEqual(response.status_code, 200) + + response = self.rest_client.post( + "/retry_test", + data=json.dumps( + { + "instructions": {"storage.objects.get": ["stall-for-1s-after-256K"]}, + "transport": "GRPC", + }, + ), + ) + self.assertEqual(response.status_code, 200) + create_rest = json.loads(response.data) + self.assertIn("id", create_rest) + + context = unittest.mock.Mock() + context.invocation_metadata = unittest.mock.Mock( + return_value=(("x-retry-test-id", create_rest.get("id")),) + ) + + r1 = storage_pb2.BidiReadObjectRequest( + read_object_spec=storage_pb2.BidiReadObjectSpec( + bucket="projects/_/buckets/bucket-name", + object="512k.txt", + ), + read_ranges=[ + storage_pb2.ReadRange( + read_offset=0, + read_length=1 * 1024 * 1024, + read_id=1, + ), + ], + ) + + start_time = time.perf_counter() + response = self.grpc.BidiReadObject([r1], context) + list(response) + elapsed = time.perf_counter() - start_time + self.assertGreater(elapsed, 1) + def test_grpc_retry_broken_stream(self): # Use the XML API to inject an object with some data. media = self._create_block(2 * UPLOAD_QUANTUM) From d19c63e6ae5fd3bacf9c965a8ba2651c11ffa3e1 Mon Sep 17 00:00:00 2001 From: raj-prince Date: Tue, 21 Jul 2026 03:13:11 +0000 Subject: [PATCH 2/8] fixing lint issue --- tests/test_testbench_retry.py | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/tests/test_testbench_retry.py b/tests/test_testbench_retry.py index db25a205..4d7d6959 100644 --- a/tests/test_testbench_retry.py +++ b/tests/test_testbench_retry.py @@ -1075,7 +1075,9 @@ def test_grpc_retry_stall_read_after_bytes(self): "/retry_test", data=json.dumps( { - "instructions": {"storage.objects.get": ["stall-for-1s-after-128K"]}, + "instructions": { + "storage.objects.get": ["stall-for-1s-after-128K"] + }, "transport": "GRPC", }, ), @@ -1125,7 +1127,8 @@ def test_grpc_retry_stall_write_after_bytes(self): storage_pb2.StartResumableWriteRequest( write_object_spec=storage_pb2.WriteObjectSpec( resource=storage_pb2.Object( - name="object-name-stall", bucket="projects/_/buckets/bucket-name" + name="object-name-stall", + bucket="projects/_/buckets/bucket-name", ) ) ), @@ -1151,7 +1154,9 @@ def test_grpc_retry_stall_write_after_bytes(self): r2 = storage_pb2.WriteObjectRequest( upload_id=start.upload_id, write_offset=len(content), - checksummed_data=storage_pb2.ChecksummedData(content=b"", crc32c=crc32c.crc32c(b"")), + checksummed_data=storage_pb2.ChecksummedData( + content=b"", crc32c=crc32c.crc32c(b"") + ), finish_write=True, ) start_time = time.perf_counter() @@ -1221,7 +1226,9 @@ def test_grpc_bidiread_retry_stall_after_bytes(self): "/retry_test", data=json.dumps( { - "instructions": {"storage.objects.get": ["stall-for-1s-after-256K"]}, + "instructions": { + "storage.objects.get": ["stall-for-1s-after-256K"] + }, "transport": "GRPC", }, ), From d96e1043afd65e9938db45c408bf163d71e01345 Mon Sep 17 00:00:00 2001 From: raj-prince Date: Tue, 21 Jul 2026 03:33:07 +0000 Subject: [PATCH 3/8] extend test for fast second call --- tests/test_testbench_retry.py | 30 ++++++++++++++++++++++++++++++ 1 file changed, 30 insertions(+) diff --git a/tests/test_testbench_retry.py b/tests/test_testbench_retry.py index 4d7d6959..1dd14a12 100644 --- a/tests/test_testbench_retry.py +++ b/tests/test_testbench_retry.py @@ -1102,6 +1102,17 @@ def test_grpc_retry_stall_read_after_bytes(self): elapsed = time.perf_counter() - start_time self.assertGreater(elapsed, 1) + start_time = time.perf_counter() + response = self.grpc.ReadObject( + storage_pb2.ReadObjectRequest( + bucket="projects/_/buckets/bucket-name", object="512k.txt" + ), + context, + ) + list(response) + elapsed = time.perf_counter() - start_time + self.assertLess(elapsed, 1) + def test_grpc_retry_stall_write_after_bytes(self): response = self.rest_client.post( "/retry_test", @@ -1213,6 +1224,19 @@ def test_grpc_retry_stall_bidiwrite_after_bytes(self): elapsed = time.perf_counter() - start_time self.assertGreater(elapsed, 1) + r2 = storage_pb2.BidiWriteObjectRequest( + upload_id=start.upload_id, + write_offset=len(content), + checksummed_data=storage_pb2.ChecksummedData( + content=b"", crc32c=crc32c.crc32c(b"") + ), + finish_write=True, + ) + start_time = time.perf_counter() + _ = list(self.grpc.BidiWriteObject([r2], context)) + elapsed = time.perf_counter() - start_time + self.assertLess(elapsed, 1) + def test_grpc_bidiread_retry_stall_after_bytes(self): media = self._create_block(5 * 1024 * 1024) response = self.rest_client.put( @@ -1262,6 +1286,12 @@ def test_grpc_bidiread_retry_stall_after_bytes(self): elapsed = time.perf_counter() - start_time self.assertGreater(elapsed, 1) + start_time = time.perf_counter() + response = self.grpc.BidiReadObject([r1], context) + list(response) + elapsed = time.perf_counter() - start_time + self.assertLess(elapsed, 1) + def test_grpc_retry_broken_stream(self): # Use the XML API to inject an object with some data. media = self._create_block(2 * UPLOAD_QUANTUM) From d7511a980227d4e075d2f116fc66fe7abfc0f855 Mon Sep 17 00:00:00 2001 From: raj-prince Date: Tue, 21 Jul 2026 11:36:50 +0000 Subject: [PATCH 4/8] Special handling for zeroK --- testbench/common.py | 9 +++++-- testbench/database.py | 4 ++- testbench/grpc_server.py | 49 +++++++++++++++++++++++++++-------- testbench/rest_server.py | 4 ++- tests/test_testbench_retry.py | 34 +++--------------------- 5 files changed, 55 insertions(+), 45 deletions(-) diff --git a/testbench/common.py b/testbench/common.py index 580afef7..56b62feb 100644 --- a/testbench/common.py +++ b/testbench/common.py @@ -995,9 +995,14 @@ def handle_stall_uploads_after_bytes( e.g. We are uploading 120K of data then, stall-2s-after-100K will stall the request. """ if len(upload.media) <= after_bytes and len(upload.media) + len(data) > after_bytes: + should_stall = True if test_id: - database.dequeue_next_instruction(test_id, "storage.objects.insert") - time.sleep(stall_time) + should_stall = ( + database.dequeue_next_instruction(test_id, "storage.objects.insert") + is not None + ) + if should_stall: + time.sleep(stall_time) def handle_retry_uploads_error_after_bytes( diff --git a/testbench/database.py b/testbench/database.py index f47f9081..d9653f48 100644 --- a/testbench/database.py +++ b/testbench/database.py @@ -752,7 +752,9 @@ def insert_retry_test(self, instructions, transport="HTTP"): def has_instructions_retry_test(self, retry_test_id, method, transport="HTTP"): with self._retry_tests_lock: - retry_test = self.get_retry_test(retry_test_id) + retry_test = self._retry_tests.get(retry_test_id, None) + if retry_test is None: + return False # Add validation for request transport as well. if (len(retry_test["instructions"].get(method, [])) > 0) and retry_test[ "transport" diff --git a/testbench/grpc_server.py b/testbench/grpc_server.py index c4875e4f..d93b5539 100644 --- a/testbench/grpc_server.py +++ b/testbench/grpc_server.py @@ -43,7 +43,15 @@ from google.storage.control.v2 import storage_control_pb2, storage_control_pb2_grpc from google.storage.v2 import storage_pb2, storage_pb2_grpc -_GRPC_SERVER_THREAD_COUNT = 2 +_GRPC_SERVER_THREAD_COUNT = 8 + + +def _should_stall_after_bytes(bytes_yielded, chunk_len, stall_after_bytes): + if chunk_len <= 0: + return False + if stall_after_bytes == 0: + return bytes_yielded == 0 + return bytes_yielded < stall_after_bytes <= bytes_yielded + chunk_len def _trimmed_content(content): @@ -633,14 +641,24 @@ def ReadObject(self, request, context): chunk_len = end - start # Apply stall once when the configured byte threshold is reached. + should_stall_now = _should_stall_after_bytes( + bytes_yielded, chunk_len, stall_after_bytes + ) if ( stall_time > 0 and not stall_applied - and bytes_yielded < stall_after_bytes - and (bytes_yielded + chunk_len) >= stall_after_bytes + and should_stall_now ): - self.db.dequeue_next_instruction(test_id, method) - time.sleep(stall_time) + should_stall = True + dequeue_result = None + if test_id: + dequeue_result = self.db.dequeue_next_instruction(test_id, method) + should_stall = dequeue_result is not None + if should_stall: + print( + f"ReadObject sleeping for {stall_time}s after {stall_after_bytes} bytes" + ) + time.sleep(stall_time) stall_applied = True # Handle retry test broken-stream failures if applicable. @@ -829,14 +847,21 @@ def read_results(): read_range["read_length"] -= excess # Apply stall once when the configured byte threshold is reached. + should_stall_now = _should_stall_after_bytes( + bytes_yielded, len(chunk), stall_after_bytes + ) if ( stall_time > 0 and not stall_applied - and bytes_yielded < stall_after_bytes - and (bytes_yielded + len(chunk)) >= stall_after_bytes + and should_stall_now ): - self.db.dequeue_next_instruction(test_id, method) - time.sleep(stall_time) + should_stall = True + dequeue_result = None + if test_id: + dequeue_result = self.db.dequeue_next_instruction(test_id, method) + should_stall = dequeue_result is not None + if should_stall: + time.sleep(stall_time) stall_applied = True bytes_yielded += len(chunk) @@ -1325,9 +1350,11 @@ def GetStorageLayout(self, request, context): return layout -def run(port, database, echo_metadata=False): +def run(port, database, echo_metadata=False, thread_count=None): server = grpc.server( - futures.ThreadPoolExecutor(max_workers=_GRPC_SERVER_THREAD_COUNT) + futures.ThreadPoolExecutor( + max_workers=thread_count if thread_count is not None else _GRPC_SERVER_THREAD_COUNT + ) ) storage_pb2_grpc.add_StorageServicer_to_server( StorageServicer(database, echo_metadata), server diff --git a/testbench/rest_server.py b/testbench/rest_server.py index ca6c123d..b5999c2c 100644 --- a/testbench/rest_server.py +++ b/testbench/rest_server.py @@ -175,7 +175,9 @@ def start_grpc(): port = flask.request.args.get("port", "0") echo_metadata = flask.request.args.get("echo-metadata", False) grpc_port, grpc_service = testbench.grpc_server.run( - int(port), db, echo_metadata=echo_metadata + int(port), + db, + echo_metadata=echo_metadata, ) return str(grpc_port) diff --git a/tests/test_testbench_retry.py b/tests/test_testbench_retry.py index 1dd14a12..2ee1fd30 100644 --- a/tests/test_testbench_retry.py +++ b/tests/test_testbench_retry.py @@ -1102,17 +1102,6 @@ def test_grpc_retry_stall_read_after_bytes(self): elapsed = time.perf_counter() - start_time self.assertGreater(elapsed, 1) - start_time = time.perf_counter() - response = self.grpc.ReadObject( - storage_pb2.ReadObjectRequest( - bucket="projects/_/buckets/bucket-name", object="512k.txt" - ), - context, - ) - list(response) - elapsed = time.perf_counter() - start_time - self.assertLess(elapsed, 1) - def test_grpc_retry_stall_write_after_bytes(self): response = self.rest_client.post( "/retry_test", @@ -1140,6 +1129,10 @@ def test_grpc_retry_stall_write_after_bytes(self): resource=storage_pb2.Object( name="object-name-stall", bucket="projects/_/buckets/bucket-name", + + + + ) ) ), @@ -1224,19 +1217,6 @@ def test_grpc_retry_stall_bidiwrite_after_bytes(self): elapsed = time.perf_counter() - start_time self.assertGreater(elapsed, 1) - r2 = storage_pb2.BidiWriteObjectRequest( - upload_id=start.upload_id, - write_offset=len(content), - checksummed_data=storage_pb2.ChecksummedData( - content=b"", crc32c=crc32c.crc32c(b"") - ), - finish_write=True, - ) - start_time = time.perf_counter() - _ = list(self.grpc.BidiWriteObject([r2], context)) - elapsed = time.perf_counter() - start_time - self.assertLess(elapsed, 1) - def test_grpc_bidiread_retry_stall_after_bytes(self): media = self._create_block(5 * 1024 * 1024) response = self.rest_client.put( @@ -1286,12 +1266,6 @@ def test_grpc_bidiread_retry_stall_after_bytes(self): elapsed = time.perf_counter() - start_time self.assertGreater(elapsed, 1) - start_time = time.perf_counter() - response = self.grpc.BidiReadObject([r1], context) - list(response) - elapsed = time.perf_counter() - start_time - self.assertLess(elapsed, 1) - def test_grpc_retry_broken_stream(self): # Use the XML API to inject an object with some data. media = self._create_block(2 * UPLOAD_QUANTUM) From afb22b58edd73f36ae7516499ebdd4b43c55fb96 Mon Sep 17 00:00:00 2001 From: raj-prince Date: Tue, 21 Jul 2026 11:53:53 +0000 Subject: [PATCH 5/8] minor refactoring --- testbench/grpc_server.py | 85 +++++++++++++++++++++------------------- 1 file changed, 45 insertions(+), 40 deletions(-) diff --git a/testbench/grpc_server.py b/testbench/grpc_server.py index d93b5539..ce15b843 100644 --- a/testbench/grpc_server.py +++ b/testbench/grpc_server.py @@ -54,6 +54,29 @@ def _should_stall_after_bytes(bytes_yielded, chunk_len, stall_after_bytes): return bytes_yielded < stall_after_bytes <= bytes_yielded + chunk_len +def _apply_grpc_read_stall_if_applicable( + database, + test_id, + method, + bytes_yielded, + chunk_len, + stall_time, + stall_after_bytes, +): + if not test_id: + return False + + if stall_time <= 0 or not _should_stall_after_bytes( + bytes_yielded, chunk_len, stall_after_bytes + ): + return False + + if database.dequeue_next_instruction(test_id, method) is None: + return False + + time.sleep(stall_time) + + def _trimmed_content(content): if len(content) > 10: content = content[:7] @@ -550,7 +573,10 @@ def precondition(_, live_version, ctx): bucket = self.db.get_bucket(request.destination.bucket, context).metadata metadata = storage_pb2.Object() metadata.MergeFrom(request.destination) - (blob, _,) = gcs.object.Object.init( + ( + blob, + _, + ) = gcs.object.Object.init( request, metadata, composed_media, bucket, True, context ) self.db.insert_object( @@ -640,26 +666,15 @@ def ReadObject(self, request, context): end = min(start + size, read_end) chunk_len = end - start - # Apply stall once when the configured byte threshold is reached. - should_stall_now = _should_stall_after_bytes( - bytes_yielded, chunk_len, stall_after_bytes + _apply_grpc_read_stall_if_applicable( + self.db, + test_id, + method, + bytes_yielded, + chunk_len, + stall_time, + stall_after_bytes, ) - if ( - stall_time > 0 - and not stall_applied - and should_stall_now - ): - should_stall = True - dequeue_result = None - if test_id: - dequeue_result = self.db.dequeue_next_instruction(test_id, method) - should_stall = dequeue_result is not None - if should_stall: - print( - f"ReadObject sleeping for {stall_time}s after {stall_after_bytes} bytes" - ) - time.sleep(stall_time) - stall_applied = True # Handle retry test broken-stream failures if applicable. if broken_stream_after_bytes and end >= broken_stream_after_bytes: @@ -846,23 +861,15 @@ def read_results(): range_end = False read_range["read_length"] -= excess - # Apply stall once when the configured byte threshold is reached. - should_stall_now = _should_stall_after_bytes( - bytes_yielded, len(chunk), stall_after_bytes + _apply_grpc_read_stall_if_applicable( + self.db, + test_id, + method, + bytes_yielded, + len(chunk), + stall_time, + stall_after_bytes, ) - if ( - stall_time > 0 - and not stall_applied - and should_stall_now - ): - should_stall = True - dequeue_result = None - if test_id: - dequeue_result = self.db.dequeue_next_instruction(test_id, method) - should_stall = dequeue_result is not None - if should_stall: - time.sleep(stall_time) - stall_applied = True bytes_yielded += len(chunk) returnable -= count @@ -1350,11 +1357,9 @@ def GetStorageLayout(self, request, context): return layout -def run(port, database, echo_metadata=False, thread_count=None): +def run(port, database, echo_metadata=False): server = grpc.server( - futures.ThreadPoolExecutor( - max_workers=thread_count if thread_count is not None else _GRPC_SERVER_THREAD_COUNT - ) + futures.ThreadPoolExecutor(max_workers=_GRPC_SERVER_THREAD_COUNT) ) storage_pb2_grpc.add_StorageServicer_to_server( StorageServicer(database, echo_metadata), server From d24d65b64d4799f5f86ed285ef58b268ad96b2b6 Mon Sep 17 00:00:00 2001 From: raj-prince Date: Tue, 21 Jul 2026 12:02:58 +0000 Subject: [PATCH 6/8] removing unnecessary change --- testbench/grpc_server.py | 7 ++----- testbench/rest_server.py | 8 +++----- tests/test_testbench_retry.py | 4 ---- 3 files changed, 5 insertions(+), 14 deletions(-) diff --git a/testbench/grpc_server.py b/testbench/grpc_server.py index ce15b843..533810ad 100644 --- a/testbench/grpc_server.py +++ b/testbench/grpc_server.py @@ -573,10 +573,7 @@ def precondition(_, live_version, ctx): bucket = self.db.get_bucket(request.destination.bucket, context).metadata metadata = storage_pb2.Object() metadata.MergeFrom(request.destination) - ( - blob, - _, - ) = gcs.object.Object.init( + (blob, _,) = gcs.object.Object.init( request, metadata, composed_media, bucket, True, context ) self.db.insert_object( @@ -833,7 +830,6 @@ def terminate_w_timer(): returnable = ( broken_stream_after_bytes if broken_stream_after_bytes else sys.maxsize ) - bytes_yielded = 0 # We don't want to have a thread blocking on request_iterator, so we # have to handle results in batches rather than concurrently. This @@ -853,6 +849,7 @@ def read_results(): for request in request_iterator: yield from responses_for_range_batch(request.read_ranges) + bytes_yielded = 0 for chunk, range_end, read_range in read_results(): count = len(chunk) excess = count - returnable diff --git a/testbench/rest_server.py b/testbench/rest_server.py index b5999c2c..b08515fe 100644 --- a/testbench/rest_server.py +++ b/testbench/rest_server.py @@ -175,9 +175,7 @@ def start_grpc(): port = flask.request.args.get("port", "0") echo_metadata = flask.request.args.get("echo-metadata", False) grpc_port, grpc_service = testbench.grpc_server.run( - int(port), - db, - echo_metadata=echo_metadata, + int(port), db, echo_metadata=echo_metadata ) return str(grpc_port) @@ -1249,10 +1247,10 @@ def delete_resumable_upload(bucket_name): # === SERVER === # # Define the WSGI application to handle HMAC key and service account requests -(PROJECTS_HANDLER_PATH, projects_app) = projects_rest_server.get_projects_app(db) +PROJECTS_HANDLER_PATH, projects_app = projects_rest_server.get_projects_app(db) # Define the WSGI application to handle IAM requests -(IAM_HANDLER_PATH, iam_app) = iam_rest_server.get_iam_app() +IAM_HANDLER_PATH, iam_app = iam_rest_server.get_iam_app() server = flask.Flask(__name__) server.debug = False diff --git a/tests/test_testbench_retry.py b/tests/test_testbench_retry.py index 2ee1fd30..4d7d6959 100644 --- a/tests/test_testbench_retry.py +++ b/tests/test_testbench_retry.py @@ -1129,10 +1129,6 @@ def test_grpc_retry_stall_write_after_bytes(self): resource=storage_pb2.Object( name="object-name-stall", bucket="projects/_/buckets/bucket-name", - - - - ) ) ), From fbed0b07ec6b31dc1f4481ee26bee34a4635c9cb Mon Sep 17 00:00:00 2001 From: raj-prince Date: Tue, 21 Jul 2026 12:07:12 +0000 Subject: [PATCH 7/8] lint fix --- testbench/grpc_server.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/testbench/grpc_server.py b/testbench/grpc_server.py index 533810ad..4ec3d169 100644 --- a/testbench/grpc_server.py +++ b/testbench/grpc_server.py @@ -573,7 +573,10 @@ def precondition(_, live_version, ctx): bucket = self.db.get_bucket(request.destination.bucket, context).metadata metadata = storage_pb2.Object() metadata.MergeFrom(request.destination) - (blob, _,) = gcs.object.Object.init( + ( + blob, + _, + ) = gcs.object.Object.init( request, metadata, composed_media, bucket, True, context ) self.db.insert_object( From fa48f2d52de52de4d1724adddb06635a05944ac2 Mon Sep 17 00:00:00 2001 From: raj-prince Date: Tue, 21 Jul 2026 12:17:34 +0000 Subject: [PATCH 8/8] lint fix --- testbench/grpc_server.py | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/testbench/grpc_server.py b/testbench/grpc_server.py index 4ec3d169..533810ad 100644 --- a/testbench/grpc_server.py +++ b/testbench/grpc_server.py @@ -573,10 +573,7 @@ def precondition(_, live_version, ctx): bucket = self.db.get_bucket(request.destination.bucket, context).metadata metadata = storage_pb2.Object() metadata.MergeFrom(request.destination) - ( - blob, - _, - ) = gcs.object.Object.init( + (blob, _,) = gcs.object.Object.init( request, metadata, composed_media, bucket, True, context ) self.db.insert_object(