From 2ed7659ed8152866f5489e77e5bad16ba6ae2216 Mon Sep 17 00:00:00 2001 From: Daniel Alley Date: Thu, 10 Sep 2026 13:33:52 -0400 Subject: [PATCH 1/3] Stream proxied S3 and Azure artifacts Bypass django-storages eager file buffering when the content app proxies S3 or Azure artifacts. Preserve the generic storage fallback and HTTP range behavior. closes #7806 Assisted By: Codex (GPT-5) Terra 5.6 --- CHANGES/7806.bugfix | 1 + CLAUDE.md | 7 + .../functional/api/test_download_policies.py | 74 +++++- pulpcore/_object_storage.py | 166 +++++++++++++ pulpcore/responses.py | 39 ++++ .../api/test_artifact_distribution.py | 4 +- pulpcore/tests/unit/test_object_storage.py | 220 ++++++++++++++++++ 7 files changed, 508 insertions(+), 3 deletions(-) create mode 100644 CHANGES/7806.bugfix create mode 100644 pulpcore/_object_storage.py create mode 100644 pulpcore/tests/unit/test_object_storage.py diff --git a/CHANGES/7806.bugfix b/CHANGES/7806.bugfix new file mode 100644 index 00000000000..a1a7a1eae7c --- /dev/null +++ b/CHANGES/7806.bugfix @@ -0,0 +1 @@ +Resolve out-of-memory conditions when serving large artifacts from S3/Azure when REDIRECT_TO_OBJECT_STORAGE=False. diff --git a/CLAUDE.md b/CLAUDE.md index e775a8afcf5..266d3ad43b4 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -38,6 +38,13 @@ pulpcore & pulp-file functional tests require both client bindings to be install **Always** use the `oci-env` to run the functional and unit tests. +The active OCI profile may mount this checkout at `/src/pulpcore`, while `oci-env test -p core` +expects `/src/core`. If that wrapper cannot find `/src/core`, run the focused test with +`oci-env exec pytest /src/pulpcore/` instead. + +`ArtifactDistribution` is matched within its `pulp_domain`. Its generated artifact URL includes +the artifact domain, so test-created artifact distributions must use that same domain. + ## Modifying template_config.yml Use the `plugin-template` tool after any changes made to `template_config.yml`. diff --git a/pulp_file/tests/functional/api/test_download_policies.py b/pulp_file/tests/functional/api/test_download_policies.py index 46020243f1b..4f8cae6cd6f 100644 --- a/pulp_file/tests/functional/api/test_download_policies.py +++ b/pulp_file/tests/functional/api/test_download_policies.py @@ -7,10 +7,15 @@ from urllib.parse import urljoin import pytest +import requests from aiohttp.client_exceptions import ClientResponseError from bs4 import BeautifulSoup -from pulpcore.client.pulp_file import FileFilePublication, FileRepositorySyncURL +from pulpcore.client.pulp_file import ( + FileFilePublication, + FileRepositorySyncURL, + RepositoryAddRemoveContent, +) from pulpcore.tests.functional.utils import download_file, get_files_in_manifest OBJECT_STORAGES = ( @@ -37,6 +42,73 @@ def _do_range_request_download_and_assert(url, range_header, expected_bytes): ) +@pytest.mark.parametrize( + "storage_class", + ( + pytest.param("storages.backends.s3.S3Storage", id="s3"), + pytest.param("storages.backends.s3boto3.S3Boto3Storage", id="s3boto3"), + pytest.param("storages.backends.azure_storage.AzureStorage", id="azure"), + ), +) +def test_proxied_object_storage_artifact_streaming( + storage_class, + pulp_settings, + domain_factory, + random_artifact_factory, + file_bindings, + file_repository_factory, + file_publication_factory, + file_distribution_factory, + distribution_base_url, + gen_object_with_cleanup, + monitor_task, +): + """Proxy object-storage content served by a distribution without redirecting.""" + if pulp_settings.STORAGES["default"]["BACKEND"] != storage_class: + pytest.skip("The functional environment does not provide this object-storage configuration") + + domain = domain_factory(storage_class=storage_class, redirect_to_object_storage=False) + artifact = random_artifact_factory(pulp_domain=domain.name, size=32) + content = gen_object_with_cleanup( + file_bindings.ContentFilesApi, + artifact=artifact.pulp_href, + relative_path=str(uuid.uuid4()), + pulp_domain=domain.name, + ) + repository = file_repository_factory(pulp_domain=domain.name) + monitor_task( + file_bindings.RepositoriesFileApi.modify( + repository.pulp_href, + RepositoryAddRemoveContent(add_content_units=[content.pulp_href]), + ).task + ) + publication = file_publication_factory( + pulp_domain=domain.name, + repository=repository.pulp_href, + ) + distribution = file_distribution_factory( + pulp_domain=domain.name, + publication=publication.pulp_href, + ) + content_url = urljoin(distribution_base_url(distribution.base_url), content.relative_path) + + full_response = requests.get(content_url, allow_redirects=False) + assert full_response.status_code == requests.codes.ok + assert not full_response.is_redirect + assert hashlib.sha256(full_response.content).hexdigest() == artifact.sha256 + assert full_response.headers["Content-Length"] == str(len(full_response.content)) + assert full_response.headers["Accept-Ranges"] == "bytes" + + range_response = requests.get( + content_url, headers={"Range": "bytes=1-4"}, allow_redirects=False + ) + assert range_response.status_code == requests.codes.partial_content + assert not range_response.is_redirect + assert range_response.content == full_response.content[1:5] + assert range_response.headers["Content-Length"] == "4" + assert range_response.headers["Content-Range"] == f"bytes 1-4/{len(full_response.content)}" + + @pytest.mark.parallel @pytest.mark.parametrize("download_policy", ["immediate", "on_demand", "streamed"]) def test_download_policy( diff --git a/pulpcore/_object_storage.py b/pulpcore/_object_storage.py new file mode 100644 index 00000000000..1af85791151 --- /dev/null +++ b/pulpcore/_object_storage.py @@ -0,0 +1,166 @@ +"""Temporary private streaming adapters for object-storage response bodies. + +django-storages file objects are seekable Django file objects. Opening one for +reading downloads the complete object into a temporary file, which is unsuitable +for the content app's bounded response loop. These adapters expose only the +synchronous `open/read/close` interface needed by :mod:`pulpcore.responses`. + +Remove these adapters when django-storages provides a stable, cross-backend, +streaming-read API that Pulp can use instead. +""" + +from typing import Any + +S3_STORAGE_CLASSES = frozenset( + ( + "storages.backends.s3.S3Storage", + "storages.backends.s3boto3.S3Boto3Storage", + ) +) +AZURE_STORAGE_CLASSES = frozenset(("storages.backends.azure_storage.AzureStorage",)) +STREAMING_STORAGE_CLASSES = S3_STORAGE_CLASSES | AZURE_STORAGE_CLASSES + +# These are the read arguments accepted by boto3's get_object operation and by +# s3transfer's download argument filter. The import is optional because pulpcore +# can be installed without the S3 extra. +try: + from s3transfer.constants import ALLOWED_DOWNLOAD_ARGS as _S3_DOWNLOAD_ARGS +except ImportError: # pragma: no cover - exercised only without the S3 extra + _S3_DOWNLOAD_ARGS = ( + "ChecksumMode", + "ExpectedBucketOwner", + "IfMatch", + "IfModifiedSince", + "IfNoneMatch", + "IfUnmodifiedSince", + "PartNumber", + "RequestPayer", + "SSECustomerAlgorithm", + "SSECustomerKey", + "SSECustomerKeyMD5", + "VersionId", + ) + + +def _filter_s3_download_params(params: dict[str, Any]) -> dict[str, Any]: + """Retain only configured values accepted by boto3 `get_object`.""" + + return {key: value for key, value in params.items() if key in _S3_DOWNLOAD_ARGS} + + +def _clean_s3_name(name): + """Match django-storages' logical-name processing before adding `location`.""" + + try: + from storages.utils import clean_name + except ImportError: # pragma: no cover - S3 storage requires django-storages + import posixpath + + cleaned_name = posixpath.normpath(name).replace("\\", "/") + if name.endswith("/") and not cleaned_name.endswith("/"): + cleaned_name += "/" + return "" if cleaned_name == "." else cleaned_name + return clean_name(name) + + +class S3Stream: + """A ranged reader backed by the configured django-storages S3 client.""" + + def __init__(self, storage, name: str, offset: int, count: int): + self.storage = storage + self.name = name + self.offset = offset + self.count = count + self.body = None + + def open(self): + name = _clean_s3_name(self.name) + key = self.storage._normalize_name(name) + # `get_object_parameters` receives the logical name; `Key` includes + # the backend's location prefix just as django-storages' `open` does. + params = _filter_s3_download_params(self.storage.get_object_parameters(name)) + params.update( + Bucket=self.storage.bucket_name, + Key=key, + Range=f"bytes={self.offset}-{self.offset + self.count - 1}", + ) + response = self.storage.connection.meta.client.get_object(**params) + self.body = response["Body"] + return self + + def read(self, size: int) -> bytes: + return self.body.read(size) + + def close(self): + if self.body is not None: + self.body.close() + self.body = None + + +class AzureStream: + """A ranged reader backed by an Azure `StorageStreamDownloader`.""" + + def __init__(self, storage, name: str, offset: int, count: int): + self.storage = storage + self.name = name + self.offset = offset + self.count = count + self.downloader = None + self.chunks = None + self.pending = b"" + + def open(self): + path = self.storage._get_valid_path(self.name) + self.downloader = self.storage.client.download_blob( + path, + offset=self.offset, + length=self.count, + timeout=self.storage.timeout, + ) + self.chunks = iter(self.downloader.chunks()) + return self + + def read(self, size: int) -> bytes: + # Azure controls the size of values yielded by `chunks()`. Keep at + # most one such value pending while presenting Pulp's smaller read size. + while len(self.pending) < size: + try: + self.pending += next(self.chunks) + except StopIteration: + break + + chunk = self.pending[:size] + self.pending = self.pending[size:] + return chunk + + def close(self): + if self.downloader is None: + return + + close = getattr(self.downloader, "close", None) + if close is None: + response = getattr(self.downloader, "_response", None) + for response_part in ( + response, + getattr(response, "http_response", None), + getattr(getattr(response, "http_response", None), "internal_response", None), + ): + close = getattr(response_part, "close", None) + if close is not None: + close() + break + else: + close() + self.downloader = None + self.chunks = None + self.pending = b"" + + +def get_stream(storage_class: str, storage, name: str, offset: int, count: int): + """Return an adapter for a supported storage class, or `None` for fallback.""" + + if storage_class in S3_STORAGE_CLASSES: + return S3Stream(storage, name, offset, count) + if storage_class in AZURE_STORAGE_CLASSES: + return AzureStream(storage, name, offset, count) + return None diff --git a/pulpcore/responses.py b/pulpcore/responses.py index 1b1fac62a0d..81f5c40e652 100644 --- a/pulpcore/responses.py +++ b/pulpcore/responses.py @@ -7,6 +7,7 @@ HTTPRequestRangeNotSatisfiable, ) +from pulpcore._object_storage import STREAMING_STORAGE_CLASSES, get_stream from pulpcore.app.models import Artifact @@ -33,6 +34,17 @@ def __init__( self._chunk_size = chunk_size async def _sendfile(self, request, fobj, offset, count): + storage_class = self._artifact.pulp_domain.storage_class + if storage_class in STREAMING_STORAGE_CLASSES: + object_stream = get_stream( + storage_class, + self._artifact.pulp_domain.get_storage(), + fobj.name, + offset, + count, + ) + return await self._sendfile_object_storage(request, object_stream, count) + # To keep memory usage low, fobj is transferred in chunks # controlled by the constructor's chunk_size argument. @@ -54,6 +66,33 @@ async def _sendfile(self, request, fobj, offset, count): await writer.drain() return writer + async def _sendfile_object_storage(self, request, object_stream, count): + """Write a blocking provider stream without blocking the content-app loop. + + Object storage SDKs are synchronous. Adapter methods therefore run in a + worker thread, while `writer.write` maintains aiohttp backpressure. The + adapter is closed even if opening it or writing its response fails. + """ + writer = await super().prepare(request) + assert writer is not None + + stream = object_stream + try: + stream = await asyncio.to_thread(object_stream.open) + remaining = count + while remaining: + chunk = await asyncio.to_thread(stream.read, min(self._chunk_size, remaining)) + if not chunk: + break + if len(chunk) > remaining: + chunk = chunk[:remaining] + await writer.write(chunk) + remaining -= len(chunk) + await writer.drain() + return writer + finally: + await asyncio.to_thread(stream.close) + async def prepare(self, request): if self._artifact is None: self._artifact = await Artifact.objects.select_related("pulp_domain").aget( diff --git a/pulpcore/tests/functional/api/test_artifact_distribution.py b/pulpcore/tests/functional/api/test_artifact_distribution.py index 14b6dc7e862..dc493981af6 100644 --- a/pulpcore/tests/functional/api/test_artifact_distribution.py +++ b/pulpcore/tests/functional/api/test_artifact_distribution.py @@ -11,7 +11,7 @@ ) -def test_artifact_distribution(random_artifact, pulp_settings): +def test_artifact_distribution(random_artifact, pulp_settings, distribution_base_url): settings = pulp_settings artifact_uuid = random_artifact.pulp_href.split("/")[-2] @@ -22,7 +22,7 @@ def test_artifact_distribution(random_artifact, pulp_settings): ) process = subprocess.run(["pulpcore-manager", "shell", "-c", commands], capture_output=True) assert process.returncode == 0 - artifact_url = process.stdout.decode().strip() + artifact_url = distribution_base_url(process.stdout.decode().strip()) response = requests.get(artifact_url) response.raise_for_status() diff --git a/pulpcore/tests/unit/test_object_storage.py b/pulpcore/tests/unit/test_object_storage.py new file mode 100644 index 00000000000..f2b41e2b9e3 --- /dev/null +++ b/pulpcore/tests/unit/test_object_storage.py @@ -0,0 +1,220 @@ +import asyncio +from unittest.mock import AsyncMock, Mock + +import pytest + +from pulpcore._object_storage import AzureStream, S3Stream, get_stream + + +class S3Body: + def __init__(self): + self.read_sizes = [] + self.closed = False + + def read(self, size): + self.read_sizes.append(size) + return b"abc"[:size] + + def close(self): + self.closed = True + + +def test_s3_stream_normalizes_name_and_filters_download_parameters(): + """Preserve S3 locations and read-specific options when bypassing `Storage.open()`. + + The temporary adapter calls the provider directly, so this protects against + losing a configured location prefix or security-related object parameters. + """ + body = S3Body() + client = Mock() + client.get_object.return_value = {"Body": body} + storage = Mock() + storage.location = "prefix" + storage.bucket_name = "bucket" + storage.connection.meta.client = client + storage._normalize_name.side_effect = lambda name: f"prefix/{name}" + storage.get_object_parameters.return_value = { + "SSECustomerAlgorithm": "AES256", + "RequestPayer": "requester", + "VersionId": "version", + "ChecksumMode": "ENABLED", + "ExpectedBucketOwner": "owner", + "ContentType": "not-a-download-argument", + } + + stream = S3Stream(storage, "./artifact/name", 3, 10) + assert stream.open() is stream + storage._normalize_name.assert_called_once_with("artifact/name") + storage.get_object_parameters.assert_called_once_with("artifact/name") + assert client.get_object.call_args.kwargs == { + "SSECustomerAlgorithm": "AES256", + "RequestPayer": "requester", + "VersionId": "version", + "ChecksumMode": "ENABLED", + "ExpectedBucketOwner": "owner", + "Bucket": "bucket", + "Key": "prefix/artifact/name", + "Range": "bytes=3-12", + } + + assert stream.read(2) == b"ab" + stream.close() + assert body.read_sizes == [2] + assert body.closed + + +def test_azure_stream_uses_ranged_chunks_and_bounds_reads(): + """Adapt Azure's provider-sized chunks to bounded reads without changing its range.""" + downloader = Mock() + downloader.chunks.return_value = iter((b"ab", b"cdef", b"gh")) + client = Mock() + client.download_blob.return_value = downloader + storage = Mock() + storage.timeout = 17 + storage.client = client + storage._get_valid_path.side_effect = lambda name: f"prefix/{name}" + + stream = AzureStream(storage, "artifact/name", 3, 8) + assert stream.open() is stream + assert client.download_blob.call_args.kwargs == { + "offset": 3, + "length": 8, + "timeout": 17, + } + assert client.download_blob.call_args.args == ("prefix/artifact/name",) + assert stream.read(3) == b"abc" + assert stream.read(3) == b"def" + assert stream.read(10) == b"gh" + + stream.close() + downloader.close.assert_called_once_with() + + +def test_unknown_storage_class_keeps_file_object_fallback(): + """Keep unsupported backends on the established generic storage code path.""" + assert get_stream("pulpcore.app.models.storage.FileSystem", Mock(), "name", 0, 1) is None + + +def test_response_closes_object_stream_after_write_failure(monkeypatch): + """Release the provider response when the client disconnects during a write.""" + asyncio.run(_test_response_closes_object_stream_after_write_failure(monkeypatch)) + + +async def _test_response_closes_object_stream_after_write_failure(monkeypatch): + from aiohttp.web import StreamResponse + + from pulpcore.responses import ArtifactResponse + + writer = Mock() + writer.write = AsyncMock(side_effect=RuntimeError("client disconnected")) + stream = Mock() + stream.open.return_value = stream + stream.read.return_value = b"abc" + + async def prepare(_self, _request): + return writer + + monkeypatch.setattr(StreamResponse, "prepare", prepare) + response = ArtifactResponse(artifact=Mock(), chunk_size=4) + + with pytest.raises(RuntimeError, match="client disconnected"): + await response._sendfile_object_storage(Mock(), stream, 3) + + stream.close.assert_called_once_with() + + +def test_response_closes_object_stream_after_open_failure(monkeypatch): + """Attempt adapter cleanup even when opening the provider stream raises.""" + asyncio.run(_test_response_closes_object_stream_after_open_failure(monkeypatch)) + + +async def _test_response_closes_object_stream_after_open_failure(monkeypatch): + from aiohttp.web import StreamResponse + + from pulpcore.responses import ArtifactResponse + + writer = Mock() + stream = Mock() + stream.open.side_effect = RuntimeError("provider unavailable") + + async def prepare(_self, _request): + return writer + + monkeypatch.setattr(StreamResponse, "prepare", prepare) + response = ArtifactResponse(artifact=Mock(), chunk_size=4) + + with pytest.raises(RuntimeError, match="provider unavailable"): + await response._sendfile_object_storage(Mock(), stream, 3) + + stream.close.assert_called_once_with() + + +def test_artifact_response_dispatches_s3_stream_without_file_reads(monkeypatch): + """Use the adapter for S3 so django-storages cannot eagerly spool the object.""" + asyncio.run(_test_artifact_response_dispatches_s3_stream_without_file_reads(monkeypatch)) + + +async def _test_artifact_response_dispatches_s3_stream_without_file_reads(monkeypatch): + from pulpcore.responses import ArtifactResponse + + storage = Mock() + domain = Mock(storage_class="storages.backends.s3.S3Storage") + domain.get_storage.return_value = storage + artifact = Mock(pulp_domain=domain) + file_object = Mock() + file_object.name = "artifact/name" + object_stream = Mock() + response = ArtifactResponse(artifact=artifact) + response._sendfile_object_storage = AsyncMock(return_value="writer") + + def get_stream(storage_class, selected_storage, name, offset, count): + assert storage_class == "storages.backends.s3.S3Storage" + assert selected_storage is storage + assert name == "artifact/name" + assert offset == 11 + assert count == 19 + return object_stream + + monkeypatch.setattr("pulpcore.responses.get_stream", get_stream) + + assert await response._sendfile("request", file_object, 11, 19) == "writer" + domain.get_storage.assert_called_once_with() + response._sendfile_object_storage.assert_awaited_once_with("request", object_stream, 19) + file_object.seek.assert_not_called() + file_object.read.assert_not_called() + + +def test_artifact_response_keeps_file_fallback_for_unsupported_storage(monkeypatch): + """Retain the previous seek-and-read behavior for non-adapter storage backends.""" + asyncio.run(_test_artifact_response_keeps_file_fallback_for_unsupported_storage(monkeypatch)) + + +async def _test_artifact_response_keeps_file_fallback_for_unsupported_storage(monkeypatch): + from aiohttp.web import StreamResponse + + from pulpcore.responses import ArtifactResponse + + writer = Mock() + writer.write = AsyncMock() + writer.drain = AsyncMock() + + async def prepare(_self, _request): + return writer + + monkeypatch.setattr(StreamResponse, "prepare", prepare) + monkeypatch.setattr( + "pulpcore.responses.get_stream", Mock(side_effect=AssertionError("adapter selected")) + ) + + domain = Mock(storage_class="pulpcore.app.models.storage.FileSystem") + artifact = Mock(pulp_domain=domain) + file_object = Mock() + file_object.name = "artifact/name" + file_object.read.return_value = b"payload" + response = ArtifactResponse(artifact=artifact, chunk_size=8) + + assert await response._sendfile("request", file_object, 3, 7) is writer + domain.get_storage.assert_not_called() + file_object.seek.assert_called_once_with(3) + file_object.read.assert_called_once_with(7) + writer.write.assert_awaited_once_with(b"payload") From 6fbacef8da657ad5b1bfb19713db086c1e35abbe Mon Sep 17 00:00:00 2001 From: Daniel Alley Date: Mon, 14 Sep 2026 13:51:18 -0400 Subject: [PATCH 2/3] Revert "Stream proxied S3 and Azure artifacts" This reverts commit 2ed7659ed8152866f5489e77e5bad16ba6ae2216. --- CHANGES/7806.bugfix | 1 - CLAUDE.md | 7 - .../functional/api/test_download_policies.py | 74 +----- pulpcore/_object_storage.py | 166 ------------- pulpcore/responses.py | 39 ---- .../api/test_artifact_distribution.py | 4 +- pulpcore/tests/unit/test_object_storage.py | 220 ------------------ 7 files changed, 3 insertions(+), 508 deletions(-) delete mode 100644 CHANGES/7806.bugfix delete mode 100644 pulpcore/_object_storage.py delete mode 100644 pulpcore/tests/unit/test_object_storage.py diff --git a/CHANGES/7806.bugfix b/CHANGES/7806.bugfix deleted file mode 100644 index a1a7a1eae7c..00000000000 --- a/CHANGES/7806.bugfix +++ /dev/null @@ -1 +0,0 @@ -Resolve out-of-memory conditions when serving large artifacts from S3/Azure when REDIRECT_TO_OBJECT_STORAGE=False. diff --git a/CLAUDE.md b/CLAUDE.md index 266d3ad43b4..e775a8afcf5 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -38,13 +38,6 @@ pulpcore & pulp-file functional tests require both client bindings to be install **Always** use the `oci-env` to run the functional and unit tests. -The active OCI profile may mount this checkout at `/src/pulpcore`, while `oci-env test -p core` -expects `/src/core`. If that wrapper cannot find `/src/core`, run the focused test with -`oci-env exec pytest /src/pulpcore/` instead. - -`ArtifactDistribution` is matched within its `pulp_domain`. Its generated artifact URL includes -the artifact domain, so test-created artifact distributions must use that same domain. - ## Modifying template_config.yml Use the `plugin-template` tool after any changes made to `template_config.yml`. diff --git a/pulp_file/tests/functional/api/test_download_policies.py b/pulp_file/tests/functional/api/test_download_policies.py index 4f8cae6cd6f..46020243f1b 100644 --- a/pulp_file/tests/functional/api/test_download_policies.py +++ b/pulp_file/tests/functional/api/test_download_policies.py @@ -7,15 +7,10 @@ from urllib.parse import urljoin import pytest -import requests from aiohttp.client_exceptions import ClientResponseError from bs4 import BeautifulSoup -from pulpcore.client.pulp_file import ( - FileFilePublication, - FileRepositorySyncURL, - RepositoryAddRemoveContent, -) +from pulpcore.client.pulp_file import FileFilePublication, FileRepositorySyncURL from pulpcore.tests.functional.utils import download_file, get_files_in_manifest OBJECT_STORAGES = ( @@ -42,73 +37,6 @@ def _do_range_request_download_and_assert(url, range_header, expected_bytes): ) -@pytest.mark.parametrize( - "storage_class", - ( - pytest.param("storages.backends.s3.S3Storage", id="s3"), - pytest.param("storages.backends.s3boto3.S3Boto3Storage", id="s3boto3"), - pytest.param("storages.backends.azure_storage.AzureStorage", id="azure"), - ), -) -def test_proxied_object_storage_artifact_streaming( - storage_class, - pulp_settings, - domain_factory, - random_artifact_factory, - file_bindings, - file_repository_factory, - file_publication_factory, - file_distribution_factory, - distribution_base_url, - gen_object_with_cleanup, - monitor_task, -): - """Proxy object-storage content served by a distribution without redirecting.""" - if pulp_settings.STORAGES["default"]["BACKEND"] != storage_class: - pytest.skip("The functional environment does not provide this object-storage configuration") - - domain = domain_factory(storage_class=storage_class, redirect_to_object_storage=False) - artifact = random_artifact_factory(pulp_domain=domain.name, size=32) - content = gen_object_with_cleanup( - file_bindings.ContentFilesApi, - artifact=artifact.pulp_href, - relative_path=str(uuid.uuid4()), - pulp_domain=domain.name, - ) - repository = file_repository_factory(pulp_domain=domain.name) - monitor_task( - file_bindings.RepositoriesFileApi.modify( - repository.pulp_href, - RepositoryAddRemoveContent(add_content_units=[content.pulp_href]), - ).task - ) - publication = file_publication_factory( - pulp_domain=domain.name, - repository=repository.pulp_href, - ) - distribution = file_distribution_factory( - pulp_domain=domain.name, - publication=publication.pulp_href, - ) - content_url = urljoin(distribution_base_url(distribution.base_url), content.relative_path) - - full_response = requests.get(content_url, allow_redirects=False) - assert full_response.status_code == requests.codes.ok - assert not full_response.is_redirect - assert hashlib.sha256(full_response.content).hexdigest() == artifact.sha256 - assert full_response.headers["Content-Length"] == str(len(full_response.content)) - assert full_response.headers["Accept-Ranges"] == "bytes" - - range_response = requests.get( - content_url, headers={"Range": "bytes=1-4"}, allow_redirects=False - ) - assert range_response.status_code == requests.codes.partial_content - assert not range_response.is_redirect - assert range_response.content == full_response.content[1:5] - assert range_response.headers["Content-Length"] == "4" - assert range_response.headers["Content-Range"] == f"bytes 1-4/{len(full_response.content)}" - - @pytest.mark.parallel @pytest.mark.parametrize("download_policy", ["immediate", "on_demand", "streamed"]) def test_download_policy( diff --git a/pulpcore/_object_storage.py b/pulpcore/_object_storage.py deleted file mode 100644 index 1af85791151..00000000000 --- a/pulpcore/_object_storage.py +++ /dev/null @@ -1,166 +0,0 @@ -"""Temporary private streaming adapters for object-storage response bodies. - -django-storages file objects are seekable Django file objects. Opening one for -reading downloads the complete object into a temporary file, which is unsuitable -for the content app's bounded response loop. These adapters expose only the -synchronous `open/read/close` interface needed by :mod:`pulpcore.responses`. - -Remove these adapters when django-storages provides a stable, cross-backend, -streaming-read API that Pulp can use instead. -""" - -from typing import Any - -S3_STORAGE_CLASSES = frozenset( - ( - "storages.backends.s3.S3Storage", - "storages.backends.s3boto3.S3Boto3Storage", - ) -) -AZURE_STORAGE_CLASSES = frozenset(("storages.backends.azure_storage.AzureStorage",)) -STREAMING_STORAGE_CLASSES = S3_STORAGE_CLASSES | AZURE_STORAGE_CLASSES - -# These are the read arguments accepted by boto3's get_object operation and by -# s3transfer's download argument filter. The import is optional because pulpcore -# can be installed without the S3 extra. -try: - from s3transfer.constants import ALLOWED_DOWNLOAD_ARGS as _S3_DOWNLOAD_ARGS -except ImportError: # pragma: no cover - exercised only without the S3 extra - _S3_DOWNLOAD_ARGS = ( - "ChecksumMode", - "ExpectedBucketOwner", - "IfMatch", - "IfModifiedSince", - "IfNoneMatch", - "IfUnmodifiedSince", - "PartNumber", - "RequestPayer", - "SSECustomerAlgorithm", - "SSECustomerKey", - "SSECustomerKeyMD5", - "VersionId", - ) - - -def _filter_s3_download_params(params: dict[str, Any]) -> dict[str, Any]: - """Retain only configured values accepted by boto3 `get_object`.""" - - return {key: value for key, value in params.items() if key in _S3_DOWNLOAD_ARGS} - - -def _clean_s3_name(name): - """Match django-storages' logical-name processing before adding `location`.""" - - try: - from storages.utils import clean_name - except ImportError: # pragma: no cover - S3 storage requires django-storages - import posixpath - - cleaned_name = posixpath.normpath(name).replace("\\", "/") - if name.endswith("/") and not cleaned_name.endswith("/"): - cleaned_name += "/" - return "" if cleaned_name == "." else cleaned_name - return clean_name(name) - - -class S3Stream: - """A ranged reader backed by the configured django-storages S3 client.""" - - def __init__(self, storage, name: str, offset: int, count: int): - self.storage = storage - self.name = name - self.offset = offset - self.count = count - self.body = None - - def open(self): - name = _clean_s3_name(self.name) - key = self.storage._normalize_name(name) - # `get_object_parameters` receives the logical name; `Key` includes - # the backend's location prefix just as django-storages' `open` does. - params = _filter_s3_download_params(self.storage.get_object_parameters(name)) - params.update( - Bucket=self.storage.bucket_name, - Key=key, - Range=f"bytes={self.offset}-{self.offset + self.count - 1}", - ) - response = self.storage.connection.meta.client.get_object(**params) - self.body = response["Body"] - return self - - def read(self, size: int) -> bytes: - return self.body.read(size) - - def close(self): - if self.body is not None: - self.body.close() - self.body = None - - -class AzureStream: - """A ranged reader backed by an Azure `StorageStreamDownloader`.""" - - def __init__(self, storage, name: str, offset: int, count: int): - self.storage = storage - self.name = name - self.offset = offset - self.count = count - self.downloader = None - self.chunks = None - self.pending = b"" - - def open(self): - path = self.storage._get_valid_path(self.name) - self.downloader = self.storage.client.download_blob( - path, - offset=self.offset, - length=self.count, - timeout=self.storage.timeout, - ) - self.chunks = iter(self.downloader.chunks()) - return self - - def read(self, size: int) -> bytes: - # Azure controls the size of values yielded by `chunks()`. Keep at - # most one such value pending while presenting Pulp's smaller read size. - while len(self.pending) < size: - try: - self.pending += next(self.chunks) - except StopIteration: - break - - chunk = self.pending[:size] - self.pending = self.pending[size:] - return chunk - - def close(self): - if self.downloader is None: - return - - close = getattr(self.downloader, "close", None) - if close is None: - response = getattr(self.downloader, "_response", None) - for response_part in ( - response, - getattr(response, "http_response", None), - getattr(getattr(response, "http_response", None), "internal_response", None), - ): - close = getattr(response_part, "close", None) - if close is not None: - close() - break - else: - close() - self.downloader = None - self.chunks = None - self.pending = b"" - - -def get_stream(storage_class: str, storage, name: str, offset: int, count: int): - """Return an adapter for a supported storage class, or `None` for fallback.""" - - if storage_class in S3_STORAGE_CLASSES: - return S3Stream(storage, name, offset, count) - if storage_class in AZURE_STORAGE_CLASSES: - return AzureStream(storage, name, offset, count) - return None diff --git a/pulpcore/responses.py b/pulpcore/responses.py index 81f5c40e652..1b1fac62a0d 100644 --- a/pulpcore/responses.py +++ b/pulpcore/responses.py @@ -7,7 +7,6 @@ HTTPRequestRangeNotSatisfiable, ) -from pulpcore._object_storage import STREAMING_STORAGE_CLASSES, get_stream from pulpcore.app.models import Artifact @@ -34,17 +33,6 @@ def __init__( self._chunk_size = chunk_size async def _sendfile(self, request, fobj, offset, count): - storage_class = self._artifact.pulp_domain.storage_class - if storage_class in STREAMING_STORAGE_CLASSES: - object_stream = get_stream( - storage_class, - self._artifact.pulp_domain.get_storage(), - fobj.name, - offset, - count, - ) - return await self._sendfile_object_storage(request, object_stream, count) - # To keep memory usage low, fobj is transferred in chunks # controlled by the constructor's chunk_size argument. @@ -66,33 +54,6 @@ async def _sendfile(self, request, fobj, offset, count): await writer.drain() return writer - async def _sendfile_object_storage(self, request, object_stream, count): - """Write a blocking provider stream without blocking the content-app loop. - - Object storage SDKs are synchronous. Adapter methods therefore run in a - worker thread, while `writer.write` maintains aiohttp backpressure. The - adapter is closed even if opening it or writing its response fails. - """ - writer = await super().prepare(request) - assert writer is not None - - stream = object_stream - try: - stream = await asyncio.to_thread(object_stream.open) - remaining = count - while remaining: - chunk = await asyncio.to_thread(stream.read, min(self._chunk_size, remaining)) - if not chunk: - break - if len(chunk) > remaining: - chunk = chunk[:remaining] - await writer.write(chunk) - remaining -= len(chunk) - await writer.drain() - return writer - finally: - await asyncio.to_thread(stream.close) - async def prepare(self, request): if self._artifact is None: self._artifact = await Artifact.objects.select_related("pulp_domain").aget( diff --git a/pulpcore/tests/functional/api/test_artifact_distribution.py b/pulpcore/tests/functional/api/test_artifact_distribution.py index dc493981af6..14b6dc7e862 100644 --- a/pulpcore/tests/functional/api/test_artifact_distribution.py +++ b/pulpcore/tests/functional/api/test_artifact_distribution.py @@ -11,7 +11,7 @@ ) -def test_artifact_distribution(random_artifact, pulp_settings, distribution_base_url): +def test_artifact_distribution(random_artifact, pulp_settings): settings = pulp_settings artifact_uuid = random_artifact.pulp_href.split("/")[-2] @@ -22,7 +22,7 @@ def test_artifact_distribution(random_artifact, pulp_settings, distribution_base ) process = subprocess.run(["pulpcore-manager", "shell", "-c", commands], capture_output=True) assert process.returncode == 0 - artifact_url = distribution_base_url(process.stdout.decode().strip()) + artifact_url = process.stdout.decode().strip() response = requests.get(artifact_url) response.raise_for_status() diff --git a/pulpcore/tests/unit/test_object_storage.py b/pulpcore/tests/unit/test_object_storage.py deleted file mode 100644 index f2b41e2b9e3..00000000000 --- a/pulpcore/tests/unit/test_object_storage.py +++ /dev/null @@ -1,220 +0,0 @@ -import asyncio -from unittest.mock import AsyncMock, Mock - -import pytest - -from pulpcore._object_storage import AzureStream, S3Stream, get_stream - - -class S3Body: - def __init__(self): - self.read_sizes = [] - self.closed = False - - def read(self, size): - self.read_sizes.append(size) - return b"abc"[:size] - - def close(self): - self.closed = True - - -def test_s3_stream_normalizes_name_and_filters_download_parameters(): - """Preserve S3 locations and read-specific options when bypassing `Storage.open()`. - - The temporary adapter calls the provider directly, so this protects against - losing a configured location prefix or security-related object parameters. - """ - body = S3Body() - client = Mock() - client.get_object.return_value = {"Body": body} - storage = Mock() - storage.location = "prefix" - storage.bucket_name = "bucket" - storage.connection.meta.client = client - storage._normalize_name.side_effect = lambda name: f"prefix/{name}" - storage.get_object_parameters.return_value = { - "SSECustomerAlgorithm": "AES256", - "RequestPayer": "requester", - "VersionId": "version", - "ChecksumMode": "ENABLED", - "ExpectedBucketOwner": "owner", - "ContentType": "not-a-download-argument", - } - - stream = S3Stream(storage, "./artifact/name", 3, 10) - assert stream.open() is stream - storage._normalize_name.assert_called_once_with("artifact/name") - storage.get_object_parameters.assert_called_once_with("artifact/name") - assert client.get_object.call_args.kwargs == { - "SSECustomerAlgorithm": "AES256", - "RequestPayer": "requester", - "VersionId": "version", - "ChecksumMode": "ENABLED", - "ExpectedBucketOwner": "owner", - "Bucket": "bucket", - "Key": "prefix/artifact/name", - "Range": "bytes=3-12", - } - - assert stream.read(2) == b"ab" - stream.close() - assert body.read_sizes == [2] - assert body.closed - - -def test_azure_stream_uses_ranged_chunks_and_bounds_reads(): - """Adapt Azure's provider-sized chunks to bounded reads without changing its range.""" - downloader = Mock() - downloader.chunks.return_value = iter((b"ab", b"cdef", b"gh")) - client = Mock() - client.download_blob.return_value = downloader - storage = Mock() - storage.timeout = 17 - storage.client = client - storage._get_valid_path.side_effect = lambda name: f"prefix/{name}" - - stream = AzureStream(storage, "artifact/name", 3, 8) - assert stream.open() is stream - assert client.download_blob.call_args.kwargs == { - "offset": 3, - "length": 8, - "timeout": 17, - } - assert client.download_blob.call_args.args == ("prefix/artifact/name",) - assert stream.read(3) == b"abc" - assert stream.read(3) == b"def" - assert stream.read(10) == b"gh" - - stream.close() - downloader.close.assert_called_once_with() - - -def test_unknown_storage_class_keeps_file_object_fallback(): - """Keep unsupported backends on the established generic storage code path.""" - assert get_stream("pulpcore.app.models.storage.FileSystem", Mock(), "name", 0, 1) is None - - -def test_response_closes_object_stream_after_write_failure(monkeypatch): - """Release the provider response when the client disconnects during a write.""" - asyncio.run(_test_response_closes_object_stream_after_write_failure(monkeypatch)) - - -async def _test_response_closes_object_stream_after_write_failure(monkeypatch): - from aiohttp.web import StreamResponse - - from pulpcore.responses import ArtifactResponse - - writer = Mock() - writer.write = AsyncMock(side_effect=RuntimeError("client disconnected")) - stream = Mock() - stream.open.return_value = stream - stream.read.return_value = b"abc" - - async def prepare(_self, _request): - return writer - - monkeypatch.setattr(StreamResponse, "prepare", prepare) - response = ArtifactResponse(artifact=Mock(), chunk_size=4) - - with pytest.raises(RuntimeError, match="client disconnected"): - await response._sendfile_object_storage(Mock(), stream, 3) - - stream.close.assert_called_once_with() - - -def test_response_closes_object_stream_after_open_failure(monkeypatch): - """Attempt adapter cleanup even when opening the provider stream raises.""" - asyncio.run(_test_response_closes_object_stream_after_open_failure(monkeypatch)) - - -async def _test_response_closes_object_stream_after_open_failure(monkeypatch): - from aiohttp.web import StreamResponse - - from pulpcore.responses import ArtifactResponse - - writer = Mock() - stream = Mock() - stream.open.side_effect = RuntimeError("provider unavailable") - - async def prepare(_self, _request): - return writer - - monkeypatch.setattr(StreamResponse, "prepare", prepare) - response = ArtifactResponse(artifact=Mock(), chunk_size=4) - - with pytest.raises(RuntimeError, match="provider unavailable"): - await response._sendfile_object_storage(Mock(), stream, 3) - - stream.close.assert_called_once_with() - - -def test_artifact_response_dispatches_s3_stream_without_file_reads(monkeypatch): - """Use the adapter for S3 so django-storages cannot eagerly spool the object.""" - asyncio.run(_test_artifact_response_dispatches_s3_stream_without_file_reads(monkeypatch)) - - -async def _test_artifact_response_dispatches_s3_stream_without_file_reads(monkeypatch): - from pulpcore.responses import ArtifactResponse - - storage = Mock() - domain = Mock(storage_class="storages.backends.s3.S3Storage") - domain.get_storage.return_value = storage - artifact = Mock(pulp_domain=domain) - file_object = Mock() - file_object.name = "artifact/name" - object_stream = Mock() - response = ArtifactResponse(artifact=artifact) - response._sendfile_object_storage = AsyncMock(return_value="writer") - - def get_stream(storage_class, selected_storage, name, offset, count): - assert storage_class == "storages.backends.s3.S3Storage" - assert selected_storage is storage - assert name == "artifact/name" - assert offset == 11 - assert count == 19 - return object_stream - - monkeypatch.setattr("pulpcore.responses.get_stream", get_stream) - - assert await response._sendfile("request", file_object, 11, 19) == "writer" - domain.get_storage.assert_called_once_with() - response._sendfile_object_storage.assert_awaited_once_with("request", object_stream, 19) - file_object.seek.assert_not_called() - file_object.read.assert_not_called() - - -def test_artifact_response_keeps_file_fallback_for_unsupported_storage(monkeypatch): - """Retain the previous seek-and-read behavior for non-adapter storage backends.""" - asyncio.run(_test_artifact_response_keeps_file_fallback_for_unsupported_storage(monkeypatch)) - - -async def _test_artifact_response_keeps_file_fallback_for_unsupported_storage(monkeypatch): - from aiohttp.web import StreamResponse - - from pulpcore.responses import ArtifactResponse - - writer = Mock() - writer.write = AsyncMock() - writer.drain = AsyncMock() - - async def prepare(_self, _request): - return writer - - monkeypatch.setattr(StreamResponse, "prepare", prepare) - monkeypatch.setattr( - "pulpcore.responses.get_stream", Mock(side_effect=AssertionError("adapter selected")) - ) - - domain = Mock(storage_class="pulpcore.app.models.storage.FileSystem") - artifact = Mock(pulp_domain=domain) - file_object = Mock() - file_object.name = "artifact/name" - file_object.read.return_value = b"payload" - response = ArtifactResponse(artifact=artifact, chunk_size=8) - - assert await response._sendfile("request", file_object, 3, 7) is writer - domain.get_storage.assert_not_called() - file_object.seek.assert_called_once_with(3) - file_object.read.assert_called_once_with(7) - writer.write.assert_awaited_once_with(b"payload") From 2b3471d253cf14133176f373b45b8e2e15a8296f Mon Sep 17 00:00:00 2001 From: Daniel Alley Date: Mon, 14 Sep 2026 14:10:46 -0400 Subject: [PATCH 3/3] Stream proxied object-storage artifacts via django-storages Use the forked django-storages open_stream API for supported remote storage backends so Pulp can proxy large artifacts without opening their eager, seekable file objects. Preserve the existing generic storage fallback and HTTP range behavior. closes #7806 Assisted By: Codex (GPT-5) Terra 5.6 --- CHANGES/7806.bugfix | 1 + .../functional/api/test_download_policies.py | 75 ++++++++- pulpcore/responses.py | 66 ++++++++ pulpcore/tests/unit/test_storage_streaming.py | 159 ++++++++++++++++++ pyproject.toml | 8 +- 5 files changed, 304 insertions(+), 5 deletions(-) create mode 100644 CHANGES/7806.bugfix create mode 100644 pulpcore/tests/unit/test_storage_streaming.py diff --git a/CHANGES/7806.bugfix b/CHANGES/7806.bugfix new file mode 100644 index 00000000000..357015a17cc --- /dev/null +++ b/CHANGES/7806.bugfix @@ -0,0 +1 @@ +Fixed proxying large S3 and Azure artifacts without buffering their complete contents. diff --git a/pulp_file/tests/functional/api/test_download_policies.py b/pulp_file/tests/functional/api/test_download_policies.py index 46020243f1b..882457a008c 100644 --- a/pulp_file/tests/functional/api/test_download_policies.py +++ b/pulp_file/tests/functional/api/test_download_policies.py @@ -7,10 +7,15 @@ from urllib.parse import urljoin import pytest +import requests from aiohttp.client_exceptions import ClientResponseError from bs4 import BeautifulSoup -from pulpcore.client.pulp_file import FileFilePublication, FileRepositorySyncURL +from pulpcore.client.pulp_file import ( + FileFilePublication, + FileRepositorySyncURL, + RepositoryAddRemoveContent, +) from pulpcore.tests.functional.utils import download_file, get_files_in_manifest OBJECT_STORAGES = ( @@ -37,6 +42,74 @@ def _do_range_request_download_and_assert(url, range_header, expected_bytes): ) +@pytest.mark.parametrize( + "storage_class", + ( + pytest.param("storages.backends.s3.S3Storage", id="s3"), + pytest.param("storages.backends.s3boto3.S3Boto3Storage", id="s3boto3"), + pytest.param("storages.backends.azure_storage.AzureStorage", id="azure"), + ), +) +def test_proxied_object_storage_artifact_streaming( + storage_class, + pulp_settings, + domain_factory, + random_artifact_factory, + file_bindings, + file_repository_factory, + file_publication_factory, + file_distribution_factory, + distribution_base_url, + gen_object_with_cleanup, + monitor_task, +): + """Serve full and ranged object-storage content through Pulp without redirects.""" + + if pulp_settings.STORAGES["default"]["BACKEND"] != storage_class: + pytest.skip("The functional environment does not provide this object-storage configuration") + + domain = domain_factory(storage_class=storage_class, redirect_to_object_storage=False) + artifact = random_artifact_factory(pulp_domain=domain.name, size=32) + content = gen_object_with_cleanup( + file_bindings.ContentFilesApi, + artifact=artifact.pulp_href, + relative_path=str(uuid.uuid4()), + pulp_domain=domain.name, + ) + repository = file_repository_factory(pulp_domain=domain.name) + monitor_task( + file_bindings.RepositoriesFileApi.modify( + repository.pulp_href, + RepositoryAddRemoveContent(add_content_units=[content.pulp_href]), + ).task + ) + publication = file_publication_factory( + pulp_domain=domain.name, + repository=repository.pulp_href, + ) + distribution = file_distribution_factory( + pulp_domain=domain.name, + publication=publication.pulp_href, + ) + content_url = urljoin(distribution_base_url(distribution.base_url), content.relative_path) + + full_response = requests.get(content_url, allow_redirects=False) + assert full_response.status_code == requests.codes.ok + assert not full_response.is_redirect + assert hashlib.sha256(full_response.content).hexdigest() == artifact.sha256 + assert full_response.headers["Content-Length"] == str(len(full_response.content)) + assert full_response.headers["Accept-Ranges"] == "bytes" + + range_response = requests.get( + content_url, headers={"Range": "bytes=1-4"}, allow_redirects=False + ) + assert range_response.status_code == requests.codes.partial_content + assert not range_response.is_redirect + assert range_response.content == full_response.content[1:5] + assert range_response.headers["Content-Length"] == "4" + assert range_response.headers["Content-Range"] == f"bytes 1-4/{len(full_response.content)}" + + @pytest.mark.parallel @pytest.mark.parametrize("download_policy", ["immediate", "on_demand", "streamed"]) def test_download_policy( diff --git a/pulpcore/responses.py b/pulpcore/responses.py index 1b1fac62a0d..8f45a50dd85 100644 --- a/pulpcore/responses.py +++ b/pulpcore/responses.py @@ -1,4 +1,5 @@ import asyncio +from contextlib import asynccontextmanager from aiohttp import hdrs from aiohttp.web import StreamResponse @@ -9,6 +10,15 @@ from pulpcore.app.models import Artifact +STREAMING_STORAGE_CLASSES = frozenset( + ( + "storages.backends.s3.S3Storage", + "storages.backends.s3boto3.S3Boto3Storage", + "storages.backends.azure_storage.AzureStorage", + "storages.backends.gcloud.GoogleCloudStorage", + ) +) + class ArtifactResponse(StreamResponse): """A response object can be used to send artifacts.""" @@ -33,6 +43,12 @@ def __init__( self._chunk_size = chunk_size async def _sendfile(self, request, fobj, offset, count): + if self._artifact.pulp_domain.storage_class in STREAMING_STORAGE_CLASSES: + storage = self._artifact.pulp_domain.get_storage() + return await self._sendfile_storage_stream( + request, storage.open_stream, fobj.name, offset, count + ) + # To keep memory usage low, fobj is transferred in chunks # controlled by the constructor's chunk_size argument. @@ -54,6 +70,56 @@ async def _sendfile(self, request, fobj, offset, count): await writer.drain() return writer + @staticmethod + @asynccontextmanager + async def _storage_stream(stream_opener, name, offset, count): + """Bridge django-storages' synchronous streaming context manager. + + The django-storages fork supplies ``open_stream()`` on the remote + backends selected above. Its provider I/O, including opening and + closing the context, must stay off the content app's event loop. + Propagate exception details to the synchronous context manager so it + retains normal ``with`` semantics. Once django-storages releases this + API upstream, replace the forked dependency with that release. + """ + + stream_context = await asyncio.to_thread(stream_opener, name, start=offset, length=count) + stream = await asyncio.to_thread(stream_context.__enter__) + try: + yield stream + except BaseException as exc: + if not await asyncio.to_thread( + stream_context.__exit__, type(exc), exc, exc.__traceback__ + ): + raise + else: + await asyncio.to_thread(stream_context.__exit__, None, None, None) + + async def _sendfile_storage_stream(self, request, stream_opener, name, offset, count): + """Write an ``open_stream()`` response in bounded chunks. + + The storage API is synchronous while aiohttp writes are asynchronous. + Bound each provider read to the remaining HTTP range and use aiohttp's + normal backpressure for every write. + """ + + writer = await super().prepare(request) + assert writer is not None + + async with self._storage_stream(stream_opener, name, offset, count) as stream: + remaining = count + while remaining: + chunk = await asyncio.to_thread(stream.read, min(self._chunk_size, remaining)) + if not chunk: + break + if len(chunk) > remaining: + chunk = chunk[:remaining] + await writer.write(chunk) + remaining -= len(chunk) + + await writer.drain() + return writer + async def prepare(self, request): if self._artifact is None: self._artifact = await Artifact.objects.select_related("pulp_domain").aget( diff --git a/pulpcore/tests/unit/test_storage_streaming.py b/pulpcore/tests/unit/test_storage_streaming.py new file mode 100644 index 00000000000..66e16cca099 --- /dev/null +++ b/pulpcore/tests/unit/test_storage_streaming.py @@ -0,0 +1,159 @@ +"""Tests for ArtifactResponse's opt-in django-storages streaming path.""" + +import asyncio +from unittest.mock import AsyncMock, MagicMock, Mock, call + +import pytest + + +@pytest.mark.parametrize( + "storage_class", + ( + "storages.backends.s3.S3Storage", + "storages.backends.s3boto3.S3Boto3Storage", + "storages.backends.azure_storage.AzureStorage", + "storages.backends.gcloud.GoogleCloudStorage", + ), +) +def test_artifact_response_uses_open_stream_for_a_domain_storage(storage_class, monkeypatch): + """Use the explicit storage API instead of opening a seekable temporary file.""" + + asyncio.run( + _test_artifact_response_uses_open_stream_for_a_domain_storage(storage_class, monkeypatch) + ) + + +async def _test_artifact_response_uses_open_stream_for_a_domain_storage(storage_class, monkeypatch): + from aiohttp.web import StreamResponse + + from pulpcore.responses import ArtifactResponse + + writer = Mock() + writer.write = AsyncMock() + writer.drain = AsyncMock() + stream = Mock() + stream.read.side_effect = (b"abc", b"def", b"") + stream_context = MagicMock() + stream_context.__enter__.return_value = stream + stream_opener = Mock(return_value=stream_context) + storage = Mock(open_stream=stream_opener) + domain = Mock(storage_class=storage_class) + domain.get_storage.return_value = storage + artifact = Mock(pulp_domain=domain) + file_object = Mock() + file_object.name = "artifact/name" + + async def prepare(_self, _request): + return writer + + monkeypatch.setattr(StreamResponse, "prepare", prepare) + response = ArtifactResponse(artifact=artifact, chunk_size=3) + + assert await response._sendfile("request", file_object, 11, 6) is writer + stream_opener.assert_called_once_with("artifact/name", start=11, length=6) + stream.read.assert_has_calls([call(3), call(3)]) + writer.write.assert_has_awaits([call(b"abc"), call(b"def")]) + stream_context.__exit__.assert_called_once_with(None, None, None) + file_object.seek.assert_not_called() + file_object.read.assert_not_called() + + +def test_artifact_response_limits_storage_stream_to_http_range(monkeypatch): + """Never send bytes beyond the headers' selected HTTP range.""" + + asyncio.run(_test_artifact_response_limits_storage_stream_to_http_range(monkeypatch)) + + +async def _test_artifact_response_limits_storage_stream_to_http_range(monkeypatch): + from aiohttp.web import StreamResponse + + from pulpcore.responses import ArtifactResponse + + writer = Mock() + writer.write = AsyncMock() + writer.drain = AsyncMock() + stream = Mock() + stream.read.return_value = b"provider returned too much" + stream_context = MagicMock() + stream_context.__enter__.return_value = stream + + async def prepare(_self, _request): + return writer + + monkeypatch.setattr(StreamResponse, "prepare", prepare) + response = ArtifactResponse(artifact=Mock(), chunk_size=10) + + assert ( + await response._sendfile_storage_stream( + "request", Mock(return_value=stream_context), "artifact/name", 0, 3 + ) + is writer + ) + writer.write.assert_awaited_once_with(b"pro") + stream_context.__exit__.assert_called_once_with(None, None, None) + + +def test_artifact_response_closes_stream_context_after_write_error(monkeypatch): + """Give the storage context the exception needed to release provider resources.""" + + asyncio.run(_test_artifact_response_closes_stream_context_after_write_error(monkeypatch)) + + +async def _test_artifact_response_closes_stream_context_after_write_error(monkeypatch): + from aiohttp.web import StreamResponse + + from pulpcore.responses import ArtifactResponse + + writer = Mock() + writer.write = AsyncMock(side_effect=RuntimeError("client disconnected")) + stream = Mock() + stream.read.return_value = b"abc" + stream_context = MagicMock() + stream_context.__enter__.return_value = stream + + async def prepare(_self, _request): + return writer + + monkeypatch.setattr(StreamResponse, "prepare", prepare) + response = ArtifactResponse(artifact=Mock(), chunk_size=3) + + with pytest.raises(RuntimeError, match="client disconnected"): + await response._sendfile_storage_stream( + "request", Mock(return_value=stream_context), "artifact/name", 0, 3 + ) + + assert stream_context.__exit__.call_args.args[0] is RuntimeError + assert str(stream_context.__exit__.call_args.args[1]) == "client disconnected" + + +def test_artifact_response_keeps_file_fallback_without_open_stream(monkeypatch): + """Retain existing storage-file serving for filesystems and unsupported backends.""" + + asyncio.run(_test_artifact_response_keeps_file_fallback_without_open_stream(monkeypatch)) + + +async def _test_artifact_response_keeps_file_fallback_without_open_stream(monkeypatch): + from aiohttp.web import StreamResponse + + from pulpcore.responses import ArtifactResponse + + writer = Mock() + writer.write = AsyncMock() + writer.drain = AsyncMock() + storage = Mock(spec=[]) + domain = Mock(storage_class="pulpcore.app.models.storage.FileSystem") + domain.get_storage.return_value = storage + artifact = Mock(pulp_domain=domain) + file_object = Mock() + file_object.read.return_value = b"payload" + + async def prepare(_self, _request): + return writer + + monkeypatch.setattr(StreamResponse, "prepare", prepare) + response = ArtifactResponse(artifact=artifact, chunk_size=8) + + assert await response._sendfile("request", file_object, 3, 7) is writer + file_object.seek.assert_called_once_with(3) + file_object.read.assert_called_once_with(7) + writer.write.assert_awaited_once_with(b"payload") diff --git a/pyproject.toml b/pyproject.toml index e242c59ab3b..3fc29f6d7d7 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -67,10 +67,10 @@ dependencies = [ ] [project.optional-dependencies] -sftp = ["django-storages[sftp]==1.14.6"] -s3 = ["django-storages[boto3]==1.14.6"] -google = ["django-storages[google]==1.14.6"] -azure = ["django-storages[azure]==1.14.6"] +sftp = ["django-storages[sftp] @ git+https://github.com/dralley/django-storages.git@open-stream"] +s3 = ["django-storages[boto3] @ git+https://github.com/dralley/django-storages.git@open-stream"] +google = ["django-storages[google] @ git+https://github.com/dralley/django-storages.git@open-stream"] +azure = ["django-storages[azure] @ git+https://github.com/dralley/django-storages.git@open-stream"] prometheus = ["django-prometheus"] saml2 = ["djangosaml2>=1.12.0,<1.13"] kafka = [