From d0661cc3ada1263bf6a6b1515e4044864ff67385 Mon Sep 17 00:00:00 2001 From: Daniel Sanche Date: Thu, 1 Oct 2026 17:39:53 -0700 Subject: [PATCH 1/3] feat(api_core): resumable upload supports request_body in constructor --- .../api_core/resumable_transfer/upload.py | 17 +++- .../resumable_transfer/upload_async.py | 16 +++- .../asyncio/test_resumable_transfer_async.py | 76 +++++++++++++++ .../tests/unit/test_resumable_transfer.py | 93 +++++++++++++++++++ 4 files changed, 198 insertions(+), 4 deletions(-) diff --git a/packages/google-api-core/google/api_core/resumable_transfer/upload.py b/packages/google-api-core/google/api_core/resumable_transfer/upload.py index 2cc3b1230929..02d6c0024a8e 100644 --- a/packages/google-api-core/google/api_core/resumable_transfer/upload.py +++ b/packages/google-api-core/google/api_core/resumable_transfer/upload.py @@ -105,6 +105,7 @@ def __init__( response_type: Optional[Any] = None, start_retry: Optional[google.api_core.retry.Retry] = None, start_timeout: float = DEFAULT_START_TIMEOUT, + request_body: Optional[Union[str, bytes]] = None, ) -> None: """Initializes a ResumableUploadSession. @@ -132,6 +133,9 @@ def __init__( ``UnseekableStreamError``) are never retried. start_timeout: Timeout in seconds for the start request. Defaults to ``60.0`` seconds. + request_body: Optional default metadata payload sent with the start request. + Used by ``upload()`` and ``iter_upload()`` whenever they are called + without an explicit ``request_body``. When ``None``, an empty payload is sent. """ self._config = config or ResumableUploadConfig() self._transport = transport @@ -139,6 +143,9 @@ def __init__( self._response_type = response_type self._start_retry = start_retry self._start_timeout = start_timeout + self._request_body: Union[str, bytes] = ( + request_body if request_body is not None else "" + ) self._response: Optional[Any] = None self._state = upload_state._ProtocolState( upload_url=upload_url, @@ -810,7 +817,7 @@ def attempt_stream() -> Generator[UploadProgress, None, None]: def upload( self, stream: Union[BinaryIO, bytes, Iterable[bytes]], - request_body: Union[str, bytes] = "", + request_body: Optional[Union[str, bytes]] = None, size: Optional[int] = None, transport: Optional[requests.Session] = None, content_type: Optional[str] = None, @@ -822,6 +829,8 @@ def upload( Args: stream: Data payload to upload (file-like stream, bytes, or iterable of bytes). request_body: Initial metadata payload sent with the start request. + When ``None``, the ``request_body`` supplied to the session constructor + (if any) is used. size: Total stream size in bytes, if known. transport: Optional requests session. content_type: Optional MIME type of the stream payload. @@ -859,7 +868,7 @@ def upload( def iter_upload( self, stream: Union[BinaryIO, bytes, Iterable[bytes]], - request_body: Union[str, bytes] = "", + request_body: Optional[Union[str, bytes]] = None, size: Optional[int] = None, transport: Optional[requests.Session] = None, content_type: Optional[str] = None, @@ -871,6 +880,8 @@ def iter_upload( Args: stream: Data payload to upload (file-like stream, bytes, or iterable of bytes). request_body: Initial metadata payload sent with the start request. + When ``None``, the ``request_body`` supplied to the session constructor + (if any) is used. size: Total stream size in bytes, if known. transport: Optional requests session. content_type: Optional MIME type of the stream payload. @@ -896,6 +907,8 @@ def iter_upload( sess = self._get_transport(transport) if content_type is not None: self._content_type = content_type + if request_body is None: + request_body = self._request_body self._reset_transfer_state() progress_queue: List[UploadProgress] = [] try: diff --git a/packages/google-api-core/google/api_core/resumable_transfer/upload_async.py b/packages/google-api-core/google/api_core/resumable_transfer/upload_async.py index 13cfcb5c6fbd..09ab05968240 100644 --- a/packages/google-api-core/google/api_core/resumable_transfer/upload_async.py +++ b/packages/google-api-core/google/api_core/resumable_transfer/upload_async.py @@ -171,6 +171,7 @@ def __init__( response_type: Optional[Any] = None, start_retry: Optional[google.api_core.retry.AsyncRetry] = None, start_timeout: float = DEFAULT_START_TIMEOUT, + request_body: Optional[Union[str, bytes]] = None, ) -> None: """Initializes an AsyncResumableUploadSession. @@ -198,6 +199,9 @@ def __init__( retried. start_timeout: Timeout in seconds for the start request. Defaults to ``60.0`` seconds. + request_body: Optional default metadata payload sent with the start request. + Used by ``upload()`` whenever it is called without an explicit + ``request_body``. When ``None``, an empty payload is sent. """ self._config = config or ResumableUploadConfig() self._transport = transport @@ -205,6 +209,9 @@ def __init__( self._response_type = response_type self._start_retry = start_retry self._start_timeout = start_timeout + self._request_body: Union[str, bytes] = ( + request_body if request_body is not None else "" + ) self._response: Optional[Any] = None self._state = upload_state._ProtocolState( upload_url=upload_url, @@ -885,7 +892,7 @@ async def attempt_stream() -> AsyncGenerator[UploadProgress, None]: def upload( self, stream: Union[AsyncIterable[bytes], BinaryIO, bytes, Iterable[bytes]], - request_body: Union[str, bytes] = "", + request_body: Optional[Union[str, bytes]] = None, size: Optional[int] = None, transport: Optional[Any] = None, content_type: Optional[str] = None, @@ -897,6 +904,8 @@ def upload( Args: stream: Data payload to upload (async iterable, binary stream, bytes, or iterable). request_body: Initial metadata payload sent with the start request. + When ``None``, the ``request_body`` supplied to the session constructor + (if any) is used. size: Total stream size in bytes, if known. transport: Optional aiohttp client session. content_type: Optional MIME type of the stream payload. @@ -925,6 +934,9 @@ def upload( if content_type is not None: self._content_type = content_type + start_body: Union[str, bytes] = ( + request_body if request_body is not None else self._request_body + ) self._reset_transfer_state() progress_queue: List[UploadProgress] = [] @@ -934,7 +946,7 @@ async def _run() -> AsyncGenerator[UploadProgress, None]: try: await self._initiate( transport=sess, - request_body=request_body, + request_body=start_body, size=computed_size, progress_queue=progress_queue, ) diff --git a/packages/google-api-core/tests/asyncio/test_resumable_transfer_async.py b/packages/google-api-core/tests/asyncio/test_resumable_transfer_async.py index b16211795ee6..9314a82b9e35 100644 --- a/packages/google-api-core/tests/asyncio/test_resumable_transfer_async.py +++ b/packages/google-api-core/tests/asyncio/test_resumable_transfer_async.py @@ -262,6 +262,82 @@ async def test_async_upload_direct_execution() -> None: assert session.upload_url == "https://upload.example.com/resumable-async" +def _make_single_chunk_async_transport() -> DummyAsyncSession: + """Builds a DummyAsyncSession returning a start response and a final chunk response.""" + start_resp = DummyAsyncResponse( + status=200, + headers={ + "X-Goog-Upload-Status": "active", + "X-Goog-Upload-URL": "https://upload.example.com/resumable-async", + }, + body=b"", + ) + chunk_resp = DummyAsyncResponse( + status=200, + headers={"X-Goog-Upload-Status": "final"}, + body=b"", + ) + return DummyAsyncSession([start_resp, chunk_resp]) + + +@pytest.mark.asyncio +async def test_async_upload_uses_session_request_body() -> None: + """Verifies the constructor request_body is sent with the start request by default.""" + async_transport = _make_single_chunk_async_transport() + session = AsyncResumableUploadSession( + upload_url="https://api.example.com/start", + transport=async_transport, + request_body='{"name": "from-init"}', + ) + + await session.upload(stream=b"payload") + + _, _, start_kwargs = async_transport.requests[0] + assert start_kwargs["data"] == b'{"name": "from-init"}' + assert start_kwargs["headers"]["X-Goog-Upload-Command"] == "start" + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "override, expected", + [ + ('{"name": "override"}', b'{"name": "override"}'), + (b'{"name": "override-bytes"}', b'{"name": "override-bytes"}'), + ("", b""), + ], +) +async def test_async_upload_request_body_argument_overrides_session_default( + override: Union[str, bytes], expected: bytes +) -> None: + """Verifies an explicit per-call request_body (even empty) overrides the constructor default.""" + async_transport = _make_single_chunk_async_transport() + session = AsyncResumableUploadSession( + upload_url="https://api.example.com/start", + transport=async_transport, + request_body='{"name": "from-init"}', + ) + + await session.upload(stream=b"payload", request_body=override) + + _, _, start_kwargs = async_transport.requests[0] + assert start_kwargs["data"] == expected + + +@pytest.mark.asyncio +async def test_async_upload_request_body_defaults_to_empty() -> None: + """Verifies the start request carries an empty payload when no request_body is configured.""" + async_transport = _make_single_chunk_async_transport() + session = AsyncResumableUploadSession( + upload_url="https://api.example.com/start", + transport=async_transport, + ) + + await session.upload(stream=b"payload") + + _, _, start_kwargs = async_transport.requests[0] + assert start_kwargs["data"] == b"" + + @pytest.mark.asyncio async def test_async_upload_multi_chunk_operation_handle() -> None: """Verifies multi-chunk upload dispatching and AsyncUploadOperation property handles.""" diff --git a/packages/google-api-core/tests/unit/test_resumable_transfer.py b/packages/google-api-core/tests/unit/test_resumable_transfer.py index d783af5432a4..94321581ddc8 100644 --- a/packages/google-api-core/tests/unit/test_resumable_transfer.py +++ b/packages/google-api-core/tests/unit/test_resumable_transfer.py @@ -366,6 +366,99 @@ def test_sync_upload_iterative_progress(): assert session.response.size == 8 +def _make_single_chunk_sync_transport() -> mock.Mock: + """Builds a mock requests.Session returning a start response and a final chunk response.""" + session_transport = mock.create_autospec(requests.Session, instance=True) + start_resp = mock.Mock( + ok=True, + status_code=200, + headers={ + "X-Goog-Upload-Status": "active", + "X-Goog-Upload-URL": "https://upload.example.com/resumable-123", + }, + ) + chunk_resp = mock.Mock( + ok=True, + status_code=200, + headers={"X-Goog-Upload-Status": "final"}, + content=b"", + ) + session_transport.request.side_effect = [start_resp, chunk_resp] + return session_transport + + +def test_sync_upload_uses_session_request_body(): + # A request_body supplied to the constructor is sent with the start + # request when upload() is called without an explicit request_body. + session_transport = _make_single_chunk_sync_transport() + session = ResumableUploadSession( + upload_url="https://api.example.com/start", + transport=session_transport, + request_body='{"name": "from-init"}', + ) + + session.upload(stream=b"payload") + + start_call = session_transport.request.call_args_list[0] + assert start_call.kwargs["data"] == b'{"name": "from-init"}' + assert start_call.kwargs["headers"]["X-Goog-Upload-Command"] == "start" + + +def test_sync_iter_upload_uses_session_request_body(): + session_transport = _make_single_chunk_sync_transport() + session = ResumableUploadSession( + upload_url="https://api.example.com/start", + transport=session_transport, + request_body=b'{"name": "from-init-bytes"}', + ) + + list(session.iter_upload(stream=b"payload")) + + start_call = session_transport.request.call_args_list[0] + assert start_call.kwargs["data"] == b'{"name": "from-init-bytes"}' + + +@pytest.mark.parametrize( + "override, expected", + [ + ('{"name": "override"}', b'{"name": "override"}'), + (b'{"name": "override-bytes"}', b'{"name": "override-bytes"}'), + ("", b""), + ], +) +def test_sync_upload_request_body_argument_overrides_session_default( + override, expected +): + # An explicit per-call request_body (including an empty one) takes + # precedence over the request_body supplied to the constructor. + session_transport = _make_single_chunk_sync_transport() + session = ResumableUploadSession( + upload_url="https://api.example.com/start", + transport=session_transport, + request_body='{"name": "from-init"}', + ) + + session.upload(stream=b"payload", request_body=override) + + start_call = session_transport.request.call_args_list[0] + assert start_call.kwargs["data"] == expected + + +def test_sync_upload_request_body_defaults_to_empty(): + # Without a request_body on the constructor or the call, the start + # request carries an empty payload. + session_transport = _make_single_chunk_sync_transport() + session = ResumableUploadSession( + upload_url="https://api.example.com/start", + transport=session_transport, + ) + + session.upload(stream=b"payload") + + start_call = session_transport.request.call_args_list[0] + assert start_call.kwargs["data"] == b"" + + def test_sync_resume(): session_transport = mock.create_autospec(requests.Session, instance=True) From 13e4a2874ea83e232b8a9646d8c1b642298be508 Mon Sep 17 00:00:00 2001 From: Daniel Sanche Date: Thu, 1 Oct 2026 17:44:05 -0700 Subject: [PATCH 2/3] added test --- .../asyncio/test_resumable_transfer_async.py | 16 ++++++++++++++++ .../tests/unit/test_resumable_transfer.py | 16 ++++++++++++++++ 2 files changed, 32 insertions(+) diff --git a/packages/google-api-core/tests/asyncio/test_resumable_transfer_async.py b/packages/google-api-core/tests/asyncio/test_resumable_transfer_async.py index 9314a82b9e35..4b799647946b 100644 --- a/packages/google-api-core/tests/asyncio/test_resumable_transfer_async.py +++ b/packages/google-api-core/tests/asyncio/test_resumable_transfer_async.py @@ -323,6 +323,22 @@ async def test_async_upload_request_body_argument_overrides_session_default( assert start_kwargs["data"] == expected +@pytest.mark.asyncio +async def test_async_upload_request_body_positional_override() -> None: + """Verifies request_body remains the second positional parameter of upload().""" + async_transport = _make_single_chunk_async_transport() + session = AsyncResumableUploadSession( + upload_url="https://api.example.com/start", + transport=async_transport, + request_body='{"name": "from-init"}', + ) + + await session.upload(b"payload", '{"name": "positional"}') + + _, _, start_kwargs = async_transport.requests[0] + assert start_kwargs["data"] == b'{"name": "positional"}' + + @pytest.mark.asyncio async def test_async_upload_request_body_defaults_to_empty() -> None: """Verifies the start request carries an empty payload when no request_body is configured.""" diff --git a/packages/google-api-core/tests/unit/test_resumable_transfer.py b/packages/google-api-core/tests/unit/test_resumable_transfer.py index 94321581ddc8..d207910f0ba9 100644 --- a/packages/google-api-core/tests/unit/test_resumable_transfer.py +++ b/packages/google-api-core/tests/unit/test_resumable_transfer.py @@ -444,6 +444,22 @@ def test_sync_upload_request_body_argument_overrides_session_default( assert start_call.kwargs["data"] == expected +def test_sync_upload_request_body_positional_override(): + # request_body stays the second positional parameter, so callers that + # pass it positionally are not broken by the constructor default. + session_transport = _make_single_chunk_sync_transport() + session = ResumableUploadSession( + upload_url="https://api.example.com/start", + transport=session_transport, + request_body='{"name": "from-init"}', + ) + + session.upload(b"payload", '{"name": "positional"}') + + start_call = session_transport.request.call_args_list[0] + assert start_call.kwargs["data"] == b'{"name": "positional"}' + + def test_sync_upload_request_body_defaults_to_empty(): # Without a request_body on the constructor or the call, the start # request carries an empty payload. From c0f9ee44f968ba92ff633b2323b7cb4638b25ca5 Mon Sep 17 00:00:00 2001 From: Daniel Sanche Date: Thu, 1 Oct 2026 17:56:27 -0700 Subject: [PATCH 3/3] removed extra variable --- .../google/api_core/resumable_transfer/upload_async.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/packages/google-api-core/google/api_core/resumable_transfer/upload_async.py b/packages/google-api-core/google/api_core/resumable_transfer/upload_async.py index 09ab05968240..538d2c050650 100644 --- a/packages/google-api-core/google/api_core/resumable_transfer/upload_async.py +++ b/packages/google-api-core/google/api_core/resumable_transfer/upload_async.py @@ -934,9 +934,8 @@ def upload( if content_type is not None: self._content_type = content_type - start_body: Union[str, bytes] = ( - request_body if request_body is not None else self._request_body - ) + if request_body is None: + request_body = self._request_body self._reset_transfer_state() progress_queue: List[UploadProgress] = [] @@ -946,7 +945,7 @@ async def _run() -> AsyncGenerator[UploadProgress, None]: try: await self._initiate( transport=sess, - request_body=start_body, + request_body=request_body, size=computed_size, progress_queue=progress_queue, )