Skip to content

Commit ef11aeb

Browse files
committed
removes metadata from reader and writer | updates docstring
1 parent 10af2a5 commit ef11aeb

6 files changed

Lines changed: 72 additions & 280 deletions

File tree

packages/google-cloud-storage/google/cloud/storage/asyncio/async_appendable_object_writer.py

Lines changed: 1 addition & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -213,7 +213,6 @@ def __init__(
213213
self.object_resource: Optional[_storage_v2.Object] = None
214214
self._flush_count = 0
215215
self.blob: Optional[Blob] = None
216-
self.metadata: Optional[List[Tuple[str, str]]] = None
217216

218217
@classmethod
219218
def from_blob(
@@ -313,8 +312,6 @@ async def open(
313312
if self._is_stream_open:
314313
raise ValueError("Underlying bidi-gRPC stream is already open")
315314

316-
self.metadata = metadata
317-
318315
if retry_policy is None:
319316
retry_policy = AsyncRetry(
320317
predicate=_is_write_retryable, on_error=self._on_open_error
@@ -337,7 +334,7 @@ def combined_on_error(exc):
337334
)
338335

339336
async def _do_open():
340-
current_metadata = list(self.metadata) if self.metadata else []
337+
current_metadata = list(metadata) if metadata else []
341338

342339
# Cleanup stream from previous failed attempt, if any.
343340
if self.write_obj_stream:
@@ -411,11 +408,6 @@ async def append(
411408
412409
:raises ValueError: If the stream is not open.
413410
"""
414-
if metadata is not None:
415-
self.metadata = metadata
416-
else:
417-
metadata = self.metadata
418-
419411
if not self._is_stream_open:
420412
raise ValueError("Stream is not open. Call open() before append().")
421413
if not data:

packages/google-cloud-storage/google/cloud/storage/asyncio/async_grpc_client.py

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,22 @@
2626
_DEFAULT_HOST = "storage.googleapis.com"
2727

2828

29+
def _validate_metadata(metadata):
30+
"""Validates that metadata is a sequence of (key, value) pairs."""
31+
if metadata is None:
32+
return
33+
if not isinstance(metadata, (list, tuple)):
34+
raise TypeError("metadata must be a list or tuple of (key, value) pairs.")
35+
for item in metadata:
36+
if not isinstance(item, (list, tuple)) or len(item) != 2:
37+
raise ValueError(
38+
"Each element in metadata must be a list or tuple of exactly 2 strings (key, value)."
39+
)
40+
key, value = item
41+
if not isinstance(key, str) or not isinstance(value, str):
42+
raise TypeError("Both key and value in metadata pairs must be strings.")
43+
44+
2945
class AsyncGrpcClient:
3046
"""An asynchronous client for interacting with Google Cloud Storage using the gRPC API.
3147
@@ -184,8 +200,18 @@ async def delete_object(
184200
:type if_metageneration_not_match: int
185201
:param if_metageneration_not_match: (Optional)
186202
203+
:type metadata: Sequence[Tuple[str, str]]
204+
:param metadata: (Optional) Additional metadata that is provided to the method.
205+
206+
:type timeout: float or None
207+
:param timeout:
208+
(Optional) The amount of time, in seconds, to wait for the request to
209+
complete.
187210
211+
:type retry: :class:`~google.api_core.retry_async.AsyncRetry` or :class:`~google.api_core.retry.Retry` or None
212+
:param retry: (Optional) Designation of what errors, if any, should be retried.
188213
"""
214+
_validate_metadata(metadata)
189215
# The gRPC API requires the bucket name to be in the format "projects/_/buckets/bucket_name"
190216
bucket_path = f"projects/_/buckets/{bucket_name}"
191217
request = storage_v2.DeleteObjectRequest(
@@ -251,9 +277,21 @@ async def get_object(
251277
:param soft_deleted:
252278
(Optional) If True, return the soft-deleted version of this object.
253279
280+
:type metadata: Sequence[Tuple[str, str]]
281+
:param metadata: (Optional) Additional metadata that is provided to the method.
282+
283+
:type timeout: float or None
284+
:param timeout:
285+
(Optional) The amount of time, in seconds, to wait for the request to
286+
complete.
287+
288+
:type retry: :class:`~google.api_core.retry_async.AsyncRetry` or :class:`~google.api_core.retry.Retry` or None
289+
:param retry: (Optional) Designation of what errors, if any, should be retried.
290+
254291
:rtype: :class:`google.cloud._storage_v2.types.Object`
255292
:returns: The object metadata resource.
256293
"""
294+
_validate_metadata(metadata)
257295
bucket_path = f"projects/_/buckets/{bucket_name}"
258296

259297
request = storage_v2.GetObjectRequest(

packages/google-cloud-storage/google/cloud/storage/asyncio/async_multi_range_downloader.py

Lines changed: 1 addition & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -233,7 +233,6 @@ def __init__(
233233
self._open_retries: int = 0
234234
self.is_finalized: bool = False
235235
self.full_obj_server_crc32c: Optional[int] = None
236-
self.metadata: Optional[List[Tuple[str, str]]] = None
237236

238237
async def __aenter__(self):
239238
"""Opens the underlying bidi-gRPC connection to read from the object."""
@@ -263,8 +262,6 @@ async def open(
263262
if self._is_stream_open:
264263
raise ValueError("Underlying bidi-gRPC stream is already open")
265264

266-
self.metadata = metadata
267-
268265
if retry_policy is None:
269266

270267
def on_error_wrapper(exc):
@@ -293,7 +290,7 @@ def combined_on_error(exc):
293290
)
294291

295292
async def _do_open():
296-
current_metadata = list(self.metadata) if self.metadata else []
293+
current_metadata = list(metadata) if metadata else []
297294

298295
# Cleanup stream from previous failed attempt, if any.
299296
if self.read_obj_str:
@@ -416,11 +413,6 @@ async def download_ranges(
416413
417414
"""
418415

419-
if metadata is not None:
420-
self.metadata = metadata
421-
else:
422-
metadata = self.metadata
423-
424416
if len(read_ranges) > 1000:
425417
raise ValueError(
426418
"Invalid input - length of read_ranges cannot be more than 1000"

packages/google-cloud-storage/tests/unit/asyncio/test_async_appendable_object_writer.py

Lines changed: 0 additions & 87 deletions
Original file line numberDiff line numberDiff line change
@@ -503,90 +503,3 @@ async def test_methods_require_open_stream_raises(self, mock_appendable_writer):
503503
for coro in methods:
504504
with pytest.raises(ValueError, match="Stream is not open"):
505505
await coro
506-
507-
@pytest.mark.asyncio
508-
async def test_append_persists_metadata_on_resumption(self, mock_appendable_writer):
509-
# Arrange
510-
mock_client = mock_appendable_writer["mock_client"]
511-
mock_stream = mock_appendable_writer["mock_stream"]
512-
513-
test_metadata = [("custom-key", "custom-value")]
514-
writer = self._make_one(mock_client)
515-
516-
# Act - Open with metadata
517-
await writer.open(metadata=test_metadata)
518-
519-
# Assert first open used metadata
520-
mock_stream.open.assert_called_once_with(metadata=test_metadata)
521-
assert writer.metadata == test_metadata
522-
523-
# Setup resumption trigger
524-
retryable_exc = exceptions.ServiceUnavailable("Retry me")
525-
mock_stream.send.side_effect = retryable_exc
526-
527-
# Reset mock_stream.open call count to verify it is called again
528-
mock_stream.open.reset_mock()
529-
530-
# Setup a fast retry policy to fail quickly in test
531-
from google.api_core.retry_async import AsyncRetry
532-
533-
fast_retry = AsyncRetry(
534-
predicate=lambda e: True,
535-
initial=0.01,
536-
maximum=0.01,
537-
multiplier=1.0,
538-
deadline=0.1,
539-
)
540-
541-
# Act - append (should trigger retry and use stored metadata)
542-
from google.api_core.exceptions import RetryError
543-
544-
try:
545-
await writer.append(b"data", retry_policy=fast_retry)
546-
except RetryError:
547-
pass
548-
549-
# Assert second open (during retry) used the same test_metadata
550-
mock_stream.open.assert_called_with(metadata=test_metadata)
551-
552-
@pytest.mark.asyncio
553-
async def test_append_updates_and_persists_metadata_on_resumption(
554-
self, mock_appendable_writer
555-
):
556-
# Arrange
557-
mock_client = mock_appendable_writer["mock_client"]
558-
mock_stream = mock_appendable_writer["mock_stream"]
559-
560-
initial_metadata = [("initial-key", "initial-value")]
561-
updated_metadata = [("updated-key", "updated-value")]
562-
writer = self._make_one(mock_client)
563-
564-
# Act - Open with initial metadata
565-
await writer.open(metadata=initial_metadata)
566-
assert writer.metadata == initial_metadata
567-
568-
# Setup resumption trigger when append is called
569-
retryable_exc = exceptions.ServiceUnavailable("Retry me")
570-
mock_stream.send.side_effect = retryable_exc
571-
mock_stream.open.reset_mock()
572-
573-
from google.api_core.retry_async import AsyncRetry
574-
fast_retry = AsyncRetry(
575-
predicate=lambda e: True,
576-
initial=0.01,
577-
maximum=0.01,
578-
multiplier=1.0,
579-
deadline=0.1,
580-
)
581-
582-
from google.api_core.exceptions import RetryError
583-
try:
584-
await writer.append(
585-
b"data", retry_policy=fast_retry, metadata=updated_metadata
586-
)
587-
except RetryError:
588-
pass
589-
590-
# Assert writer.metadata was updated and used during stream reopening
591-
assert writer.metadata == updated_metadata
592-
mock_stream.open.assert_called_with(metadata=updated_metadata)

packages/google-cloud-storage/tests/unit/asyncio/test_async_grpc_client.py

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -399,3 +399,35 @@ async def test_get_object_with_all_parameters(self, mock_async_storage_client):
399399
assert call_kwargs["metadata"] == metadata
400400
assert call_kwargs["timeout"] == timeout
401401
assert call_kwargs["retry"] == retry
402+
403+
@pytest.mark.asyncio
404+
@pytest.mark.parametrize(
405+
"invalid_metadata, expected_exc",
406+
[
407+
("not-a-sequence", TypeError),
408+
([("key_only",)], ValueError),
409+
([("too", "many", "items")], ValueError),
410+
([(123, "val")], TypeError),
411+
([("key", 456)], TypeError),
412+
],
413+
)
414+
async def test_delete_object_invalid_metadata(self, invalid_metadata, expected_exc):
415+
client = async_grpc_client.AsyncGrpcClient(credentials=_make_credentials())
416+
with pytest.raises(expected_exc):
417+
await client.delete_object("bucket", "object", metadata=invalid_metadata)
418+
419+
@pytest.mark.asyncio
420+
@pytest.mark.parametrize(
421+
"invalid_metadata, expected_exc",
422+
[
423+
("not-a-sequence", TypeError),
424+
([("key_only",)], ValueError),
425+
([("too", "many", "items")], ValueError),
426+
([(123, "val")], TypeError),
427+
([("key", 456)], TypeError),
428+
],
429+
)
430+
async def test_get_object_invalid_metadata(self, invalid_metadata, expected_exc):
431+
client = async_grpc_client.AsyncGrpcClient(credentials=_make_credentials())
432+
with pytest.raises(expected_exc):
433+
await client.get_object("bucket", "object", metadata=invalid_metadata)

0 commit comments

Comments
 (0)