Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -132,13 +133,19 @@ 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
self._content_type = content_type
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,
Expand Down Expand Up @@ -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,
Expand All @@ -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.
Expand Down Expand Up @@ -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,
Expand All @@ -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.
Expand All @@ -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:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -198,13 +199,19 @@ 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
self._content_type = content_type
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,
Expand Down Expand Up @@ -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,
Expand All @@ -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.
Expand Down Expand Up @@ -925,6 +934,8 @@ def upload(

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] = []
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -262,6 +262,98 @@ 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_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."""
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."""
Expand Down
109 changes: 109 additions & 0 deletions packages/google-api-core/tests/unit/test_resumable_transfer.py
Original file line number Diff line number Diff line change
Expand Up @@ -366,6 +366,115 @@ 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_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.
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)

Expand Down
Loading