test_iter_content.py 9.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293
  1. """
  2. Provider-agnostic unit tests for chunked streaming reads
  3. (``BucketObject.iter_content`` / ``save_content``).
  4. Every provider streams object content through ``iter_content``, but only the
  5. configured test provider gets exercised by the object store service suite. The
  6. chunk-size contract those implementations have to honour is pinned here
  7. against in-memory fakes so it has coverage in CI without cloud credentials.
  8. """
  9. import unittest
  10. from io import BytesIO
  11. from cloudbridge.base.resources import BaseBucketObject
  12. from cloudbridge.interfaces.exceptions import InvalidValueException
  13. class _FakeProvider:
  14. def __init__(self, config=None, bucket_objects=None):
  15. self._config = config or {}
  16. self.storage = _FakeStorage(bucket_objects)
  17. def _get_config_value(self, key, default_value=None):
  18. return self._config.get(key, default_value)
  19. class _FakeStorage:
  20. def __init__(self, bucket_objects=None):
  21. self._bucket_objects = bucket_objects
  22. class _NoRangeService:
  23. def download_range(self, bucket, name, offset, length):
  24. raise AssertionError("no range should have been requested")
  25. class _FakeAzureContainer:
  26. """Stands in for AzureBucket: AzureBucketObject reaches its blob client
  27. through ``container._bucket.get_blob_client(name)``."""
  28. def __init__(self, blob_client):
  29. self._bucket = self
  30. self._blob_client = blob_client
  31. def get_blob_client(self, name):
  32. return self._blob_client
  33. class _FakeBlobProperties:
  34. def __init__(self, name):
  35. self.name = name
  36. class _FakeSwiftContainer:
  37. name = "bucket"
  38. class _StreamingObject(BaseBucketObject):
  39. """A BaseBucketObject that streams from an in-memory buffer."""
  40. def __init__(self, provider, content):
  41. super(_StreamingObject, self).__init__(provider)
  42. self._content = content
  43. self.chunk_sizes_seen = []
  44. @property
  45. def id(self):
  46. return "obj"
  47. @property
  48. def name(self):
  49. return "obj"
  50. @property
  51. def size(self):
  52. return len(self._content)
  53. @property
  54. def bucket(self):
  55. return "BUCKET"
  56. def iter_content(self, chunk_size=None):
  57. size = self._iter_chunk_size(chunk_size)
  58. self.chunk_sizes_seen.append(size)
  59. return (self._content[i:i + size]
  60. for i in range(0, len(self._content), size))
  61. class IterChunkSizeTestCase(unittest.TestCase):
  62. """The resolver that turns an optional chunk_size into a concrete one."""
  63. def _obj(self, config=None, content=b""):
  64. return _StreamingObject(_FakeProvider(config), content)
  65. def test_defaults_to_class_constant_when_unset(self):
  66. obj = self._obj()
  67. self.assertEqual(obj._iter_chunk_size(),
  68. BaseBucketObject.CB_ITER_CHUNK_SIZE)
  69. def test_default_is_one_mebibyte(self):
  70. # Large enough that per-chunk overhead disappears against network
  71. # throughput, small enough to stay cheap per concurrent stream.
  72. self.assertEqual(BaseBucketObject.CB_ITER_CHUNK_SIZE, 1024 * 1024)
  73. def test_provider_config_overrides_class_constant(self):
  74. obj = self._obj({'iter_chunk_size': 8192})
  75. self.assertEqual(obj._iter_chunk_size(), 8192)
  76. def test_explicit_argument_overrides_provider_config(self):
  77. obj = self._obj({'iter_chunk_size': 8192})
  78. self.assertEqual(obj._iter_chunk_size(4096), 4096)
  79. def test_rejects_zero_chunk_size(self):
  80. obj = self._obj()
  81. with self.assertRaises(InvalidValueException):
  82. obj._iter_chunk_size(0)
  83. def test_rejects_negative_chunk_size(self):
  84. obj = self._obj()
  85. with self.assertRaises(InvalidValueException):
  86. obj._iter_chunk_size(-1)
  87. class SaveContentTestCase(unittest.TestCase):
  88. """save_content is defined in terms of iter_content."""
  89. def _obj(self, content, config=None):
  90. return _StreamingObject(_FakeProvider(config), content)
  91. def test_writes_whole_content_to_target_stream(self):
  92. content = bytes(range(256)) * 40
  93. obj = self._obj(content)
  94. target = BytesIO()
  95. obj.save_content(target)
  96. self.assertEqual(target.getvalue(), content)
  97. def test_passes_chunk_size_through_to_iter_content(self):
  98. obj = self._obj(b"x" * 100)
  99. obj.save_content(BytesIO(), chunk_size=16)
  100. self.assertEqual(obj.chunk_sizes_seen, [16])
  101. def test_uses_default_chunk_size_when_unset(self):
  102. obj = self._obj(b"x" * 10)
  103. obj.save_content(BytesIO())
  104. self.assertEqual(obj.chunk_sizes_seen,
  105. [BaseBucketObject.CB_ITER_CHUNK_SIZE])
  106. def test_handles_empty_object(self):
  107. obj = self._obj(b"")
  108. target = BytesIO()
  109. obj.save_content(target)
  110. self.assertEqual(target.getvalue(), b"")
  111. def test_does_not_require_a_readable_iter_content(self):
  112. # iter_content promises Iterable[bytes] and nothing more; providers
  113. # that return a bare generator must still work with save_content.
  114. obj = self._obj(b"abc")
  115. self.assertFalse(hasattr(obj.iter_content(), 'read'))
  116. target = BytesIO()
  117. obj.save_content(target)
  118. self.assertEqual(target.getvalue(), b"abc")
  119. class ProviderIterContentTestCase(unittest.TestCase):
  120. """
  121. The per-provider chunking loops.
  122. The object store service suite only ever exercises the one provider it is
  123. configured against - in CI, the AWS-backed mock - so the loops the other
  124. providers use to turn an SDK handle into sized chunks are pinned here
  125. against fake SDK objects instead.
  126. """
  127. def test_azure_reads_chunk_size_slices_from_one_download(self):
  128. from cloudbridge.providers.azure.resources import AzureBucketObject
  129. content = bytes(range(256)) * 40 # 10 KiB, newline-free by design
  130. reads = []
  131. downloads = []
  132. class _Downloader:
  133. def __init__(self):
  134. self.offset = 0
  135. def read(self, size):
  136. reads.append(size)
  137. data = content[self.offset:self.offset + size]
  138. self.offset += len(data)
  139. return data
  140. class _BlobClient:
  141. def download_blob(self):
  142. downloads.append(1)
  143. return _Downloader()
  144. obj = AzureBucketObject(
  145. _FakeProvider(), _FakeAzureContainer(_BlobClient()),
  146. _FakeBlobProperties("obj"))
  147. chunks = list(obj.iter_content(chunk_size=1024))
  148. self.assertEqual(b"".join(chunks), content)
  149. self.assertEqual([len(c) for c in chunks], [1024] * 10)
  150. self.assertEqual(set(reads), {1024},
  151. "Chunk size must be passed straight to read().")
  152. self.assertEqual(len(downloads), 1,
  153. "The whole object should stream from a single "
  154. "download, not one request per chunk.")
  155. def test_azure_uses_resolved_default_chunk_size(self):
  156. from cloudbridge.providers.azure.resources import AzureBucketObject
  157. reads = []
  158. class _Downloader:
  159. def read(self, size):
  160. reads.append(size)
  161. return b""
  162. class _BlobClient:
  163. def download_blob(self):
  164. return _Downloader()
  165. obj = AzureBucketObject(
  166. _FakeProvider({'iter_chunk_size': 4096}),
  167. _FakeAzureContainer(_BlobClient()), _FakeBlobProperties("obj"))
  168. self.assertEqual(list(obj.iter_content()), [])
  169. self.assertEqual(reads, [4096])
  170. def test_gcp_fetches_successive_ranges_of_chunk_size(self):
  171. from cloudbridge.providers.gcp.resources import GCPBucketObject
  172. content = bytes(range(256)) * 40 # 10 KiB
  173. ranges = []
  174. class _BucketObjects:
  175. def download_range(self, bucket, name, offset, length):
  176. ranges.append((offset, length))
  177. return content[offset:offset + length]
  178. obj = GCPBucketObject(
  179. _FakeProvider(bucket_objects=_BucketObjects()), "BUCKET",
  180. {'name': 'obj', 'size': str(len(content))})
  181. chunks = list(obj.iter_content(chunk_size=4096))
  182. self.assertEqual(b"".join(chunks), content)
  183. self.assertEqual(
  184. ranges, [(0, 4096), (4096, 4096), (8192, 2048)],
  185. "Ranges must tile the object exactly and the last must be "
  186. "clamped to the object size, not overrun it.")
  187. def test_gcp_reads_nothing_for_an_empty_object(self):
  188. from cloudbridge.providers.gcp.resources import GCPBucketObject
  189. obj = GCPBucketObject(
  190. _FakeProvider(bucket_objects=_NoRangeService()), "BUCKET",
  191. {'name': 'obj', 'size': '0'})
  192. self.assertEqual(list(obj.iter_content()), [])
  193. def test_gcp_rejects_bad_chunk_size_before_any_request(self):
  194. from cloudbridge.providers.gcp.resources import GCPBucketObject
  195. obj = GCPBucketObject(
  196. _FakeProvider(bucket_objects=_NoRangeService()), "BUCKET",
  197. {'name': 'obj', 'size': '100'})
  198. with self.assertRaises(InvalidValueException):
  199. obj.iter_content(chunk_size=0)
  200. def test_openstack_passes_chunk_size_as_resp_chunk_size(self):
  201. from cloudbridge.providers.openstack.resources import (
  202. OpenStackBucketObject)
  203. calls = []
  204. class _Swift:
  205. def get_object(self, container, name, resp_chunk_size=None):
  206. calls.append(resp_chunk_size)
  207. return {}, iter([b"data"])
  208. provider = _FakeProvider()
  209. provider.swift = _Swift()
  210. obj = OpenStackBucketObject(
  211. provider, _FakeSwiftContainer(), {'name': 'obj'})
  212. self.assertEqual(list(obj.iter_content(chunk_size=8192)), [b"data"])
  213. self.assertEqual(calls, [8192],
  214. "resp_chunk_size is what makes swiftclient stream "
  215. "rather than return the whole object.")
  216. if __name__ == "__main__":
  217. unittest.main()